Plugins: Support MT app registration (#113348)

This commit is contained in:
Todd Treece
2025-12-02 09:59:46 -05:00
committed by GitHub
parent 77e13f7ef8
commit bdf529c545
22 changed files with 385 additions and 240 deletions
+4 -5
View File
@@ -40,7 +40,6 @@ import (
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/org"
"github.com/grafana/grafana/pkg/services/org/orgtest"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync/installsyncfakes"
"github.com/grafana/grafana/pkg/services/pluginsintegration/managedplugins"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginaccesscontrol"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginassets"
@@ -528,7 +527,7 @@ func callGetPluginAsset(sc *scenarioContext) {
func pluginAssetScenario(t *testing.T, desc string, url string, urlPattern string,
cfg *setting.Cfg, pluginRegistry registry.Service, fn scenarioFunc) {
t.Run(fmt.Sprintf("%s %s", desc, url), func(t *testing.T) {
store, err := pluginstore.NewPluginStoreForTest(pluginRegistry, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
store, err := pluginstore.NewPluginStoreForTest(pluginRegistry, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
hs := HTTPServer{
@@ -643,7 +642,7 @@ func Test_PluginsList_AccessControl(t *testing.T) {
for _, tc := range tcs {
t.Run(tc.desc, func(t *testing.T) {
server := SetupAPITestServer(t, func(hs *HTTPServer) {
store, err := pluginstore.NewPluginStoreForTest(pluginRegistry, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
store, err := pluginstore.NewPluginStoreForTest(pluginRegistry, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
hs.Cfg = setting.NewCfg()
@@ -833,7 +832,7 @@ func Test_PluginsSettings(t *testing.T) {
for _, tc := range tcs {
t.Run(tc.desc, func(t *testing.T) {
server := SetupAPITestServer(t, func(hs *HTTPServer) {
store, err := pluginstore.NewPluginStoreForTest(pluginRegistry, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
store, err := pluginstore.NewPluginStoreForTest(pluginRegistry, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
hs.Cfg = setting.NewCfg()
@@ -903,7 +902,7 @@ func Test_UpdatePluginSetting(t *testing.T) {
t.Run("should return an error when trying to disable an auto-enabled plugin", func(t *testing.T) {
server := SetupAPITestServer(t, func(hs *HTTPServer) {
store, err := pluginstore.NewPluginStoreForTest(pluginRegistry, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
store, err := pluginstore.NewPluginStoreForTest(pluginRegistry, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
hs.Cfg = setting.NewCfg()
+4 -74
View File
@@ -1,29 +1,15 @@
package plugins
import (
"context"
"fmt"
"os"
"github.com/grafana/grafana-app-sdk/k8s"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apiserver/pkg/authorization/authorizer"
"k8s.io/apiserver/pkg/registry/generic"
"k8s.io/apiserver/pkg/registry/rest"
restclient "k8s.io/client-go/rest"
"github.com/grafana/grafana-app-sdk/app"
appsdkapiserver "github.com/grafana/grafana-app-sdk/k8s/apiserver"
"github.com/grafana/grafana-app-sdk/simple"
pluginsappapis "github.com/grafana/grafana/apps/plugins/pkg/apis"
pluginsv0alpha1 "github.com/grafana/grafana/apps/plugins/pkg/apis/plugins/v0alpha1"
pluginsapp "github.com/grafana/grafana/apps/plugins/pkg/app"
"github.com/grafana/grafana/apps/plugins/pkg/app/meta"
"github.com/grafana/grafana/pkg/configprovider"
"github.com/grafana/grafana/pkg/services/apiserver"
"github.com/grafana/grafana/pkg/services/apiserver/appinstaller"
"github.com/grafana/grafana/pkg/services/apiserver/endpoints/request"
)
var (
@@ -32,17 +18,10 @@ var (
)
type AppInstaller struct {
metaManager *meta.ProviderManager
cfgProvider configprovider.ConfigProvider
restConfigProvider apiserver.RestConfigProvider
appsdkapiserver.AppInstaller
}
func RegisterAppInstaller(
cfgProvider configprovider.ConfigProvider,
restConfigProvider apiserver.RestConfigProvider,
) (*AppInstaller, error) {
func ProvideAppInstaller() (*AppInstaller, error) {
grafanaComAPIURL := os.Getenv("GRAFANA_COM_API_URL")
if grafanaComAPIURL == "" {
grafanaComAPIURL = "https://grafana.com/api/plugins"
@@ -51,66 +30,17 @@ func RegisterAppInstaller(
coreProvider := meta.NewCoreProvider()
cloudProvider := meta.NewCloudProvider(grafanaComAPIURL)
metaProviderManager := meta.NewProviderManager(coreProvider, cloudProvider)
specificConfig := &pluginsapp.PluginAppConfig{
MetaProviderManager: metaProviderManager,
}
provider := simple.NewAppProvider(pluginsappapis.LocalManifest(), specificConfig, pluginsapp.New)
appConfig := app.Config{
KubeConfig: restclient.Config{}, // this will be overridden by the installer's InitializeApp method
ManifestData: *pluginsappapis.LocalManifest().ManifestData,
SpecificConfig: specificConfig,
}
i, err := appsdkapiserver.NewDefaultAppInstaller(provider, appConfig, pluginsappapis.NewGoTypeAssociator())
i, err := pluginsapp.ProvideAppInstaller(metaProviderManager)
if err != nil {
return nil, err
}
return &AppInstaller{
metaManager: metaProviderManager,
cfgProvider: cfgProvider,
restConfigProvider: restConfigProvider,
AppInstaller: i,
AppInstaller: i,
}, nil
}
func (p *AppInstaller) InstallAPIs(
server appsdkapiserver.GenericAPIServer,
restOptsGetter generic.RESTOptionsGetter,
) error {
ctx := context.Background()
cfg, err := p.cfgProvider.Get(ctx)
if err != nil {
return err
}
// Create a client factory function that will be called lazily when the client is needed.
// This avoids deadlock issues since the restConfigProvider/API server will not be ready during API installation.
clientFactory := func(ctx context.Context) (*pluginsv0alpha1.PluginClient, error) {
kubeConfig, err := p.restConfigProvider.GetRestConfig(ctx)
if err != nil {
return nil, fmt.Errorf("failed to get rest config: %w", err)
}
clientGenerator := k8s.NewClientRegistry(*kubeConfig, k8s.DefaultClientConfig())
client, err := pluginsv0alpha1.NewPluginClientFromGenerator(clientGenerator)
if err != nil {
return nil, fmt.Errorf("failed to create plugin client: %w", err)
}
return client, nil
}
pluginMetaGVR := pluginsv0alpha1.PluginMetaKind().GroupVersionResource()
replacedStorage := map[schema.GroupVersionResource]rest.Storage{
pluginMetaGVR: pluginsapp.NewPluginMetaStorage(p.metaManager, clientFactory, request.GetNamespaceMapper(cfg)),
}
wrappedServer := &customStorageWrapper{
wrapped: server,
replace: replacedStorage,
}
return p.AppInstaller.InstallAPIs(wrappedServer, restOptsGetter)
}
// GetAuthorizer returns the authorizer for the plugins app.
func (p *AppInstaller) GetAuthorizer() authorizer.Authorizer {
return pluginsapp.GetAuthorizer()
-39
View File
@@ -1,39 +0,0 @@
package plugins
import (
"fmt"
restful "github.com/emicklei/go-restful/v3"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apiserver/pkg/registry/rest"
genericserver "k8s.io/apiserver/pkg/server"
appsdkapiserver "github.com/grafana/grafana-app-sdk/k8s/apiserver"
)
var _ appsdkapiserver.GenericAPIServer = (*customStorageWrapper)(nil)
type customStorageWrapper struct {
wrapped appsdkapiserver.GenericAPIServer
replace map[schema.GroupVersionResource]rest.Storage
}
func (c *customStorageWrapper) InstallAPIGroup(
apiGroupInfo *genericserver.APIGroupInfo,
) error {
if apiGroupInfo == nil || apiGroupInfo.VersionedResourcesStorageMap == nil {
return fmt.Errorf("apiGroupInfo cannot be nil")
}
for gvr, storage := range c.replace {
if _, ok := apiGroupInfo.VersionedResourcesStorageMap[gvr.Version]; !ok {
apiGroupInfo.VersionedResourcesStorageMap[gvr.Version] = map[string]rest.Storage{}
}
apiGroupInfo.VersionedResourcesStorageMap[gvr.Version][gvr.Resource] = storage
}
return c.wrapped.InstallAPIGroup(apiGroupInfo)
}
// RegisteredWebServices implements apiserver.GenericAPIServer.
func (c *customStorageWrapper) RegisteredWebServices() []*restful.WebService {
return []*restful.WebService{}
}
+1 -1
View File
@@ -21,7 +21,7 @@ var WireSet = wire.NewSet(
ProvideBuilderRunners,
playlist.RegisterAppInstaller,
investigations.RegisterApp,
plugins.RegisterAppInstaller,
plugins.ProvideAppInstaller,
shorturl.RegisterAppInstaller,
correlations.RegisterAppInstaller,
rules.RegisterAppInstaller,
@@ -4,12 +4,16 @@ import (
"github.com/grafana/grafana/pkg/infra/tracing"
"github.com/grafana/grafana/pkg/modules"
"github.com/grafana/grafana/pkg/services/accesscontrol"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugininstaller"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore"
"github.com/grafana/grafana/pkg/services/provisioning"
)
const (
// InstallSync is the module name for the install sync service.
InstallSync = installsync.ServiceName
// PluginStore is the module name for the plugin store service.
PluginStore = pluginstore.ServiceName
@@ -46,10 +50,11 @@ func dependencyMap() map[string][]string {
Tracing: {},
GrafanaAPIServer: {Tracing},
PluginStore: {GrafanaAPIServer},
PluginInstaller: {PluginStore},
InstallSync: {PluginStore},
PluginInstaller: {InstallSync},
FixedRolesLoader: {PluginInstaller},
Provisioning: {PluginStore, PluginInstaller, FixedRolesLoader},
Core: {GrafanaAPIServer, PluginStore, PluginInstaller, FixedRolesLoader, Provisioning},
Core: {GrafanaAPIServer, PluginStore, PluginInstaller, FixedRolesLoader, Provisioning, InstallSync},
BackgroundServices: {Core},
}
}
@@ -30,6 +30,7 @@ import (
"github.com/grafana/grafana/pkg/services/notifications"
plugindashboardsservice "github.com/grafana/grafana/pkg/services/plugindashboards/service"
"github.com/grafana/grafana/pkg/services/pluginsintegration/angulardetectorsprovider"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync"
"github.com/grafana/grafana/pkg/services/pluginsintegration/keyretriever/dynamic"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginexternal"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugininstaller"
@@ -73,6 +74,7 @@ func ProvideBackgroundServiceRegistry(
dashboardServiceImpl *service.DashboardServiceImpl,
secretsGarbageCollectionWorker *secretsgarbagecollectionworker.Worker,
fixedRolesLoader *accesscontrol.FixedRolesLoader,
installSync installsync.Syncer,
// Need to make sure these are initialized, is there a better place to put them?
_ dashboardsnapshots.Service,
_ serviceaccounts.Service,
@@ -121,6 +123,7 @@ func ProvideBackgroundServiceRegistry(
dashboardServiceImpl,
secretsGarbageCollectionWorker,
fixedRolesLoader,
installSync,
)
}
+16 -16
View File
@@ -579,12 +579,7 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api
}
errorRegistry := pluginerrs.ProvideErrorTracker()
loaderLoader := loader.ProvideService(pluginManagementCfg, discovery, bootstrap, validate, initialize, terminate, errorRegistry)
clientGenerator := apiserver.ProvideClientGenerator(eventualRestConfigProvider)
syncer, err := installsync.ProvideSyncer(featureToggles, clientGenerator, orgService, configProvider, serverLockService)
if err != nil {
return nil, err
}
pluginstoreService, err := pluginstore.ProvideService(inMemory, sourcesService, loaderLoader, syncer, featureToggles)
pluginstoreService, err := pluginstore.ProvideService(inMemory, sourcesService, loaderLoader, featureToggles)
if err != nil {
return nil, err
}
@@ -789,7 +784,7 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api
if err != nil {
return nil, err
}
appInstaller, err := plugins.RegisterAppInstaller(configProvider, eventualRestConfigProvider)
appInstaller, err := plugins.ProvideAppInstaller()
if err != nil {
return nil, err
}
@@ -854,6 +849,11 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api
dashboardUpdater := service8.ProvideDashboardUpdater(inProcBus, pluginstoreService, service14, importDashboardService, service13, pluginService, dashboardService)
worker := garbagecollectionworker.ProvideWorker(cfg, secureValueMetadataStorage, keeperMetadataStorage, ossKeeperService)
fixedRolesLoader := accesscontrol.ProvideFixedRolesLoader(acimplService, featureToggles)
clientGenerator := apiserver.ProvideClientGenerator(eventualRestConfigProvider)
syncer, err := installsync.ProvideSyncer(featureToggles, clientGenerator, orgService, configProvider, serverLockService, eventualRestConfigProvider, pluginstoreService)
if err != nil {
return nil, err
}
healthService, err := grpcserver.ProvideHealthService(cfg, grpcserverProvider)
if err != nil {
return nil, err
@@ -935,7 +935,7 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api
}
ossUserProtectionImpl := authinfoimpl.ProvideOSSUserProtectionService()
registration := authnimpl.ProvideRegistration(cfg, authnService, orgService, userAuthTokenService, acimplService, permissionRegistry, apikeyService, userService, authService, ossUserProtectionImpl, loginattemptimplService, quotaService, authinfoimplService, renderingService, featureToggles, oauthtokenService, socialService, remoteCache, ldapImpl, ossImpl, tracingService, tempuserService, notificationService)
backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, worker, fixedRolesLoader, serviceImpl, serviceAccountsProxy, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration)
backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, worker, fixedRolesLoader, syncer, serviceImpl, serviceAccountsProxy, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration)
usageStatsProvidersRegistry := usagestatssvcs.ProvideUsageStatsProvidersRegistry(acimplService, userService)
server, err := New(opts, cfg, httpServer, acimplService, provisioningServiceImpl, backgroundServiceRegistry, usageStatsProvidersRegistry, statscollectorService, tracingService, featureToggles, registerer)
if err != nil {
@@ -1231,12 +1231,7 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac
}
errorRegistry := pluginerrs.ProvideErrorTracker()
loaderLoader := loader.ProvideService(pluginManagementCfg, discovery, bootstrap, validate, initialize, terminate, errorRegistry)
clientGenerator := apiserver.ProvideClientGenerator(eventualRestConfigProvider)
syncer, err := installsync.ProvideSyncer(featureToggles, clientGenerator, orgService, configProvider, serverLockService)
if err != nil {
return nil, err
}
pluginstoreService, err := pluginstore.ProvideService(inMemory, sourcesService, loaderLoader, syncer, featureToggles)
pluginstoreService, err := pluginstore.ProvideService(inMemory, sourcesService, loaderLoader, featureToggles)
if err != nil {
return nil, err
}
@@ -1443,7 +1438,7 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac
if err != nil {
return nil, err
}
appInstaller, err := plugins.RegisterAppInstaller(configProvider, eventualRestConfigProvider)
appInstaller, err := plugins.ProvideAppInstaller()
if err != nil {
return nil, err
}
@@ -1508,6 +1503,11 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac
dashboardUpdater := service8.ProvideDashboardUpdater(inProcBus, pluginstoreService, service14, importDashboardService, service13, pluginService, dashboardService)
worker := garbagecollectionworker.ProvideWorker(cfg, secureValueMetadataStorage, keeperMetadataStorage, ossKeeperService)
fixedRolesLoader := accesscontrol.ProvideFixedRolesLoader(acimplService, featureToggles)
clientGenerator := apiserver.ProvideClientGenerator(eventualRestConfigProvider)
syncer, err := installsync.ProvideSyncer(featureToggles, clientGenerator, orgService, configProvider, serverLockService, eventualRestConfigProvider, pluginstoreService)
if err != nil {
return nil, err
}
healthService, err := grpcserver.ProvideHealthService(cfg, grpcserverProvider)
if err != nil {
return nil, err
@@ -1589,7 +1589,7 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac
}
ossUserProtectionImpl := authinfoimpl.ProvideOSSUserProtectionService()
registration := authnimpl.ProvideRegistration(cfg, authnService, orgService, userAuthTokenService, acimplService, permissionRegistry, apikeyService, userService, authService, ossUserProtectionImpl, loginattemptimplService, quotaService, authinfoimplService, renderingService, featureToggles, oauthtokentestService, socialService, remoteCache, ldapImpl, ossImpl, tracingService, tempuserService, notificationServiceMock)
backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, worker, fixedRolesLoader, serviceImpl, serviceAccountsProxy, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration)
backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, worker, fixedRolesLoader, syncer, serviceImpl, serviceAccountsProxy, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration)
usageStatsProvidersRegistry := usagestatssvcs.ProvideUsageStatsProvidersRegistry(acimplService, userService)
server, err := New(opts, cfg, httpServer, acimplService, provisioningServiceImpl, backgroundServiceRegistry, usageStatsProvidersRegistry, statscollectorService, tracingService, featureToggles, registerer)
if err != nil {
@@ -1,20 +1,29 @@
package client
import (
"context"
"fmt"
"strings"
"time"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/client-go/discovery"
"k8s.io/client-go/rest"
)
var (
defaultPollInterval = 500 * time.Millisecond
defaultAvailabilityTimeout = 30 * time.Second
)
type DiscoveryClient interface {
discovery.DiscoveryInterface
GetResourceForKind(gvk schema.GroupVersionKind) (schema.GroupVersionResource, error)
GetKindForResource(gvr schema.GroupVersionResource) (schema.GroupVersionKind, error)
GetPreferredVesion(gr schema.GroupResource) (schema.GroupVersionResource, schema.GroupVersionKind, error)
GetPreferredVersionForKind(gk schema.GroupKind) (schema.GroupVersionResource, schema.GroupVersionKind, error)
WaitForAvailability(ctx context.Context, gv schema.GroupVersion) error
}
type DiscoveryClientImpl struct {
@@ -134,3 +143,13 @@ func (d *DiscoveryClientImpl) GetPreferredVersionForKind(gk schema.GroupKind) (s
}
return schema.GroupVersionResource{}, schema.GroupVersionKind{}, fmt.Errorf("preferred version not found for kind %s in group %s", gk.Kind, gk.Group)
}
func (d *DiscoveryClientImpl) WaitForAvailability(ctx context.Context, gv schema.GroupVersion) error {
return wait.PollUntilContextTimeout(ctx, defaultPollInterval, defaultAvailabilityTimeout, true, func(ctx context.Context) (bool, error) {
_, err := d.ServerResourcesForGroupVersion(gv.String())
if err != nil {
return false, nil
}
return true, nil
})
}
@@ -4,21 +4,29 @@ import (
"context"
"github.com/grafana/grafana/apps/plugins/pkg/app/install"
"github.com/grafana/grafana/pkg/plugins"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore"
)
var _ installsync.Syncer = &FakeSyncer{}
type FakeSyncer struct {
SyncFunc func(ctx context.Context, source install.Source, installedPlugins []*plugins.Plugin) error
SyncFunc func(ctx context.Context, source install.Source, installedPlugins []pluginstore.Plugin) error
}
func NewFakeSyncer() *FakeSyncer {
return &FakeSyncer{}
}
func (f *FakeSyncer) Sync(ctx context.Context, source install.Source, installedPlugins []*plugins.Plugin) error {
func (f *FakeSyncer) IsDisabled() bool {
return false
}
func (f *FakeSyncer) Run(ctx context.Context) error {
return nil
}
func (f *FakeSyncer) Sync(ctx context.Context, source install.Source, installedPlugins []pluginstore.Plugin) error {
if f.SyncFunc != nil {
return f.SyncFunc(ctx, source, installedPlugins)
}
@@ -4,18 +4,25 @@ import (
"context"
"time"
"github.com/grafana/dskit/services"
"github.com/grafana/grafana-app-sdk/logging"
"github.com/grafana/grafana-app-sdk/resource"
pluginsv0alpha1 "github.com/grafana/grafana/apps/plugins/pkg/apis/plugins/v0alpha1"
"github.com/grafana/grafana/apps/plugins/pkg/app/install"
"github.com/grafana/grafana/pkg/apimachinery/identity"
"github.com/grafana/grafana/pkg/configprovider"
"github.com/grafana/grafana/pkg/plugins"
"github.com/grafana/grafana/pkg/registry"
"github.com/grafana/grafana/pkg/services/apiserver"
"github.com/grafana/grafana/pkg/services/apiserver/client"
"github.com/grafana/grafana/pkg/services/apiserver/endpoints/request"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/org"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore"
)
const (
ServiceName = "plugins.installsync"
syncerLockActionName = "plugin-install-api-sync"
)
@@ -25,7 +32,9 @@ var (
// Syncer is the interface for syncing plugin installations to the Kubernetes-style API.
type Syncer interface {
Sync(ctx context.Context, source install.Source, installedPlugins []*plugins.Plugin) error
registry.BackgroundService
registry.CanBeDisabled
Sync(ctx context.Context, source install.Source, installedPlugins []pluginstore.Plugin) error
}
// ServerLock is the interface for acquiring distributed locks.
@@ -34,14 +43,22 @@ type ServerLock interface {
}
type syncer struct {
featureToggles featuremgmt.FeatureToggles
clientGenerator resource.ClientGenerator
installRegistrar *install.InstallRegistrar
orgService org.Service
namespaceMapper request.NamespaceMapper
serverLock ServerLock
services.NamedService
featureToggles featuremgmt.FeatureToggles
clientGenerator resource.ClientGenerator
installRegistrar *install.InstallRegistrar
orgService org.Service
namespaceMapper request.NamespaceMapper
serverLock ServerLock
restConfigProvider apiserver.RestConfigProvider
pluginsStoreService pluginstore.Store
}
var _ Syncer = (*syncer)(nil)
var _ registry.BackgroundService = (*syncer)(nil)
var _ registry.CanBeDisabled = (*syncer)(nil)
var _ services.NamedService = (*syncer)(nil)
// newSyncer creates a new syncer with the provided dependencies.
func newSyncer(
featureToggles featuremgmt.FeatureToggles,
@@ -50,15 +67,21 @@ func newSyncer(
orgService org.Service,
namespaceMapper request.NamespaceMapper,
serverLock ServerLock,
restConfigProvider apiserver.RestConfigProvider,
pluginsStoreService pluginstore.Store,
) *syncer {
return &syncer{
clientGenerator: clientGenerator,
featureToggles: featureToggles,
installRegistrar: installRegistrar,
orgService: orgService,
namespaceMapper: namespaceMapper,
serverLock: serverLock,
s := syncer{
clientGenerator: clientGenerator,
featureToggles: featureToggles,
installRegistrar: installRegistrar,
orgService: orgService,
namespaceMapper: namespaceMapper,
serverLock: serverLock,
restConfigProvider: restConfigProvider,
pluginsStoreService: pluginsStoreService,
}
s.NamedService = services.NewBasicService(nil, s.running, nil).WithName(ServiceName)
return &s
}
// ProvideSyncer creates a new Syncer for syncing plugin installations to the API.
@@ -68,6 +91,8 @@ func ProvideSyncer(
orgService org.Service,
cfgProvider configprovider.ConfigProvider,
serverLock ServerLock,
restConfigProvider apiserver.RestConfigProvider,
pluginsStoreService pluginstore.Store,
) (Syncer, error) {
cfg, err := cfgProvider.Get(context.Background())
if err != nil {
@@ -83,18 +108,49 @@ func ProvideSyncer(
orgService,
namespaceMapper,
serverLock,
restConfigProvider,
pluginsStoreService,
), nil
}
func (s *syncer) Sync(ctx context.Context, source install.Source, installedPlugins []*plugins.Plugin) error {
func (s *syncer) IsDisabled() bool {
//nolint:staticcheck // not yet migrated to OpenFeature
if !s.featureToggles.IsEnabled(ctx, featuremgmt.FlagPluginInstallAPISync) {
return nil
syncEnabled := s.featureToggles.IsEnabled(context.Background(), featuremgmt.FlagPluginInstallAPISync)
//nolint:staticcheck // not yet migrated to OpenFeature
serviceLoadingEnabled := s.featureToggles.IsEnabled(context.Background(), featuremgmt.FlagPluginStoreServiceLoading)
return !syncEnabled || !serviceLoadingEnabled
}
func (s *syncer) Run(ctx context.Context) error {
if err := s.StartAsync(ctx); err != nil {
return err
}
return s.AwaitTerminated(context.Background())
}
func (s *syncer) running(ctx context.Context) error {
ctxLog := logging.FromContext(ctx)
restConfig, err := s.restConfigProvider.GetRestConfig(ctx)
if err != nil {
return err
}
discoveryClient, err := client.NewDiscoveryClient(restConfig)
if err != nil {
ctxLog.Warn("Failed to create discovery client, skipping plugin sync", "error", err)
}
if err := discoveryClient.WaitForAvailability(ctx, pluginsv0alpha1.PluginKind().GroupVersionKind().GroupVersion()); err != nil {
ctxLog.Warn("Failed to wait for plugin API availability, skipping plugin sync", "error", err)
}
//nolint:staticcheck // not yet migrated to OpenFeature
if !s.featureToggles.IsEnabled(ctx, featuremgmt.FlagPluginStoreServiceLoading) {
logging.DefaultLogger.Warn("pluginInstallAPISync is enabled, but pluginStoreServiceLoading is disabled. skipping plugin sync.")
if err := s.Sync(ctx, install.SourcePluginStore, s.pluginsStoreService.Plugins(ctx)); err != nil {
ctxLog.Warn("Failed to sync plugins", "error", err)
}
<-ctx.Done()
return nil
}
func (s *syncer) Sync(ctx context.Context, source install.Source, installedPlugins []pluginstore.Plugin) error {
if s.IsDisabled() {
return nil
}
@@ -113,13 +169,14 @@ func (s *syncer) Sync(ctx context.Context, source install.Source, installedPlugi
return syncErr
}
func (s *syncer) syncAllNamespaces(ctx context.Context, source install.Source, installedPlugins []*plugins.Plugin) error {
func (s *syncer) syncAllNamespaces(ctx context.Context, source install.Source, installedPlugins []pluginstore.Plugin) error {
orgs, err := s.orgService.Search(ctx, &org.SearchOrgsQuery{})
if err != nil {
return err
}
for _, org := range orgs {
ctx = identity.WithServiceIdentityForSingleNamespaceContext(ctx, s.namespaceMapper(org.ID))
err := s.syncNamespace(ctx, s.namespaceMapper(org.ID), source, installedPlugins)
if err != nil {
return err
@@ -129,7 +186,7 @@ func (s *syncer) syncAllNamespaces(ctx context.Context, source install.Source, i
return nil
}
func (s *syncer) syncNamespace(ctx context.Context, namespace string, source install.Source, installedPlugins []*plugins.Plugin) error {
func (s *syncer) syncNamespace(ctx context.Context, namespace string, source install.Source, installedPlugins []pluginstore.Plugin) error {
client, err := s.installRegistrar.GetClient()
if err != nil {
return err
@@ -140,9 +197,9 @@ func (s *syncer) syncNamespace(ctx context.Context, namespace string, source ins
return err
}
installedMap := make(map[string]*plugins.Plugin)
installedMap := make(map[string]struct{})
for _, p := range installedPlugins {
installedMap[p.ID] = p
installedMap[p.ID] = struct{}{}
}
// unregister plugins that are not installed
@@ -18,6 +18,7 @@ import (
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/org"
"github.com/grafana/grafana/pkg/services/org/orgtest"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore"
)
func TestSyncer_Sync(t *testing.T) {
@@ -25,7 +26,7 @@ func TestSyncer_Sync(t *testing.T) {
name string
pluginInstallAPISyncEnabled bool
pluginStoreServiceEnabled bool
installedPlugins []*plugins.Plugin
installedPlugins []pluginstore.Plugin
orgs []*org.OrgDTO
orgServiceError error
serverLockError error
@@ -36,7 +37,7 @@ func TestSyncer_Sync(t *testing.T) {
name: "plugin install API sync feature toggle disabled",
pluginInstallAPISyncEnabled: false,
pluginStoreServiceEnabled: true,
installedPlugins: []*plugins.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore}},
installedPlugins: []pluginstore.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin"}, Class: plugins.ClassCore}},
orgs: []*org.OrgDTO{{ID: 1, Name: "Org 1"}},
expectedError: nil,
expectSyncCalls: 0,
@@ -45,7 +46,7 @@ func TestSyncer_Sync(t *testing.T) {
name: "plugin store service feature toggle disabled",
pluginInstallAPISyncEnabled: true,
pluginStoreServiceEnabled: false,
installedPlugins: []*plugins.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore}},
installedPlugins: []pluginstore.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin"}, Class: plugins.ClassCore}},
orgs: []*org.OrgDTO{{ID: 1, Name: "Org 1"}},
expectedError: nil,
expectSyncCalls: 0,
@@ -54,7 +55,7 @@ func TestSyncer_Sync(t *testing.T) {
name: "both feature toggles enabled, no orgs",
pluginInstallAPISyncEnabled: true,
pluginStoreServiceEnabled: true,
installedPlugins: []*plugins.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore}},
installedPlugins: []pluginstore.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin"}, Class: plugins.ClassCore}},
orgs: []*org.OrgDTO{},
expectedError: nil,
expectSyncCalls: 0,
@@ -63,7 +64,7 @@ func TestSyncer_Sync(t *testing.T) {
name: "both feature toggles enabled, empty installed plugins",
pluginInstallAPISyncEnabled: true,
pluginStoreServiceEnabled: true,
installedPlugins: []*plugins.Plugin{},
installedPlugins: []pluginstore.Plugin{},
orgs: []*org.OrgDTO{{ID: 1, Name: "Org 1"}},
expectedError: nil,
expectSyncCalls: 0,
@@ -72,7 +73,7 @@ func TestSyncer_Sync(t *testing.T) {
name: "both feature toggles enabled, single org",
pluginInstallAPISyncEnabled: true,
pluginStoreServiceEnabled: true,
installedPlugins: []*plugins.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore}},
installedPlugins: []pluginstore.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin"}, Class: plugins.ClassCore}},
orgs: []*org.OrgDTO{{ID: 1, Name: "Org 1"}},
expectedError: nil,
expectSyncCalls: 1,
@@ -81,7 +82,7 @@ func TestSyncer_Sync(t *testing.T) {
name: "both feature toggles enabled, multiple orgs",
pluginInstallAPISyncEnabled: true,
pluginStoreServiceEnabled: true,
installedPlugins: []*plugins.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore}},
installedPlugins: []pluginstore.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin"}, Class: plugins.ClassCore}},
orgs: []*org.OrgDTO{
{ID: 1, Name: "Org 1"},
{ID: 2, Name: "Org 2"},
@@ -94,7 +95,7 @@ func TestSyncer_Sync(t *testing.T) {
name: "org service error",
pluginInstallAPISyncEnabled: true,
pluginStoreServiceEnabled: true,
installedPlugins: []*plugins.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore}},
installedPlugins: []pluginstore.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin"}, Class: plugins.ClassCore}},
orgs: nil,
orgServiceError: errors.New("org service error"),
expectedError: errors.New("org service error"),
@@ -104,7 +105,7 @@ func TestSyncer_Sync(t *testing.T) {
name: "server lock error",
pluginInstallAPISyncEnabled: true,
pluginStoreServiceEnabled: true,
installedPlugins: []*plugins.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore}},
installedPlugins: []pluginstore.Plugin{{JSONData: plugins.JSONData{ID: "test-plugin"}, Class: plugins.ClassCore}},
orgs: []*org.OrgDTO{{ID: 1, Name: "Org 1"}},
serverLockError: errors.New("lock error"),
expectedError: errors.New("lock error"),
@@ -156,6 +157,8 @@ func TestSyncer_Sync(t *testing.T) {
orgService,
func(orgID int64) string { return "org-1" },
serverLock,
nil,
nil,
)
// Execute
@@ -177,7 +180,7 @@ func TestSyncer_Sync(t *testing.T) {
func TestSyncer_syncNamespace(t *testing.T) {
tests := []struct {
name string
installedPlugins []*plugins.Plugin
installedPlugins []pluginstore.Plugin
apiPlugins []pluginsv0alpha1.Plugin
clientListError error
expectedError error
@@ -188,7 +191,7 @@ func TestSyncer_syncNamespace(t *testing.T) {
}{
{
name: "no installed plugins, no API plugins",
installedPlugins: []*plugins.Plugin{},
installedPlugins: []pluginstore.Plugin{},
apiPlugins: []pluginsv0alpha1.Plugin{},
expectedError: nil,
expectedRegCalls: 0,
@@ -196,7 +199,7 @@ func TestSyncer_syncNamespace(t *testing.T) {
},
{
name: "installed plugins only",
installedPlugins: []*plugins.Plugin{
installedPlugins: []pluginstore.Plugin{
{JSONData: plugins.JSONData{ID: "plugin-1", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore},
{JSONData: plugins.JSONData{ID: "plugin-2", Info: plugins.Info{Version: "2.0.0"}}, Class: plugins.ClassExternal},
},
@@ -208,7 +211,7 @@ func TestSyncer_syncNamespace(t *testing.T) {
},
{
name: "API plugins only",
installedPlugins: []*plugins.Plugin{},
installedPlugins: []pluginstore.Plugin{},
apiPlugins: []pluginsv0alpha1.Plugin{
{
ObjectMeta: metav1.ObjectMeta{
@@ -236,7 +239,7 @@ func TestSyncer_syncNamespace(t *testing.T) {
},
{
name: "mixed - some match",
installedPlugins: []*plugins.Plugin{
installedPlugins: []pluginstore.Plugin{
{JSONData: plugins.JSONData{ID: "plugin-1", Info: plugins.Info{Version: "1.0.0"}}, Class: plugins.ClassCore},
{JSONData: plugins.JSONData{ID: "plugin-2", Info: plugins.Info{Version: "2.0.0"}}, Class: plugins.ClassExternal},
{JSONData: plugins.JSONData{ID: "plugin-3", Info: plugins.Info{Version: "3.0.0"}}, Class: plugins.ClassExternal},
@@ -269,7 +272,7 @@ func TestSyncer_syncNamespace(t *testing.T) {
},
{
name: "list error",
installedPlugins: []*plugins.Plugin{},
installedPlugins: []pluginstore.Plugin{},
apiPlugins: []pluginsv0alpha1.Plugin{},
clientListError: errors.New("list error"),
expectedError: errors.New("list error"),
@@ -327,6 +330,8 @@ func TestSyncer_syncNamespace(t *testing.T) {
orgtest.NewOrgServiceFake(),
func(orgID int64) string { return "org-1" },
&fakeServerLock{},
nil,
nil,
)
// Execute
@@ -378,6 +383,8 @@ func TestInstallRegistrar_GetClient(t *testing.T) {
orgtest.NewOrgServiceFake(),
func(orgID int64) string { return "org-1" },
&fakeServerLock{},
nil,
nil,
)
// First call
@@ -14,7 +14,6 @@ import (
"github.com/grafana/grafana/pkg/plugins/manager/registry"
"github.com/grafana/grafana/pkg/services/datasources"
fakeDatasources "github.com/grafana/grafana/pkg/services/datasources/fakes"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync/installsyncfakes"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginconfig"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginsettings"
@@ -42,7 +41,7 @@ func TestGet(t *testing.T) {
cfg := setting.NewCfg()
ds := &fakeDatasources.FakeDataSourceService{}
db := &dbtest.FakeDB{ExpectedError: pluginsettings.ErrPluginSettingNotFound}
store, err := pluginstore.NewPluginStoreForTest(preg, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
store, err := pluginstore.NewPluginStoreForTest(preg, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
pcp := plugincontext.ProvideService(cfg, localcache.ProvideService(),
store, &fakeDatasources.FakeCacheService{},
@@ -10,7 +10,6 @@ import (
"github.com/grafana/grafana/pkg/plugins/manager/pluginfakes"
"github.com/grafana/grafana/pkg/plugins/manager/registry"
"github.com/grafana/grafana/pkg/plugins/repo"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync/installsyncfakes"
"github.com/grafana/grafana/pkg/services/pluginsintegration/managedplugins"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginchecker"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore"
@@ -27,7 +26,7 @@ func TestService_IsDisabled(t *testing.T) {
&setting.Cfg{
PreinstallPluginsAsync: []setting.InstallPlugin{{ID: "myplugin"}},
},
pluginstore.New(registry.NewInMemory(), &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer()),
pluginstore.New(registry.NewInMemory(), &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}),
&pluginfakes.FakePluginInstaller{},
prometheus.NewRegistry(),
&pluginfakes.FakePluginRepo{},
@@ -160,7 +159,7 @@ func TestService_Run(t *testing.T) {
}
installed := 0
installedFromURL := 0
store, err := pluginstore.NewPluginStoreForTest(preg, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
store, err := pluginstore.NewPluginStoreForTest(preg, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
s, err := ProvideService(
&setting.Cfg{
@@ -8,14 +8,12 @@ import (
"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"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync"
)
var _ Store = (*Service)(nil)
@@ -34,18 +32,17 @@ type Store interface {
type Service struct {
services.NamedService
pluginRegistry registry.Service
pluginLoader loader.Service
pluginSources sources.Registry
installsRegistrar installsync.Syncer
loadOnStartup bool
pluginRegistry registry.Service
pluginLoader loader.Service
pluginSources sources.Registry
loadOnStartup bool
}
func ProvideService(pluginRegistry registry.Service, pluginSources sources.Registry,
pluginLoader loader.Service, installsRegistrar installsync.Syncer, features featuremgmt.FeatureToggles) (*Service, error) {
pluginLoader loader.Service, features featuremgmt.FeatureToggles) (*Service, error) {
//nolint:staticcheck // not yet migrated to OpenFeature
if features.IsEnabledGlobally(featuremgmt.FlagPluginStoreServiceLoading) {
s := New(pluginRegistry, pluginLoader, pluginSources, installsRegistrar)
s := New(pluginRegistry, pluginLoader, pluginSources)
s.loadOnStartup = true
return s, nil
}
@@ -56,24 +53,18 @@ 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, installsRegistrar), nil
return New(pluginRegistry, pluginLoader, pluginSources), nil
}
func (s *Service) Run(ctx context.Context) error {
@@ -84,8 +75,8 @@ func (s *Service) Run(ctx context.Context) error {
return s.AwaitTerminated(stopCtx)
}
func NewPluginStoreForTest(pluginRegistry registry.Service, pluginLoader loader.Service, pluginSources sources.Registry, installsRegistrar installsync.Syncer) (*Service, error) {
s := New(pluginRegistry, pluginLoader, pluginSources, installsRegistrar)
func NewPluginStoreForTest(pluginRegistry registry.Service, pluginLoader loader.Service, pluginSources sources.Registry) (*Service, error) {
s := New(pluginRegistry, pluginLoader, pluginSources)
s.loadOnStartup = true
if err := s.StartAsync(context.Background()); err != nil {
return nil, err
@@ -96,12 +87,11 @@ func NewPluginStoreForTest(pluginRegistry registry.Service, pluginLoader loader.
return s, nil
}
func New(pluginRegistry registry.Service, pluginLoader loader.Service, pluginSources sources.Registry, installsRegistrar installsync.Syncer) *Service {
func New(pluginRegistry registry.Service, pluginLoader loader.Service, pluginSources sources.Registry) *Service {
s := &Service{
pluginRegistry: pluginRegistry,
pluginLoader: pluginLoader,
pluginSources: pluginSources,
installsRegistrar: installsRegistrar,
pluginRegistry: pluginRegistry,
pluginLoader: pluginLoader,
pluginSources: pluginSources,
}
s.NamedService = services.NewBasicService(s.starting, s.running, s.stopping).WithName(ServiceName)
return s
@@ -116,21 +106,15 @@ 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,13 +7,11 @@ 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) {
@@ -78,7 +76,7 @@ func TestStore_ProvideService(t *testing.T) {
features = featuremgmt.WithFeatures()
}
service, err := ProvideService(pluginfakes.NewFakePluginRegistry(), srcs, l, installsyncfakes.NewFakeSyncer(), features)
service, err := ProvideService(pluginfakes.NewFakePluginRegistry(), srcs, l, features)
require.Equal(t, tt.expectedLoadOnStartup, service.loadOnStartup)
require.Equal(t, tt.expectedBeforeStart, loadedSrcs)
require.NoError(t, err)
@@ -91,47 +89,6 @@ 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) {
@@ -145,7 +102,7 @@ func TestStore_Plugin(t *testing.T) {
p1.ID: p1,
p2.ID: p2,
},
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
p, exists := ps.Plugin(context.Background(), p1.ID)
@@ -175,7 +132,7 @@ func TestStore_Plugins(t *testing.T) {
p4.ID: p4,
p5.ID: p5,
},
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
ToGrafanaDTO(p1)
@@ -219,7 +176,7 @@ func TestStore_Routes(t *testing.T) {
p5.ID: p5,
p6.ID: p6,
},
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
sr := func(p *plugins.Plugin) *plugins.StaticRoute {
@@ -249,7 +206,7 @@ func TestProcessManager_shutdown(t *testing.T) {
unloaded = true
return nil, nil
},
}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
}, &pluginfakes.FakeSourceRegistry{})
ctx, cancel := context.WithCancel(context.Background())
@@ -282,7 +239,7 @@ func TestProcessManager_shutdown(t *testing.T) {
UnloadFunc: func(_ context.Context, plugin *plugins.Plugin) (*plugins.Plugin, error) {
return nil, expectedErr
},
}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
err = ps.stopping(nil)
@@ -302,7 +259,7 @@ func TestStore_availablePlugins(t *testing.T) {
p1.ID: p1,
p2.ID: p2,
},
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{}, installsyncfakes.NewFakeSyncer())
}, &pluginfakes.FakeLoader{}, &pluginfakes.FakeSourceRegistry{})
require.NoError(t, err)
aps := ps.availablePlugins(context.Background())
@@ -27,7 +27,6 @@ import (
"github.com/grafana/grafana/pkg/plugins/manager/sources"
"github.com/grafana/grafana/pkg/plugins/pluginassets"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/pluginsintegration/installsync/installsyncfakes"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pipeline"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginconfig"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginerrs"
@@ -65,7 +64,7 @@ func CreateIntegrationTestCtx(t *testing.T, cfg *setting.Cfg, coreRegistry *core
Terminator: term,
})
ps, err := pluginstore.NewPluginStoreForTest(reg, l, sources.ProvideService(cfg, pCfg), installsyncfakes.NewFakeSyncer())
ps, err := pluginstore.NewPluginStoreForTest(reg, l, sources.ProvideService(cfg, pCfg))
require.NoError(t, err)
return &IntegrationTestCtx{