Plugins: API sync (#112452)

This commit is contained in:
Todd Treece
2025-10-24 08:09:26 -04:00
committed by GitHub
parent 8b12bbcc55
commit dc77da11cf
25 changed files with 1214 additions and 47 deletions
@@ -6,13 +6,16 @@ import (
"time"
"github.com/grafana/dskit/services"
"golang.org/x/sync/errgroup"
"github.com/grafana/grafana/apps/plugins/pkg/app/install"
"github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/plugins"
"github.com/grafana/grafana/pkg/plugins/manager/loader"
"github.com/grafana/grafana/pkg/plugins/manager/registry"
"github.com/grafana/grafana/pkg/plugins/manager/sources"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"golang.org/x/sync/errgroup"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync"
)
var _ Store = (*Service)(nil)
@@ -31,16 +34,17 @@ type Store interface {
type Service struct {
services.NamedService
pluginRegistry registry.Service
pluginLoader loader.Service
pluginSources sources.Registry
loadOnStartup bool
pluginRegistry registry.Service
pluginLoader loader.Service
pluginSources sources.Registry
installsRegistrar installsync.Syncer
loadOnStartup bool
}
func ProvideService(pluginRegistry registry.Service, pluginSources sources.Registry,
pluginLoader loader.Service, features featuremgmt.FeatureToggles) (*Service, error) {
pluginLoader loader.Service, installsRegistrar installsync.Syncer, features featuremgmt.FeatureToggles) (*Service, error) {
if features.IsEnabledGlobally(featuremgmt.FlagPluginStoreServiceLoading) {
s := New(pluginRegistry, pluginLoader, pluginSources)
s := New(pluginRegistry, pluginLoader, pluginSources, installsRegistrar)
s.loadOnStartup = true
return s, nil
}
@@ -51,19 +55,24 @@ func ProvideService(pluginRegistry registry.Service, pluginSources sources.Regis
logger := log.New("plugin.store")
logger.Info("Loading plugins...")
loadedPluginsToSync := make([]*plugins.Plugin, 0)
for _, ps := range pluginSources.List(ctx) {
loadedPlugins, err := pluginLoader.Load(ctx, ps)
if err != nil {
logger.Error("Loading plugin source failed", "source", ps.PluginClass(ctx), "error", err)
return nil, err
}
loadedPluginsToSync = append(loadedPluginsToSync, loadedPlugins...)
totalPlugins += len(loadedPlugins)
}
if err := installsRegistrar.Sync(ctx, install.SourcePluginStore, loadedPluginsToSync); err != nil {
logger.Error("Syncing plugin installations failed", "error", err)
}
logger.Info("Plugins loaded", "count", totalPlugins, "duration", time.Since(start))
return New(pluginRegistry, pluginLoader, pluginSources), nil
return New(pluginRegistry, pluginLoader, pluginSources, installsRegistrar), nil
}
func (s *Service) Run(ctx context.Context) error {
@@ -74,8 +83,8 @@ func (s *Service) Run(ctx context.Context) error {
return s.AwaitTerminated(stopCtx)
}
func NewPluginStoreForTest(pluginRegistry registry.Service, pluginLoader loader.Service, pluginSources sources.Registry) (*Service, error) {
s := New(pluginRegistry, pluginLoader, pluginSources)
func NewPluginStoreForTest(pluginRegistry registry.Service, pluginLoader loader.Service, pluginSources sources.Registry, installsRegistrar installsync.Syncer) (*Service, error) {
s := New(pluginRegistry, pluginLoader, pluginSources, installsRegistrar)
s.loadOnStartup = true
if err := s.StartAsync(context.Background()); err != nil {
return nil, err
@@ -86,11 +95,12 @@ func NewPluginStoreForTest(pluginRegistry registry.Service, pluginLoader loader.
return s, nil
}
func New(pluginRegistry registry.Service, pluginLoader loader.Service, pluginSources sources.Registry) *Service {
func New(pluginRegistry registry.Service, pluginLoader loader.Service, pluginSources sources.Registry, installsRegistrar installsync.Syncer) *Service {
s := &Service{
pluginRegistry: pluginRegistry,
pluginLoader: pluginLoader,
pluginSources: pluginSources,
pluginRegistry: pluginRegistry,
pluginLoader: pluginLoader,
pluginSources: pluginSources,
installsRegistrar: installsRegistrar,
}
s.NamedService = services.NewBasicService(s.starting, s.running, s.stopping).WithName(ServiceName)
return s
@@ -105,15 +115,21 @@ func (s *Service) starting(ctx context.Context) error {
logger := log.New(ServiceName)
logger.Info("Loading plugins...")
loadedPluginsToSync := make([]*plugins.Plugin, 0)
for _, ps := range s.pluginSources.List(ctx) {
loadedPlugins, err := s.pluginLoader.Load(ctx, ps)
if err != nil {
logger.Error("Loading plugin source failed", "source", ps.PluginClass(ctx), "error", err)
return err
}
loadedPluginsToSync = append(loadedPluginsToSync, loadedPlugins...)
totalPlugins += len(loadedPlugins)
}
if err := s.installsRegistrar.Sync(ctx, install.SourcePluginStore, loadedPluginsToSync); err != nil {
logger.Error("Syncing plugin installations failed", "error", err)
}
logger.Info("Plugins loaded", "count", totalPlugins, "duration", time.Since(start))
return nil
}
@@ -7,11 +7,13 @@ import (
"github.com/stretchr/testify/require"
"github.com/grafana/grafana/apps/plugins/pkg/app/install"
"github.com/grafana/grafana/pkg/plugins"
"github.com/grafana/grafana/pkg/plugins/backendplugin"
"github.com/grafana/grafana/pkg/plugins/log"
"github.com/grafana/grafana/pkg/plugins/manager/pluginfakes"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync/installsyncfakes"
)
func TestStore_ProvideService(t *testing.T) {
@@ -76,7 +78,7 @@ func TestStore_ProvideService(t *testing.T) {
features = featuremgmt.WithFeatures()
}
service, err := ProvideService(pluginfakes.NewFakePluginRegistry(), srcs, l, features)
service, err := ProvideService(pluginfakes.NewFakePluginRegistry(), srcs, l, installsyncfakes.NewFakeSyncer(), features)
require.Equal(t, tt.expectedLoadOnStartup, service.loadOnStartup)
require.Equal(t, tt.expectedBeforeStart, loadedSrcs)
require.NoError(t, err)
@@ -89,6 +91,47 @@ func TestStore_ProvideService(t *testing.T) {
})
}
})
t.Run("Plugin installs are synced", func(t *testing.T) {
registrar := installsyncfakes.NewFakeSyncer()
registered := []*plugins.Plugin{}
registrar.SyncFunc = func(ctx context.Context, source install.Source, installedPlugins []*plugins.Plugin) error {
registered = append(registered, installedPlugins...)
return nil
}
srcs := &pluginfakes.FakeSourceRegistry{ListFunc: func(_ context.Context) []plugins.PluginSource {
return []plugins.PluginSource{
&pluginfakes.FakePluginSource{
PluginClassFunc: func(ctx context.Context) plugins.Class {
return plugins.ClassExternal
},
DiscoverFunc: func(ctx context.Context) ([]*plugins.FoundBundle, error) {
return []*plugins.FoundBundle{
{
Primary: plugins.FoundPlugin{JSONData: plugins.JSONData{ID: "test-plugin"}},
},
}, nil
},
DefaultSignatureFunc: func(ctx context.Context) (plugins.Signature, bool) {
return plugins.Signature{}, false
},
},
}
}}
l := &pluginfakes.FakeLoader{
LoadFunc: func(ctx context.Context, src plugins.PluginSource) ([]*plugins.Plugin, error) {
return []*plugins.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin"}}}, nil
},
}
service, err := ProvideService(pluginfakes.NewFakePluginRegistry(), srcs, l, registrar, featuremgmt.WithFeatures())
require.NoError(t, err)
ctx := context.Background()
err = service.StartAsync(ctx)
require.NoError(t, err)
err = service.AwaitRunning(ctx)
require.NoError(t, err)
require.Len(t, registered, 1)
})
}
func TestStore_Plugin(t *testing.T) {
@@ -102,7 +145,7 @@ func TestStore_Plugin(t *testing.T) {
p1.ID: p1,
p2.ID: p2,
},
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
require.NoError(t, err)
p, exists := ps.Plugin(context.Background(), p1.ID)
@@ -132,7 +175,7 @@ func TestStore_Plugins(t *testing.T) {
p4.ID: p4,
p5.ID: p5,
},
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
require.NoError(t, err)
ToGrafanaDTO(p1)
@@ -176,7 +219,7 @@ func TestStore_Routes(t *testing.T) {
p5.ID: p5,
p6.ID: p6,
},
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
require.NoError(t, err)
sr := func(p *plugins.Plugin) *plugins.StaticRoute {
@@ -206,7 +249,7 @@ func TestProcessManager_shutdown(t *testing.T) {
unloaded = true
return nil, nil
},
}, &pluginfakes.FakeSourceRegistry{})
}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
ctx, cancel := context.WithCancel(context.Background())
@@ -239,7 +282,7 @@ func TestProcessManager_shutdown(t *testing.T) {
UnloadFunc: func(_ context.Context, plugin *plugins.Plugin) (*plugins.Plugin, error) {
return nil, expectedErr
},
}, &pluginfakes.FakeSourceRegistry{})
}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
require.NoError(t, err)
err = ps.stopping(nil)
@@ -259,7 +302,7 @@ func TestStore_availablePlugins(t *testing.T) {
p1.ID: p1,
p2.ID: p2,
},
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
require.NoError(t, err)
aps := ps.availablePlugins(context.Background())