From 84ae9ea71bb907056af8ab5af6c4b752fc07f992 Mon Sep 17 00:00:00 2001 From: Stephanie Hingtgen Date: Wed, 3 Sep 2025 11:30:41 -0600 Subject: [PATCH] Provisioning: Add scaffolding for repo controller (#110543) --- apps/provisioning/pkg/controller/repo.go | 158 ++++++++++++++++++++ pkg/operators/provisioning/README.md | 12 +- pkg/operators/provisioning/config.go | 122 +++++++++++++++ pkg/operators/provisioning/jobs_operator.go | 107 ++----------- pkg/operators/provisioning/repo_operator.go | 65 ++++++++ pkg/operators/register.go | 6 + 6 files changed, 374 insertions(+), 96 deletions(-) create mode 100644 apps/provisioning/pkg/controller/repo.go create mode 100644 pkg/operators/provisioning/config.go create mode 100644 pkg/operators/provisioning/repo_operator.go diff --git a/apps/provisioning/pkg/controller/repo.go b/apps/provisioning/pkg/controller/repo.go new file mode 100644 index 00000000000..4bec1f36f30 --- /dev/null +++ b/apps/provisioning/pkg/controller/repo.go @@ -0,0 +1,158 @@ +package controller + +import ( + "context" + "log/slog" + "os" + "time" + + "github.com/grafana/grafana-app-sdk/logging" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" + + provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1" + typedclient "github.com/grafana/grafana/apps/provisioning/pkg/generated/clientset/versioned/typed/provisioning/v0alpha1" + informerv0alpha1 "github.com/grafana/grafana/apps/provisioning/pkg/generated/informers/externalversions/provisioning/v0alpha1" + listers "github.com/grafana/grafana/apps/provisioning/pkg/generated/listers/provisioning/v0alpha1" +) + +type RepositoryController struct { + client typedclient.ProvisioningV0alpha1Interface + repoLister listers.RepositoryLister + repoSynced cache.InformerSynced + logger logging.Logger + queue workqueue.TypedRateLimitingInterface[string] +} + +func NewRepositoryController( + provisioningClient typedclient.ProvisioningV0alpha1Interface, + repoInformer informerv0alpha1.RepositoryInformer, +) (*RepositoryController, error) { + controller := &RepositoryController{ + client: provisioningClient, + repoLister: repoInformer.Lister(), + repoSynced: repoInformer.Informer().HasSynced, + logger: logging.NewSLogLogger(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{ + Level: slog.LevelDebug, + })), + queue: workqueue.NewTypedRateLimitingQueue[string](workqueue.DefaultTypedControllerRateLimiter[string]()), + } + + _, err := repoInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ + AddFunc: controller.enqueue, + UpdateFunc: func(oldObj, newObj interface{}) { + controller.enqueue(newObj) + }, + DeleteFunc: controller.enqueue, + }) + if err != nil { + return nil, err + } + + return controller, nil +} + +func (c *RepositoryController) Run(ctx context.Context) { + defer c.queue.ShutDown() + + if !cache.WaitForCacheSync(ctx.Done(), c.repoSynced) { + c.logger.Error("Failed to sync informer cache") + return + } + + go func() { + wait.UntilWithContext(ctx, c.runWorker, time.Second) + c.logger.Info("Worker stopped") + }() + + <-ctx.Done() +} + +func (c *RepositoryController) enqueue(obj interface{}) { + key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj) + if err != nil { + c.logger.Error("Couldn't get key for object", "error", err) + return + } + switch repo := obj.(type) { + case *provisioning.Repository: + var eventType string + if repo.DeletionTimestamp != nil { + eventType = "delete" + } else { + eventType = "add/update" + } + c.logger.Debug("Received repository event", + "event_type", eventType, + "key", key, + "namespace", repo.Namespace, + "name", repo.Name, + "generation", repo.Generation) + } + + c.queue.Add(key) +} + +func (c *RepositoryController) runWorker(ctx context.Context) { + for c.processNextWorkItem(ctx) { + } +} + +func (c *RepositoryController) processNextWorkItem(ctx context.Context) bool { + key, quit := c.queue.Get() + if quit { + return false + } + defer c.queue.Done(key) + + logger := c.logger.With("key", key) + logger.Debug("Processing work item from queue") + + err := c.processRepository(ctx, key) + if err == nil { + c.queue.Forget(key) + logger.Debug("Successfully processed work item") + return true + } + + logger.Error("Failed to process repository", "error", err) + c.queue.AddRateLimited(key) + return true +} + +func (c *RepositoryController) processRepository(ctx context.Context, key string) error { + namespace, name, err := cache.SplitMetaNamespaceKey(key) + if err != nil { + return err + } + + repo, err := c.repoLister.Repositories(namespace).Get(name) + if err != nil { + return err + } + c.logger.Debug("Processing repository", + "namespace", repo.Namespace, + "name", repo.Name, + "type", repo.Spec.Type, + "generation", repo.Generation, + "observedGeneration", repo.Status.ObservedGeneration) + + if repo.Generation != repo.Status.ObservedGeneration { + repo.Status.ObservedGeneration = repo.Generation + + _, err = c.client.Repositories(repo.Namespace).UpdateStatus(ctx, repo, metav1.UpdateOptions{}) + if err != nil { + return err + } + + // TODO: do a lot more here :) + c.logger.Debug("Updated repository status", + "namespace", repo.Namespace, + "name", repo.Name, + "observedGeneration", repo.Status.ObservedGeneration) + } + + return nil +} diff --git a/pkg/operators/provisioning/README.md b/pkg/operators/provisioning/README.md index ee64c4583a7..07f3f07eb68 100644 --- a/pkg/operators/provisioning/README.md +++ b/pkg/operators/provisioning/README.md @@ -1,4 +1,8 @@ -# Jobs Controller +# Provisioning Controllers + +Git sync has two different controllers: the jobs controller and the repo controller. + +## Jobs Controller > [!WARNING] > This controller has current limitations: @@ -143,3 +147,9 @@ curl -X POST https://localhost:6446/apis/provisioning.grafana.app/v0alpha1/names 3. In a full setup with the concurrent driver, workers claim and process jobs, updating status and writing history. 4. Entries move to `HistoricJobs`; if cleanup is enabled, older entries are pruned based on `--history-expiration`. + +## [WIP] Repository Controller + +This controller is responsible for watching repositories. It will eventually do health checks, queue sync jobs, and create/delete github hooks. + +To run locally, run `GF_DEFAULT_TARGET=operator GF_OPERATOR_NAME=provisioning-repo ./bin/darwin-arm64/grafana server target --config=conf/operator.ini` \ No newline at end of file diff --git a/pkg/operators/provisioning/config.go b/pkg/operators/provisioning/config.go new file mode 100644 index 00000000000..7c3758577f5 --- /dev/null +++ b/pkg/operators/provisioning/config.go @@ -0,0 +1,122 @@ +package provisioning + +import ( + "crypto/x509" + "fmt" + "net/http" + "os" + "time" + + "github.com/grafana/authlib/authn" + "k8s.io/client-go/rest" + "k8s.io/client-go/transport" + + "github.com/grafana/grafana/pkg/setting" + + authrt "github.com/grafana/grafana/apps/provisioning/pkg/auth" + client "github.com/grafana/grafana/apps/provisioning/pkg/generated/clientset/versioned" +) + +// provisioningControllerConfig contains the configuration that overlaps for the jobs and repo controllers +type provisioningControllerConfig struct { + provisioningClient *client.Clientset + resyncInterval time.Duration +} + +// expects: +// [grpc_client_authentication] +// token = +// token_exchange_url = +// [operator] +// provisioning_server_url = +// tls_insecure = +// tls_cert_file = +// tls_key_file = +// tls_ca_file = +// resync_interval = +func setupFromConfig(cfg *setting.Cfg) (controllerCfg *provisioningControllerConfig, err error) { + if cfg == nil { + return nil, fmt.Errorf("no configuration available") + } + + gRPCAuth := cfg.SectionWithEnvOverrides("grpc_client_authentication") + token := gRPCAuth.Key("token").String() + if token == "" { + return nil, fmt.Errorf("token is required in [grpc_client_authentication] section") + } + tokenExchangeURL := gRPCAuth.Key("token_exchange_url").String() + if tokenExchangeURL == "" { + return nil, fmt.Errorf("token_exchange_url is required in [grpc_client_authentication] section") + } + + operatorSec := cfg.SectionWithEnvOverrides("operator") + provisioningServerURL := operatorSec.Key("provisioning_server_url").String() + 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() + tlsCAFile := operatorSec.Key("tls_ca_file").String() + + tokenExchangeClient, err := authn.NewTokenExchangeClient(authn.TokenExchangeConfig{ + TokenExchangeURL: tokenExchangeURL, + Token: token, + }) + if err != nil { + return nil, fmt.Errorf("failed to create token exchange client: %w", err) + } + + tlsConfig, err := buildTLSConfig(tlsInsecure, tlsCertFile, tlsKeyFile, tlsCAFile) + if err != nil { + return nil, fmt.Errorf("failed to build TLS configuration: %w", err) + } + + config := &rest.Config{ + APIPath: "/apis", + Host: provisioningServerURL, + WrapTransport: transport.WrapperFunc(func(rt http.RoundTripper) http.RoundTripper { + return authrt.NewRoundTripper(tokenExchangeClient, rt) + }), + TLSClientConfig: tlsConfig, + } + + provisioningClient, err := client.NewForConfig(config) + if err != nil { + return nil, fmt.Errorf("failed to create provisioning client: %w", err) + } + + return &provisioningControllerConfig{ + provisioningClient: provisioningClient, + resyncInterval: operatorSec.Key("resync_interval").MustDuration(60 * time.Second), + }, nil +} + +func buildTLSConfig(insecure bool, certFile, keyFile, caFile string) (rest.TLSClientConfig, error) { + tlsConfig := rest.TLSClientConfig{ + Insecure: insecure, + } + + if certFile != "" && keyFile != "" { + tlsConfig.CertFile = certFile + tlsConfig.KeyFile = keyFile + } + + if caFile != "" { + // caFile is set in operator.ini file + // nolint:gosec + caCert, err := os.ReadFile(caFile) + if err != nil { + return tlsConfig, fmt.Errorf("failed to read CA certificate file: %w", err) + } + + caCertPool := x509.NewCertPool() + if !caCertPool.AppendCertsFromPEM(caCert) { + return tlsConfig, fmt.Errorf("failed to parse CA certificate") + } + + tlsConfig.CAData = caCert + } + + return tlsConfig, nil +} diff --git a/pkg/operators/provisioning/jobs_operator.go b/pkg/operators/provisioning/jobs_operator.go index bf32bd850f1..84265091d78 100644 --- a/pkg/operators/provisioning/jobs_operator.go +++ b/pkg/operators/provisioning/jobs_operator.go @@ -2,43 +2,31 @@ package provisioning import ( "context" - "crypto/x509" "fmt" "log/slog" - "net/http" "os" "os/signal" "syscall" "time" - "github.com/grafana/authlib/authn" "github.com/grafana/grafana-app-sdk/logging" "github.com/urfave/cli/v2" - "k8s.io/client-go/rest" "k8s.io/client-go/tools/cache" - "k8s.io/client-go/transport" "github.com/grafana/grafana/pkg/services/apiserver/standalone" "github.com/grafana/grafana/pkg/setting" - authrt "github.com/grafana/grafana/apps/provisioning/pkg/auth" "github.com/grafana/grafana/apps/provisioning/pkg/controller" - client "github.com/grafana/grafana/apps/provisioning/pkg/generated/clientset/versioned" informer "github.com/grafana/grafana/apps/provisioning/pkg/generated/informers/externalversions" ) -type controllerConfig struct { - provisioningClient *client.Clientset - historyExpiration time.Duration -} - func RunJobController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.Cfg) error { logger := logging.NewSLogLogger(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{ Level: slog.LevelDebug, })).With("logger", "provisioning-job-controller") logger.Info("Starting provisioning job controller") - controllerCfg, err := setupFromConfig(cfg) + controllerCfg, err := getJobsControllerConfig(cfg) if err != nil { return fmt.Errorf("failed to setup operator: %w", err) } @@ -57,7 +45,7 @@ func RunJobController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.Cf // Jobs informer and controller (resync ~60s like in register.go) jobInformerFactory := informer.NewSharedInformerFactoryWithOptions( controllerCfg.provisioningClient, - 60*time.Second, + controllerCfg.resyncInterval, ) jobInformer := jobInformerFactory.Provisioning().V0alpha1().Jobs() jobController, err := controller.NewJobController(jobInformer) @@ -113,89 +101,18 @@ func RunJobController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.Cf return nil } -func setupFromConfig(cfg *setting.Cfg) (controllerCfg *controllerConfig, err error) { - if cfg == nil { - return nil, fmt.Errorf("no configuration available") - } +type jobsControllerConfig struct { + provisioningControllerConfig + historyExpiration time.Duration +} - gRPCAuth := cfg.SectionWithEnvOverrides("grpc_client_authentication") - token := gRPCAuth.Key("token").String() - if token == "" { - return nil, fmt.Errorf("token is required in [grpc_client_authentication] section") - } - tokenExchangeURL := gRPCAuth.Key("token_exchange_url").String() - if tokenExchangeURL == "" { - return nil, fmt.Errorf("token_exchange_url is required in [grpc_client_authentication] section") - } - - operatorSec := cfg.SectionWithEnvOverrides("operator") - provisioningServerURL := operatorSec.Key("provisioning_server_url").String() - 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() - tlsCAFile := operatorSec.Key("tls_ca_file").String() - - tokenExchangeClient, err := authn.NewTokenExchangeClient(authn.TokenExchangeConfig{ - TokenExchangeURL: tokenExchangeURL, - Token: token, - }) +func getJobsControllerConfig(cfg *setting.Cfg) (*jobsControllerConfig, error) { + controllerCfg, err := setupFromConfig(cfg) if err != nil { - return nil, fmt.Errorf("failed to create token exchange client: %w", err) + return nil, err } - - tlsConfig, err := buildTLSConfig(tlsInsecure, tlsCertFile, tlsKeyFile, tlsCAFile) - if err != nil { - return nil, fmt.Errorf("failed to build TLS configuration: %w", err) - } - - config := &rest.Config{ - APIPath: "/apis", - Host: provisioningServerURL, - WrapTransport: transport.WrapperFunc(func(rt http.RoundTripper) http.RoundTripper { - return authrt.NewRoundTripper(tokenExchangeClient, rt) - }), - TLSClientConfig: tlsConfig, - } - - provisioningClient, err := client.NewForConfig(config) - if err != nil { - return nil, fmt.Errorf("failed to create provisioning client: %w", err) - } - - return &controllerConfig{ - provisioningClient: provisioningClient, - historyExpiration: operatorSec.Key("history_expiration").MustDuration(0), + return &jobsControllerConfig{ + provisioningControllerConfig: *controllerCfg, + historyExpiration: cfg.SectionWithEnvOverrides("operator").Key("history_expiration").MustDuration(0), }, nil } - -func buildTLSConfig(insecure bool, certFile, keyFile, caFile string) (rest.TLSClientConfig, error) { - tlsConfig := rest.TLSClientConfig{ - Insecure: insecure, - } - - if certFile != "" && keyFile != "" { - tlsConfig.CertFile = certFile - tlsConfig.KeyFile = keyFile - } - - if caFile != "" { - // caFile is set in operator.ini file - // nolint:gosec - caCert, err := os.ReadFile(caFile) - if err != nil { - return tlsConfig, fmt.Errorf("failed to read CA certificate file: %w", err) - } - - caCertPool := x509.NewCertPool() - if !caCertPool.AppendCertsFromPEM(caCert) { - return tlsConfig, fmt.Errorf("failed to parse CA certificate") - } - - tlsConfig.CAData = caCert - } - - return tlsConfig, nil -} diff --git a/pkg/operators/provisioning/repo_operator.go b/pkg/operators/provisioning/repo_operator.go new file mode 100644 index 00000000000..0d3cbe02c37 --- /dev/null +++ b/pkg/operators/provisioning/repo_operator.go @@ -0,0 +1,65 @@ +package provisioning + +import ( + "context" + "fmt" + "log/slog" + "os" + "os/signal" + "syscall" + + "github.com/grafana/grafana-app-sdk/logging" + "github.com/urfave/cli/v2" + "k8s.io/client-go/tools/cache" + + "github.com/grafana/grafana/apps/provisioning/pkg/controller" + "github.com/grafana/grafana/pkg/services/apiserver/standalone" + "github.com/grafana/grafana/pkg/setting" + + informer "github.com/grafana/grafana/apps/provisioning/pkg/generated/informers/externalversions" +) + +func RunRepoController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.Cfg) error { + logger := logging.NewSLogLogger(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{ + Level: slog.LevelDebug, + })).With("logger", "provisioning-repo-controller") + logger.Info("Starting provisioning repo controller") + + controllerCfg, err := setupFromConfig(cfg) + if err != nil { + return fmt.Errorf("failed to setup operator: %w", err) + } + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + sigChan := make(chan os.Signal, 1) + signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM) + go func() { + <-sigChan + fmt.Println("Received shutdown signal, stopping controllers") + cancel() + }() + + informerFactory := informer.NewSharedInformerFactoryWithOptions( + controllerCfg.provisioningClient, + controllerCfg.resyncInterval, + ) + + repoInformer := informerFactory.Provisioning().V0alpha1().Repositories() + controller, err := controller.NewRepositoryController( + controllerCfg.provisioningClient.ProvisioningV0alpha1(), + repoInformer, + ) + if err != nil { + return fmt.Errorf("failed to create repository controller: %w", err) + } + + informerFactory.Start(ctx.Done()) + if !cache.WaitForCacheSync(ctx.Done(), repoInformer.Informer().HasSynced) { + return fmt.Errorf("failed to sync informer cache") + } + + controller.Run(ctx) + return nil +} diff --git a/pkg/operators/register.go b/pkg/operators/register.go index d298eef6452..a69327c26dc 100644 --- a/pkg/operators/register.go +++ b/pkg/operators/register.go @@ -11,4 +11,10 @@ func init() { Description: "Watch provisioning jobs and manage job history cleanup", RunFunc: provisioning.RunJobController, }) + + server.RegisterOperator(server.Operator{ + Name: "provisioning-repo", + Description: "Watch provisioning repositories", + RunFunc: provisioning.RunRepoController, + }) }