From e2913815d33ce15616dc70708bb907f808d86eef Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Roberto=20Jim=C3=A9nez=20S=C3=A1nchez?= Date: Fri, 5 Sep 2025 19:52:57 +0200 Subject: [PATCH] Provisioning: Build resource clients for operators (#110699) --------- Co-authored-by: Stephanie Hingtgen --- pkg/operators/provisioning/config.go | 47 +++++++++++++++++++ pkg/operators/provisioning/jobs_operator.go | 46 ++++++++++++++---- pkg/operators/provisioning/repo_operator.go | 15 ++---- .../apis/provisioning/resources/client.go | 25 ++++++++-- 4 files changed, 111 insertions(+), 22 deletions(-) diff --git a/pkg/operators/provisioning/config.go b/pkg/operators/provisioning/config.go index c92995b08d3..4f48fba8ea9 100644 --- a/pkg/operators/provisioning/config.go +++ b/pkg/operators/provisioning/config.go @@ -1,6 +1,7 @@ package provisioning import ( + "context" "crypto/x509" "fmt" "net/http" @@ -13,6 +14,7 @@ import ( "k8s.io/client-go/transport" "github.com/grafana/grafana/pkg/infra/tracing" + "github.com/grafana/grafana/pkg/services/apiserver" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/unified" "github.com/grafana/grafana/pkg/storage/unified/resource" @@ -33,6 +35,7 @@ type provisioningControllerConfig struct { resyncInterval time.Duration repoFactory repository.Factory unified resources.ResourceStore + clients resources.ClientFactory } // expects: @@ -52,6 +55,8 @@ type provisioningControllerConfig struct { // audiences = // [operator] // provisioning_server_url = +// dashboards_server_url = +// folders_server_url = // tls_insecure = // tls_cert_file = // tls_key_file = @@ -83,6 +88,7 @@ func setupFromConfig(cfg *setting.Cfg) (controllerCfg *provisioningControllerCon if provisioningServerURL == "" { return nil, fmt.Errorf("provisioning_server_url is required in [operator] section") } + tlsInsecure := operatorSec.Key("tls_insecure").MustBool(false) tlsCertFile := operatorSec.Key("tls_cert_file").String() tlsKeyFile := operatorSec.Key("tls_key_file").String() @@ -138,10 +144,37 @@ func setupFromConfig(cfg *setting.Cfg) (controllerCfg *provisioningControllerCon return nil, fmt.Errorf("failed to setup unified storage: %w", err) } + dashboardsServerURL := operatorSec.Key("dashboards_server_url").String() + if provisioningServerURL == "" { + return nil, fmt.Errorf("dashboards_server_url is required in [operator] section") + } + foldersServerURL := operatorSec.Key("folders_server_url").String() + if provisioningServerURL == "" { + return nil, fmt.Errorf("folders_server_url is required in [operator] section") + } + + apiServerURLs := []string{dashboardsServerURL, foldersServerURL} + configProviders := make([]apiserver.RestConfigProvider, len(apiServerURLs)) + + for i, url := range apiServerURLs { + config := &rest.Config{ + APIPath: "/apis", + Host: url, + WrapTransport: transport.WrapperFunc(func(rt http.RoundTripper) http.RoundTripper { + return authrt.NewRoundTripper(tokenExchangeClient, rt) + }), + TLSClientConfig: tlsConfig, + } + configProviders[i] = NewDirectConfigProvider(config) + } + + clients := resources.NewClientFactoryForMultipleAPIServers(configProviders) + return &provisioningControllerConfig{ provisioningClient: provisioningClient, repoFactory: repoFactory, unified: unified, + clients: clients, resyncInterval: operatorSec.Key("resync_interval").MustDuration(60 * time.Second), }, nil } @@ -300,3 +333,17 @@ func setupUnifiedStorageClient(cfg *setting.Cfg, tracer tracing.Tracer, resource return client, nil } + +// directConfigProvider is a simple RestConfigProvider that always returns the same rest.Config +// it implements apiserver.RestConfigProvider +type directConfigProvider struct { + cfg *rest.Config +} + +func NewDirectConfigProvider(cfg *rest.Config) apiserver.RestConfigProvider { + return &directConfigProvider{cfg: cfg} +} + +func (r *directConfigProvider) GetRestConfig(ctx context.Context) (*rest.Config, error) { + return r.cfg, nil +} diff --git a/pkg/operators/provisioning/jobs_operator.go b/pkg/operators/provisioning/jobs_operator.go index 61a42852aad..434ea3f7635 100644 --- a/pkg/operators/provisioning/jobs_operator.go +++ b/pkg/operators/provisioning/jobs_operator.go @@ -15,9 +15,11 @@ import ( "github.com/grafana/grafana/pkg/apimachinery/identity" "github.com/grafana/grafana/pkg/apimachinery/utils" + "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" "github.com/grafana/grafana/pkg/services/apiserver/standalone" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/unified/resourcepb" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "github.com/grafana/grafana/apps/provisioning/pkg/controller" informer "github.com/grafana/grafana/apps/provisioning/pkg/generated/informers/externalversions" @@ -45,10 +47,10 @@ func RunJobController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.Cf cancel() }() - // Use unified storage client for testing purposes. + // Use unified storage client and API clients for testing purposes. // TODO: remove this once the processing logic is in place // https://github.com/grafana/git-ui-sync-project/issues/467 - go temporaryPeriodicCountManagedObjects(ctx, logger, controllerCfg) + go temporaryPeriodicTestClients(ctx, logger, controllerCfg) // Jobs informer and controller (resync ~60s like in register.go) jobInformerFactory := informer.NewSharedInformerFactoryWithOptions( @@ -125,12 +127,12 @@ func getJobsControllerConfig(cfg *setting.Cfg) (*jobsControllerConfig, error) { }, nil } -// Use unified storage client for testing purposes. +// Use unified storage and API clients for testing purposes. // TODO: remove this once the processing logic is in place // https://github.com/grafana/git-ui-sync-project/issues/467 -func temporaryPeriodicCountManagedObjects(ctx context.Context, logger logging.Logger, controllerCfg *jobsControllerConfig) { +func temporaryPeriodicTestClients(ctx context.Context, logger logging.Logger, controllerCfg *jobsControllerConfig) { tick := time.NewTicker(controllerCfg.resyncInterval) - logger.Info("starting periodic managed resource lister", "interval", controllerCfg.resyncInterval.String()) + logger.Info("starting periodic using clients", "interval", controllerCfg.resyncInterval.String()) fetchAndLog := func(ctx context.Context) { ctx, _, err := identity.WithProvisioningIdentity(ctx, "*") // "*" grants us access to all namespaces. if err != nil { @@ -146,11 +148,39 @@ func temporaryPeriodicCountManagedObjects(ctx context.Context, logger logging.Lo } else { if len(resp.Items) == 0 { logger.Info("no managed objects found") - return + } else { + for _, obj := range resp.Items { + logger.Info("manage object counts", "item", obj) + } + } + } + + // List all supported resources + client, err := controllerCfg.clients.Clients(ctx, "") + if err != nil { + logger.Error("failed to get resource clients", "error", err) + return + } + + for kind, gvr := range resources.SupportedProvisioningResources { + logger := logger.With("kind", kind, "gvr", gvr.String()) + logger.Info("fetching resources") + + resourceClient, gvk, err := client.ForResource(ctx, gvr) + if err != nil { + logger.Error("failed to get resource client", "error", err) + continue } - for _, obj := range resp.Items { - logger.Info("manage object counts", "item", obj) + logger = logger.With("gvk", gvk.String()) + list, err := resourceClient.List(ctx, metav1.ListOptions{}) + if err != nil { + logger.Error("failed to list resources", "error", err) + continue + } + + for _, item := range list.Items { + logger.Info("resource", "name", item.GetName(), "namespace", item.GetNamespace()) } } } diff --git a/pkg/operators/provisioning/repo_operator.go b/pkg/operators/provisioning/repo_operator.go index 1830b745849..95b5f842a51 100644 --- a/pkg/operators/provisioning/repo_operator.go +++ b/pkg/operators/provisioning/repo_operator.go @@ -16,6 +16,7 @@ import ( "github.com/grafana/grafana/pkg/registry/apis/provisioning/controller" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" + "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" "github.com/grafana/grafana/pkg/services/apiserver/standalone" "github.com/grafana/grafana/pkg/setting" @@ -49,15 +50,7 @@ func RunRepoController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.C controllerCfg.resyncInterval, ) - /* - // TODO: wire all of this up in order to allow the finalizers to work - clients := resources.NewClientFactory(apiserver.WithoutRestConfig) - store, err := resource.NewResourceClient(nil, nil, nil, nil, nil) - if err != nil { - return fmt.Errorf("create resource client: %w", err) - } - resourceLister := resources.NewResourceListerForMigrations(nil, nil, nil) - */ + resourceLister := resources.NewResourceListerForMigrations(controllerCfg.unified, nil, nil) jobs, err := jobs.NewJobStore(controllerCfg.provisioningClient.ProvisioningV0alpha1(), 30*time.Second) if err != nil { return fmt.Errorf("create API client job store: %w", err) @@ -70,8 +63,8 @@ func RunRepoController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.C controllerCfg.provisioningClient.ProvisioningV0alpha1(), repoInformer, controllerCfg.repoFactory, - nil, // resourceLister -- TODO: needed for finalizers - nil, // clients -- TODO: needed for finalizers + resourceLister, + controllerCfg.clients, jobs, nil, // dualwrite -- standalone operator assumes it is backed by unified storage healthChecker, diff --git a/pkg/registry/apis/provisioning/resources/client.go b/pkg/registry/apis/provisioning/resources/client.go index 3ff93ce15fe..dbf5dd56465 100644 --- a/pkg/registry/apis/provisioning/resources/client.go +++ b/pkg/registry/apis/provisioning/resources/client.go @@ -119,11 +119,30 @@ func NewClientFactory(configProvider apiserver.RestConfigProvider) ClientFactory return &clientFactory{newSingleAPIClients(configProvider)} } -func (f *clientFactory) Clients(ctx context.Context, namespace string) (ResourceClients, error) { - if namespace == "" { - return nil, fmt.Errorf("missing namespace") +// NewClientFactoryForMultipleAPIServers creates a ClientFactory for multiple API servers +func NewClientFactoryForMultipleAPIServers(configProviders []apiserver.RestConfigProvider) ClientFactory { + clientFactories := make([]ClientFactory, len(configProviders)) + + for i, configProvider := range configProviders { + clientFactory := NewClientFactory(configProvider) + clientFactories[i] = clientFactory } + return &multiClientFactory{clientFactories: clientFactories} +} + +type multiClientFactory struct { + clientFactories []ClientFactory +} + +func (m *multiClientFactory) Clients(ctx context.Context, namespace string) (ResourceClients, error) { + for _, clientFactory := range m.clientFactories { + return clientFactory.Clients(ctx, namespace) + } + return nil, fmt.Errorf("no client factories available") +} + +func (f *clientFactory) Clients(ctx context.Context, namespace string) (ResourceClients, error) { return &resourceClients{ namespace: namespace, clientsProvider: f.clientsProvider,