/
githubmirror
/
grafana
Обзор
Документация
Войти
/
githubmirror
/
grafana
Код
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
pkg/services/pluginsintegration/installsync/syncer.go
227 строк
7 KB
Todd Treece
Plugins: Cache install sync reads (#129623)
31 июл 2026, 15:27
Не верифицирован
31 июл 2026, 15:27
06ec4f8
Код
Авторство
О чём код?
package installsync import ( "context" "errors" "fmt" "slices" "strings" "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/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" ) var ( lockTimeout = 10 * time.Minute ) // Syncer is the interface for syncing plugin installations to the Kubernetes-style API. type Syncer interface { registry.BackgroundService registry.CanBeDisabled Sync(ctx context.Context, source install.Source, installedPlugins []pluginstore.Plugin) error } // ServerLock is the interface for acquiring distributed locks. type ServerLock interface { LockExecuteAndRelease(ctx context.Context, actionName string, maxInterval time.Duration, fn func(ctx context.Context)) error } type syncer struct { services.NamedService featureToggles featuremgmt.FeatureToggles 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, installRegistrar *install.InstallRegistrar, orgService org.Service, namespaceMapper request.NamespaceMapper, serverLock ServerLock, restConfigProvider apiserver.RestConfigProvider, pluginsStoreService pluginstore.Store, ) *syncer { s := syncer{ 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. func ProvideSyncer( featureToggles featuremgmt.FeatureToggles, clientGenerator resource.ClientGenerator, 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 { return nil, err } installRegistrar := install.NewInstallRegistrar(logging.DefaultLogger, clientGenerator) namespaceMapper := request.GetNamespaceMapper(cfg) return newSyncer( featureToggles, installRegistrar, orgService, namespaceMapper, serverLock, restConfigProvider, pluginsStoreService, ), nil } func (s *syncer) IsDisabled() bool { //nolint:staticcheck // not yet migrated to OpenFeature 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) } 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 } if len(installedPlugins) == 0 { return nil } var syncErr error lockErr := s.serverLock.LockExecuteAndRelease(ctx, syncerLockActionName, lockTimeout, func(ctx context.Context) { syncErr = s.syncAllNamespaces(ctx, source, installedPlugins) }) if lockErr != nil { return lockErr } return syncErr } 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 } // a namespace failure must not strand the ones after it: the sync runs // once per process, so skipped namespaces stay stale until restart var errs []error for _, org := range orgs { namespace := s.namespaceMapper(org.ID) nsCtx := identity.WithServiceIdentityForSingleNamespaceContext(ctx, namespace) if err := s.syncNamespace(nsCtx, namespace, source, installedPlugins); err != nil { errs = append(errs, fmt.Errorf("sync namespace %q: %w", namespace, err)) } } return errors.Join(errs...) } func (s *syncer) syncNamespace(ctx context.Context, namespace string, source install.Source, installedPlugins []pluginstore.Plugin) error { primaryPlugins := filterPrimaryPlugins(installedPlugins) desired := make([]install.PluginInstall, 0, len(primaryPlugins)) for _, p := range primaryPlugins { desired = append(desired, install.PluginInstall{ ID: p.ID, Version: p.Info.Version, Source: source, Dependencies: pluginDependencyIDs(p), }) } return s.installRegistrar.SyncNamespace(ctx, namespace, source, desired) } func pluginDependencyIDs(p pluginstore.Plugin) []string { dependencies := p.Dependencies.Plugins ids := make([]string, 0, len(dependencies)) for _, dependency := range dependencies { id := strings.TrimSpace(dependency.ID) if id != "" && !slices.Contains(ids, id) { ids = append(ids, id) } } return ids } func filterPrimaryPlugins(plugins []pluginstore.Plugin) []pluginstore.Plugin { primaryPlugins := make([]pluginstore.Plugin, 0, len(plugins)) for _, p := range plugins { if p.Parent != nil || p.IncludedInAppID != "" { continue } primaryPlugins = append(primaryPlugins, p) } return primaryPlugins }