Provisioning: Reuse controller from registry (#110639)
This commit is contained in:
@@ -1,175 +0,0 @@
|
||||
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"
|
||||
"github.com/grafana/grafana/apps/provisioning/pkg/repository"
|
||||
)
|
||||
|
||||
type RepositoryController struct {
|
||||
client typedclient.ProvisioningV0alpha1Interface
|
||||
repoLister listers.RepositoryLister
|
||||
repoSynced cache.InformerSynced
|
||||
logger logging.Logger
|
||||
queue workqueue.TypedRateLimitingInterface[string]
|
||||
repoFactory repository.Factory
|
||||
}
|
||||
|
||||
func NewRepositoryController(
|
||||
provisioningClient typedclient.ProvisioningV0alpha1Interface,
|
||||
repoInformer informerv0alpha1.RepositoryInformer,
|
||||
repoFactory repository.Factory,
|
||||
) (*RepositoryController, error) {
|
||||
controller := &RepositoryController{
|
||||
repoFactory: repoFactory,
|
||||
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)
|
||||
|
||||
// These lines are here only for testing purposes until we use the real controller
|
||||
built, err := c.repoFactory.Build(ctx, repo)
|
||||
if err != nil {
|
||||
c.logger.Error("Failed to build repository instance", "error", err, "namespace", repo.Namespace, "name", repo.Name)
|
||||
} else {
|
||||
results, err := built.Test(ctx)
|
||||
if err != nil {
|
||||
c.logger.Error("Repository test failed", "error", err, "namespace", repo.Namespace, "name", repo.Name)
|
||||
} else {
|
||||
c.logger.Debug("Repository test results", "results", results, "namespace", repo.Namespace, "name", repo.Name)
|
||||
}
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
@@ -7,16 +7,20 @@ import (
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-app-sdk/logging"
|
||||
appcontroller "github.com/grafana/grafana/apps/provisioning/pkg/controller"
|
||||
"github.com/urfave/cli/v2"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
|
||||
"github.com/grafana/grafana/apps/provisioning/pkg/controller"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/controller"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs"
|
||||
"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"
|
||||
"github.com/grafana/grafana/apps/provisioning/pkg/repository"
|
||||
)
|
||||
|
||||
func RunRepoController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.Cfg) error {
|
||||
@@ -25,7 +29,7 @@ func RunRepoController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.C
|
||||
})).With("logger", "provisioning-repo-controller")
|
||||
logger.Info("Starting provisioning repo controller")
|
||||
|
||||
controllerCfg, err := setupFromConfig(cfg)
|
||||
controllerCfg, err := getRepoControllerConfig(cfg)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to setup operator: %w", err)
|
||||
}
|
||||
@@ -46,11 +50,38 @@ 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)
|
||||
*/
|
||||
jobs, err := jobs.NewJobStore(controllerCfg.provisioningClient.ProvisioningV0alpha1(), 30*time.Second)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create API client job store: %w", err)
|
||||
}
|
||||
tester := &repository.Tester{}
|
||||
statusPatcher := appcontroller.NewRepositoryStatusPatcher(controllerCfg.provisioningClient.ProvisioningV0alpha1())
|
||||
healthChecker := controller.NewHealthChecker(
|
||||
tester,
|
||||
statusPatcher,
|
||||
)
|
||||
|
||||
repoInformer := informerFactory.Provisioning().V0alpha1().Repositories()
|
||||
controller, err := controller.NewRepositoryController(
|
||||
controllerCfg.provisioningClient.ProvisioningV0alpha1(),
|
||||
repoInformer,
|
||||
controllerCfg.repoFactory,
|
||||
nil, // resourceLister -- TODO: needed for finalizers
|
||||
nil, // clients -- TODO: needed for finalizers
|
||||
tester,
|
||||
jobs,
|
||||
nil, // dualwrite -- standalone operator assumes it is backed by unified storage
|
||||
healthChecker,
|
||||
statusPatcher,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to create repository controller: %w", err)
|
||||
@@ -61,6 +92,22 @@ func RunRepoController(opts standalone.BuildInfo, c *cli.Context, cfg *setting.C
|
||||
return fmt.Errorf("failed to sync informer cache")
|
||||
}
|
||||
|
||||
controller.Run(ctx)
|
||||
controller.Run(ctx, controllerCfg.workerCount)
|
||||
return nil
|
||||
}
|
||||
|
||||
type repoControllerConfig struct {
|
||||
provisioningControllerConfig
|
||||
workerCount int
|
||||
}
|
||||
|
||||
func getRepoControllerConfig(cfg *setting.Cfg) (*repoControllerConfig, error) {
|
||||
controllerCfg, err := setupFromConfig(cfg)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &repoControllerConfig{
|
||||
provisioningControllerConfig: *controllerCfg,
|
||||
workerCount: cfg.SectionWithEnvOverrides("operator").Key("worker_count").MustInt(1),
|
||||
}, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user