From bd550d2f063673dc14e72d2785d9281c505ba9b7 Mon Sep 17 00:00:00 2001 From: Stephanie Hingtgen Date: Mon, 22 Sep 2025 08:54:50 -0600 Subject: [PATCH] Provisioning: Wire up prometheus (#111444) --- apps/provisioning/go.mod | 2 +- apps/provisioning/pkg/controller/status.go | 9 ++-- .../pkg/controller/status_test.go | 5 ++- pkg/operators/provisioning/config.go | 7 ++-- pkg/operators/provisioning/jobs_operator.go | 22 ++++++---- pkg/operators/provisioning/repo_operator.go | 14 ++++--- .../apis/provisioning/controller/health.go | 5 ++- .../provisioning/controller/health_test.go | 15 +++---- .../provisioning/controller/repository.go | 5 +++ .../provisioning/jobs/concurrent_driver.go | 4 ++ .../apis/provisioning/jobs/delete/worker.go | 5 ++- .../provisioning/jobs/delete/worker_test.go | 41 ++++++++++--------- .../apis/provisioning/jobs/export/worker.go | 4 ++ .../provisioning/jobs/export/worker_test.go | 35 ++++++++-------- .../apis/provisioning/jobs/move/worker.go | 5 ++- .../provisioning/jobs/move/worker_test.go | 39 +++++++++--------- .../apis/provisioning/jobs/persistentstore.go | 12 ++++-- .../apis/provisioning/jobs/sync/worker.go | 6 +++ .../provisioning/jobs/sync/worker_test.go | 6 ++- pkg/registry/apis/provisioning/register.go | 18 +++++--- .../apis/provisioning/webhooks/register.go | 7 +++- .../apis/provisioning/webhooks/webhook.go | 4 ++ pkg/server/wire_gen.go | 4 +- 23 files changed, 169 insertions(+), 105 deletions(-) diff --git a/apps/provisioning/go.mod b/apps/provisioning/go.mod index 1ffeb70c4f3..732f69c3671 100644 --- a/apps/provisioning/go.mod +++ b/apps/provisioning/go.mod @@ -11,6 +11,7 @@ require ( github.com/grafana/grafana/pkg/apimachinery v0.0.0-20250804150913-990f1c69ecc2 github.com/grafana/nanogit v0.0.0-20250723104447-68f58f5ecec0 github.com/migueleliasweb/go-github-mock v1.1.0 + github.com/prometheus/client_golang v1.23.2 github.com/stretchr/testify v1.11.1 golang.org/x/oauth2 v0.30.0 k8s.io/apimachinery v0.34.1 @@ -53,7 +54,6 @@ require ( github.com/patrickmn/go-cache v2.1.0+incompatible // indirect github.com/pkg/errors v0.9.1 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect - github.com/prometheus/client_golang v1.23.2 // indirect github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/common v0.66.1 // indirect github.com/prometheus/procfs v0.16.1 // indirect diff --git a/apps/provisioning/pkg/controller/status.go b/apps/provisioning/pkg/controller/status.go index 40ed29624c1..ad63a169415 100644 --- a/apps/provisioning/pkg/controller/status.go +++ b/apps/provisioning/pkg/controller/status.go @@ -7,17 +7,20 @@ import ( provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1" client "github.com/grafana/grafana/apps/provisioning/pkg/generated/clientset/versioned/typed/provisioning/v0alpha1" + "github.com/prometheus/client_golang/prometheus" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" ) type RepositoryStatusPatcher struct { - client client.ProvisioningV0alpha1Interface + client client.ProvisioningV0alpha1Interface + registry prometheus.Registerer } -func NewRepositoryStatusPatcher(client client.ProvisioningV0alpha1Interface) *RepositoryStatusPatcher { +func NewRepositoryStatusPatcher(client client.ProvisioningV0alpha1Interface, registry prometheus.Registerer) *RepositoryStatusPatcher { return &RepositoryStatusPatcher{ - client: client, + client: client, + registry: registry, } } diff --git a/apps/provisioning/pkg/controller/status_test.go b/apps/provisioning/pkg/controller/status_test.go index 77e7f737ea5..12c21fb8c5b 100644 --- a/apps/provisioning/pkg/controller/status_test.go +++ b/apps/provisioning/pkg/controller/status_test.go @@ -8,6 +8,7 @@ import ( provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1" "github.com/grafana/grafana/apps/provisioning/pkg/generated/clientset/versioned/typed/provisioning/v0alpha1/fake" + "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/require" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" @@ -16,7 +17,7 @@ import ( func TestNewRepositoryStatusPatcher(t *testing.T) { client := &fake.FakeProvisioningV0alpha1{} - patcher := NewRepositoryStatusPatcher(client) + patcher := NewRepositoryStatusPatcher(client, prometheus.DefaultRegisterer) require.NotNil(t, patcher) require.Equal(t, client, patcher.client) } @@ -103,7 +104,7 @@ func TestRepositoryStatusPatcher_Patch(t *testing.T) { client.AddReactor("patch", "repositories", tt.reactorFunc) } - patcher := NewRepositoryStatusPatcher(&client) + patcher := NewRepositoryStatusPatcher(&client, prometheus.DefaultRegisterer) err := patcher.Patch(context.Background(), tt.repo, tt.patchOperations...) if tt.expectedError != "" { diff --git a/pkg/operators/provisioning/config.go b/pkg/operators/provisioning/config.go index 98d1fb78d5e..273070ac156 100644 --- a/pkg/operators/provisioning/config.go +++ b/pkg/operators/provisioning/config.go @@ -69,7 +69,7 @@ type provisioningControllerConfig struct { // local_permitted_prefixes = // [provisioning] // repository_types = -func setupFromConfig(cfg *setting.Cfg) (controllerCfg *provisioningControllerConfig, err error) { +func setupFromConfig(cfg *setting.Cfg, registry prometheus.Registerer) (controllerCfg *provisioningControllerConfig, err error) { if cfg == nil { return nil, fmt.Errorf("no configuration available") } @@ -130,7 +130,7 @@ func setupFromConfig(cfg *setting.Cfg) (controllerCfg *provisioningControllerCon return nil, fmt.Errorf("failed to setup decrypter: %w", err) } - repoFactory, err := setupRepoFactory(cfg, decrypter, provisioningClient) + repoFactory, err := setupRepoFactory(cfg, decrypter, provisioningClient, registry) if err != nil { return nil, fmt.Errorf("failed to setup repository getter: %w", err) } @@ -220,6 +220,7 @@ func setupRepoFactory( cfg *setting.Cfg, decrypter repository.Decrypter, provisioningClient *client.Clientset, + registry prometheus.Registerer, ) (repository.Factory, error) { operatorSec := cfg.SectionWithEnvOverrides("operator") provisioningSec := cfg.SectionWithEnvOverrides("provisioning") @@ -246,7 +247,7 @@ func setupRepoFactory( var webhook *webhooks.WebhookExtraBuilder provisioningAppURL := operatorSec.Key("provisioning_server_public_url").String() if provisioningAppURL != "" { - webhook = webhooks.ProvideWebhooks(provisioningAppURL) + webhook = webhooks.ProvideWebhooks(provisioningAppURL, registry) } extras = append(extras, github.Extra( diff --git a/pkg/operators/provisioning/jobs_operator.go b/pkg/operators/provisioning/jobs_operator.go index 1da6600c724..49fa2676717 100644 --- a/pkg/operators/provisioning/jobs_operator.go +++ b/pkg/operators/provisioning/jobs_operator.go @@ -10,6 +10,7 @@ import ( "time" "github.com/grafana/grafana-app-sdk/logging" + "github.com/prometheus/client_golang/prometheus" "k8s.io/client-go/tools/cache" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" @@ -33,7 +34,7 @@ func RunJobController(deps server.OperatorDependencies) error { })).With("logger", "provisioning-job-controller") logger.Info("Starting provisioning job controller") - controllerCfg, err := setupJobsControllerFromConfig(deps.Config) + controllerCfg, err := setupJobsControllerFromConfig(deps.Config, deps.Registerer) if err != nil { return fmt.Errorf("failed to setup operator: %w", err) } @@ -94,12 +95,12 @@ func RunJobController(deps server.OperatorDependencies) error { // } jobHistoryWriter := jobs.NewAPIClientHistoryWriter(controllerCfg.provisioningClient.ProvisioningV0alpha1()) - jobStore, err := jobs.NewJobStore(controllerCfg.provisioningClient.ProvisioningV0alpha1(), 30*time.Second) + jobStore, err := jobs.NewJobStore(controllerCfg.provisioningClient.ProvisioningV0alpha1(), 30*time.Second, deps.Registerer) if err != nil { return fmt.Errorf("create API client job store: %w", err) } - workers, err := setupWorkers(controllerCfg) + workers, err := setupWorkers(controllerCfg, deps.Registerer) if err != nil { return fmt.Errorf("setup workers: %w", err) } @@ -120,6 +121,7 @@ func RunJobController(deps server.OperatorDependencies) error { repoGetter, jobHistoryWriter, jobController.InsertNotifications(), + deps.Registerer, workers..., ) if err != nil { @@ -151,8 +153,8 @@ type jobsControllerConfig struct { historyExpiration time.Duration } -func setupJobsControllerFromConfig(cfg *setting.Cfg) (*jobsControllerConfig, error) { - controllerCfg, err := setupFromConfig(cfg) +func setupJobsControllerFromConfig(cfg *setting.Cfg, registry prometheus.Registerer) (*jobsControllerConfig, error) { + controllerCfg, err := setupFromConfig(cfg, registry) if err != nil { return nil, err } @@ -163,12 +165,12 @@ func setupJobsControllerFromConfig(cfg *setting.Cfg) (*jobsControllerConfig, err }, nil } -func setupWorkers(controllerCfg *jobsControllerConfig) ([]jobs.Worker, error) { +func setupWorkers(controllerCfg *jobsControllerConfig, registry prometheus.Registerer) ([]jobs.Worker, error) { clients := controllerCfg.clients parsers := resources.NewParserFactory(clients) resourceLister := resources.NewResourceLister(controllerCfg.unified) repositoryResources := resources.NewRepositoryResourcesFactory(parsers, clients, resourceLister) - statusPatcher := controller.NewRepositoryStatusPatcher(controllerCfg.provisioningClient.ProvisioningV0alpha1()) + statusPatcher := controller.NewRepositoryStatusPatcher(controllerCfg.provisioningClient.ProvisioningV0alpha1(), registry) workers := make([]jobs.Worker, 0) @@ -180,6 +182,7 @@ func setupWorkers(controllerCfg *jobsControllerConfig) ([]jobs.Worker, error) { nil, // HACK: we have updated the worker to check for nil statusPatcher.Patch, syncer, + registry, ) workers = append(workers, syncWorker) @@ -190,6 +193,7 @@ func setupWorkers(controllerCfg *jobsControllerConfig) ([]jobs.Worker, error) { repositoryResources, export.ExportAll, stageIfPossible, + registry, ) workers = append(workers, exportWorker) @@ -204,11 +208,11 @@ func setupWorkers(controllerCfg *jobsControllerConfig) ([]jobs.Worker, error) { workers = append(workers, migrationWorker) // Delete - deleteWorker := deletepkg.NewWorker(syncWorker, stageIfPossible, repositoryResources) + deleteWorker := deletepkg.NewWorker(syncWorker, stageIfPossible, repositoryResources, registry) workers = append(workers, deleteWorker) // Move - moveWorker := move.NewWorker(syncWorker, stageIfPossible, repositoryResources) + moveWorker := move.NewWorker(syncWorker, stageIfPossible, repositoryResources, registry) workers = append(workers, moveWorker) return workers, nil diff --git a/pkg/operators/provisioning/repo_operator.go b/pkg/operators/provisioning/repo_operator.go index 6162007f8e4..b64cb991777 100644 --- a/pkg/operators/provisioning/repo_operator.go +++ b/pkg/operators/provisioning/repo_operator.go @@ -11,6 +11,7 @@ import ( "github.com/grafana/grafana-app-sdk/logging" appcontroller "github.com/grafana/grafana/apps/provisioning/pkg/controller" + "github.com/prometheus/client_golang/prometheus" "k8s.io/client-go/tools/cache" "github.com/grafana/grafana/pkg/registry/apis/provisioning/controller" @@ -28,7 +29,7 @@ func RunRepoController(deps server.OperatorDependencies) error { })).With("logger", "provisioning-repo-controller") logger.Info("Starting provisioning repo controller") - controllerCfg, err := getRepoControllerConfig(deps.Config) + controllerCfg, err := getRepoControllerConfig(deps.Config, deps.Registerer) if err != nil { return fmt.Errorf("failed to setup operator: %w", err) } @@ -50,12 +51,12 @@ func RunRepoController(deps server.OperatorDependencies) error { ) resourceLister := resources.NewResourceLister(controllerCfg.unified) - jobs, err := jobs.NewJobStore(controllerCfg.provisioningClient.ProvisioningV0alpha1(), 30*time.Second) + jobs, err := jobs.NewJobStore(controllerCfg.provisioningClient.ProvisioningV0alpha1(), 30*time.Second, deps.Registerer) if err != nil { return fmt.Errorf("create API client job store: %w", err) } - statusPatcher := appcontroller.NewRepositoryStatusPatcher(controllerCfg.provisioningClient.ProvisioningV0alpha1()) - healthChecker := controller.NewHealthChecker(statusPatcher) + statusPatcher := appcontroller.NewRepositoryStatusPatcher(controllerCfg.provisioningClient.ProvisioningV0alpha1(), deps.Registerer) + healthChecker := controller.NewHealthChecker(statusPatcher, deps.Registerer) repoInformer := informerFactory.Provisioning().V0alpha1().Repositories() controller, err := controller.NewRepositoryController( @@ -68,6 +69,7 @@ func RunRepoController(deps server.OperatorDependencies) error { nil, // dualwrite -- standalone operator assumes it is backed by unified storage healthChecker, statusPatcher, + deps.Registerer, ) if err != nil { return fmt.Errorf("failed to create repository controller: %w", err) @@ -87,8 +89,8 @@ type repoControllerConfig struct { workerCount int } -func getRepoControllerConfig(cfg *setting.Cfg) (*repoControllerConfig, error) { - controllerCfg, err := setupFromConfig(cfg) +func getRepoControllerConfig(cfg *setting.Cfg, registry prometheus.Registerer) (*repoControllerConfig, error) { + controllerCfg, err := setupFromConfig(cfg, registry) if err != nil { return nil, err } diff --git a/pkg/registry/apis/provisioning/controller/health.go b/pkg/registry/apis/provisioning/controller/health.go index f6bad8fff6b..0d24e457274 100644 --- a/pkg/registry/apis/provisioning/controller/health.go +++ b/pkg/registry/apis/provisioning/controller/health.go @@ -7,6 +7,7 @@ import ( provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1" "github.com/grafana/grafana/apps/provisioning/pkg/repository" + "github.com/prometheus/client_golang/prometheus" ) // StatusPatcher defines the interface for updating repository status @@ -19,12 +20,14 @@ type StatusPatcher interface { // HealthChecker provides unified health checking for repositories type HealthChecker struct { statusPatcher StatusPatcher + registry prometheus.Registerer } // NewHealthChecker creates a new health checker -func NewHealthChecker(statusPatcher StatusPatcher) *HealthChecker { +func NewHealthChecker(statusPatcher StatusPatcher, registry prometheus.Registerer) *HealthChecker { return &HealthChecker{ statusPatcher: statusPatcher, + registry: registry, } } diff --git a/pkg/registry/apis/provisioning/controller/health_test.go b/pkg/registry/apis/provisioning/controller/health_test.go index a42cd35c7aa..c9081875b2a 100644 --- a/pkg/registry/apis/provisioning/controller/health_test.go +++ b/pkg/registry/apis/provisioning/controller/health_test.go @@ -6,6 +6,7 @@ import ( "testing" "time" + "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -18,7 +19,7 @@ import ( func TestNewHealthChecker(t *testing.T) { mockPatcher := mocks.NewStatusPatcher(t) - hc := NewHealthChecker(mockPatcher) + hc := NewHealthChecker(mockPatcher, prometheus.DefaultRegisterer) assert.NotNil(t, hc) assert.Equal(t, mockPatcher, hc.statusPatcher) @@ -135,7 +136,7 @@ func TestShouldCheckHealth(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { mockPatcher := mocks.NewStatusPatcher(t) - hc := NewHealthChecker(mockPatcher) + hc := NewHealthChecker(mockPatcher, prometheus.DefaultRegisterer) result := hc.ShouldCheckHealth(tt.repo) assert.Equal(t, tt.expected, result) @@ -222,7 +223,7 @@ func TestHasRecentFailure(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { mockPatcher := mocks.NewStatusPatcher(t) - hc := NewHealthChecker(mockPatcher) + hc := NewHealthChecker(mockPatcher, prometheus.DefaultRegisterer) result := hc.HasRecentFailure(tt.healthStatus, tt.failureType) assert.Equal(t, tt.expected, result) @@ -264,7 +265,7 @@ func TestRecordFailure(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { mockPatcher := mocks.NewStatusPatcher(t) - hc := NewHealthChecker(mockPatcher) + hc := NewHealthChecker(mockPatcher, prometheus.DefaultRegisterer) repo := &provisioning.Repository{ Status: provisioning.RepositoryStatus{ @@ -309,7 +310,7 @@ func TestRecordFailure(t *testing.T) { func TestRecordFailureFunction(t *testing.T) { mockPatcher := mocks.NewStatusPatcher(t) - hc := NewHealthChecker(mockPatcher) + hc := NewHealthChecker(mockPatcher, prometheus.DefaultRegisterer) testErr := errors.New("test error") result := hc.recordFailure(provisioning.HealthFailureHook, testErr) @@ -446,7 +447,7 @@ func TestRefreshHealth(t *testing.T) { testError: tt.testError, } - hc := NewHealthChecker(mockPatcher) + hc := NewHealthChecker(mockPatcher, prometheus.DefaultRegisterer) if tt.expectPatch { if tt.patchError != nil { @@ -556,7 +557,7 @@ func TestHasHealthStatusChanged(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { mockPatcher := mocks.NewStatusPatcher(t) - hc := NewHealthChecker(mockPatcher) + hc := NewHealthChecker(mockPatcher, prometheus.DefaultRegisterer) result := hc.hasHealthStatusChanged(tt.old, tt.new) assert.Equal(t, tt.expected, result) diff --git a/pkg/registry/apis/provisioning/controller/repository.go b/pkg/registry/apis/provisioning/controller/repository.go index 1b59c5ac2ef..9ddaf4fa476 100644 --- a/pkg/registry/apis/provisioning/controller/repository.go +++ b/pkg/registry/apis/provisioning/controller/repository.go @@ -25,6 +25,7 @@ import ( "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" + "github.com/prometheus/client_golang/prometheus" ) const loggerName = "provisioning-repository-controller" @@ -59,6 +60,8 @@ type RepositoryController struct { keyFunc func(obj any) (string, error) queue workqueue.TypedRateLimitingInterface[*queueItem] + + registry prometheus.Registerer } // NewRepositoryController creates new RepositoryController. @@ -72,6 +75,7 @@ func NewRepositoryController( dualwrite dualwrite.Service, healthChecker *HealthChecker, statusPatcher StatusPatcher, + registry prometheus.Registerer, ) (*RepositoryController, error) { rc := &RepositoryController{ client: provisioningClient, @@ -93,6 +97,7 @@ func NewRepositoryController( jobs: jobs, logger: logging.DefaultLogger.With("logger", loggerName), dualwrite: dualwrite, + registry: registry, } _, err := repoInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ diff --git a/pkg/registry/apis/provisioning/jobs/concurrent_driver.go b/pkg/registry/apis/provisioning/jobs/concurrent_driver.go index 0504fc414d5..4da84e4ec84 100644 --- a/pkg/registry/apis/provisioning/jobs/concurrent_driver.go +++ b/pkg/registry/apis/provisioning/jobs/concurrent_driver.go @@ -7,6 +7,7 @@ import ( "time" "github.com/grafana/grafana-app-sdk/logging" + "github.com/prometheus/client_golang/prometheus" ) // ConcurrentJobDriver manages multiple jobDriver instances for concurrent job processing. @@ -20,6 +21,7 @@ type ConcurrentJobDriver struct { repoGetter RepoGetter historicJobs HistoryWriter workers []Worker + registry prometheus.Registerer notifications chan struct{} } @@ -31,6 +33,7 @@ func NewConcurrentJobDriver( repoGetter RepoGetter, historicJobs HistoryWriter, notifications chan struct{}, + registry prometheus.Registerer, workers ...Worker, ) (*ConcurrentJobDriver, error) { if numDrivers <= 0 { @@ -66,6 +69,7 @@ func NewConcurrentJobDriver( historicJobs: historicJobs, workers: workers, notifications: notifications, + registry: registry, }, nil } diff --git a/pkg/registry/apis/provisioning/jobs/delete/worker.go b/pkg/registry/apis/provisioning/jobs/delete/worker.go index 11d78e73502..738876d1517 100644 --- a/pkg/registry/apis/provisioning/jobs/delete/worker.go +++ b/pkg/registry/apis/provisioning/jobs/delete/worker.go @@ -12,19 +12,22 @@ import ( "github.com/grafana/grafana/apps/provisioning/pkg/repository" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" + "github.com/prometheus/client_golang/prometheus" ) type Worker struct { syncWorker jobs.Worker wrapFn repository.WrapWithStageFn resourcesFactory resources.RepositoryResourcesFactory + registry prometheus.Registerer } -func NewWorker(syncWorker jobs.Worker, wrapFn repository.WrapWithStageFn, resourcesFactory resources.RepositoryResourcesFactory) *Worker { +func NewWorker(syncWorker jobs.Worker, wrapFn repository.WrapWithStageFn, resourcesFactory resources.RepositoryResourcesFactory, registry prometheus.Registerer) *Worker { return &Worker{ syncWorker: syncWorker, wrapFn: wrapFn, resourcesFactory: resourcesFactory, + registry: registry, } } diff --git a/pkg/registry/apis/provisioning/jobs/delete/worker_test.go b/pkg/registry/apis/provisioning/jobs/delete/worker_test.go index 1bcfed8d4cc..c5fd53282f8 100644 --- a/pkg/registry/apis/provisioning/jobs/delete/worker_test.go +++ b/pkg/registry/apis/provisioning/jobs/delete/worker_test.go @@ -14,6 +14,7 @@ import ( "github.com/grafana/grafana/apps/provisioning/pkg/repository" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" + "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" ) @@ -73,7 +74,7 @@ func TestDeleteWorker_IsSupported(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) result := worker.IsSupported(context.Background(), tt.job) require.Equal(t, tt.expected, result) }) @@ -87,7 +88,7 @@ func TestDeleteWorker_ProcessMissingDeleteSettings(t *testing.T) { }, } - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), nil, job, nil) require.EqualError(t, err, "missing delete settings") } @@ -116,7 +117,7 @@ func TestDeleteWorker_ProcessNotReaderWriter(t *testing.T) { mockProgress.On("SetTotal", mock.Anything, 1).Return() mockProgress.On("StrictMaxErrors", 1).Return() - worker := NewWorker(nil, mockWrapFn.Execute, nil) + worker := NewWorker(nil, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "delete files from repository: delete job submitted targeting repository that is not a ReaderWriter") } @@ -139,7 +140,7 @@ func TestDeleteWorker_ProcessWrapFnError(t *testing.T) { mockProgress.On("SetTotal", mock.Anything, 1).Return() mockProgress.On("StrictMaxErrors", 1).Return() - worker := NewWorker(nil, mockWrapFn.Execute, nil) + worker := NewWorker(nil, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "delete files from repository: stage failed") } @@ -187,7 +188,7 @@ func TestDeleteWorker_ProcessDeleteFilesSuccess(t *testing.T) { return result.Path == "test/path2" && result.Action == repository.FileActionDeleted && result.Error == nil })).Return() - worker := NewWorker(nil, mockWrapFn.Execute, nil) + worker := NewWorker(nil, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) } @@ -225,7 +226,7 @@ func TestDeleteWorker_ProcessDeleteFilesWithError(t *testing.T) { })).Return() mockProgress.On("TooManyErrors").Return(errors.New("too many errors")) - worker := NewWorker(nil, mockWrapFn.Execute, nil) + worker := NewWorker(nil, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "delete files from repository: too many errors") } @@ -269,7 +270,7 @@ func TestDeleteWorker_ProcessWithSyncWorker(t *testing.T) { return syncJob.Spec.Pull != nil && !syncJob.Spec.Pull.Incremental }), mockProgress).Return(nil) - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, nil) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) } @@ -309,7 +310,7 @@ func TestDeleteWorker_ProcessSyncWorkerError(t *testing.T) { syncError := errors.New("sync failed") mockSyncWorker.On("Process", mock.Anything, mockRepo, mock.Anything, mockProgress).Return(syncError) - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, nil) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "pull resources: sync failed") } @@ -378,7 +379,7 @@ func TestDeleteWorker_deleteFiles(t *testing.T) { } } - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) err := worker.deleteFiles(context.Background(), mockRepo, mockProgress, opts, tt.paths...) if tt.expectedError != "" { @@ -473,7 +474,7 @@ func TestDeleteWorker_ProcessWithResourceRefs(t *testing.T) { return result.Path == "folders/test-folder.json" && result.Action == repository.FileActionDeleted && result.Error == nil })).Return() - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) @@ -531,7 +532,7 @@ func TestDeleteWorker_ProcessResourceRefsOnly(t *testing.T) { return result.Path == "dashboards/test-dashboard.json" && result.Action == repository.FileActionDeleted && result.Error == nil })).Return() - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) } @@ -596,7 +597,7 @@ func TestDeleteWorker_ProcessResourceResolutionError(t *testing.T) { return syncJob.Spec.Pull != nil && !syncJob.Spec.Pull.Incremental }), mockProgress).Return(nil) - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) // Should succeed even with resource resolution error } @@ -635,7 +636,7 @@ func TestDeleteWorker_ProcessResourcesFactoryError(t *testing.T) { mockProgress.On("StrictMaxErrors", 1).Return() mockProgress.On("SetMessage", mock.Anything, "Resolving resource paths").Return() - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "delete files from repository: create repository resources client: failed to create repository resources client") } @@ -671,7 +672,7 @@ func TestDeleteWorker_ProcessResourceRefsNotReaderWriter(t *testing.T) { mockProgress.On("SetTotal", mock.Anything, 1).Return() mockProgress.On("StrictMaxErrors", 1).Return() - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "delete files from repository: delete job submitted targeting repository that is not a ReaderWriter") } @@ -724,7 +725,7 @@ func TestDeleteWorker_ProcessResourceResolutionTooManyErrors(t *testing.T) { })).Return() mockProgress.On("TooManyErrors").Return(errors.New("too many errors")) - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "delete files from repository: too many errors") } @@ -821,7 +822,7 @@ func TestDeleteWorker_ProcessMixedResourcesWithPartialFailure(t *testing.T) { return result.Path == "folders/valid-folder.json" && result.Error == nil })).Return() - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) // Should succeed overall, with only the failed resource recorded as error } @@ -919,7 +920,7 @@ func TestDeleteWorker_ProcessWithPathDeduplication(t *testing.T) { return result.Path == "dashboards/unique-dashboard.json" && result.Action == repository.FileActionDeleted && result.Error == nil })).Return() - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) @@ -1036,7 +1037,7 @@ func TestDeleteWorker_RefURLsSetWithRef(t *testing.T) { }, } - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepoWithURLs, job, mockProgress) require.NoError(t, err) @@ -1092,7 +1093,7 @@ func TestDeleteWorker_RefURLsNotSetWithoutRef(t *testing.T) { }, } - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepoWithURLs, job, mockProgress) require.NoError(t, err) @@ -1142,7 +1143,7 @@ func TestDeleteWorker_RefURLsNotSetForNonURLRepository(t *testing.T) { }, } - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) diff --git a/pkg/registry/apis/provisioning/jobs/export/worker.go b/pkg/registry/apis/provisioning/jobs/export/worker.go index e7f5404d6cc..f7def54ac1f 100644 --- a/pkg/registry/apis/provisioning/jobs/export/worker.go +++ b/pkg/registry/apis/provisioning/jobs/export/worker.go @@ -10,6 +10,7 @@ import ( "github.com/grafana/grafana/apps/provisioning/pkg/repository" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" + "github.com/prometheus/client_golang/prometheus" ) //go:generate mockery --name ExportFn --structname MockExportFn --inpackage --filename mock_export_fn.go --with-expecter @@ -23,6 +24,7 @@ type ExportWorker struct { repositoryResources resources.RepositoryResourcesFactory exportFn ExportFn wrapWithStageFn WrapWithStageFn + registry prometheus.Registerer } func NewExportWorker( @@ -30,12 +32,14 @@ func NewExportWorker( repositoryResources resources.RepositoryResourcesFactory, exportFn ExportFn, wrapWithStageFn WrapWithStageFn, + registry prometheus.Registerer, ) *ExportWorker { return &ExportWorker{ clientFactory: clientFactory, repositoryResources: repositoryResources, exportFn: exportFn, wrapWithStageFn: wrapWithStageFn, + registry: registry, } } diff --git a/pkg/registry/apis/provisioning/jobs/export/worker_test.go b/pkg/registry/apis/provisioning/jobs/export/worker_test.go index 21cb19d0bda..88b69691667 100644 --- a/pkg/registry/apis/provisioning/jobs/export/worker_test.go +++ b/pkg/registry/apis/provisioning/jobs/export/worker_test.go @@ -7,6 +7,7 @@ import ( "testing" "time" + "github.com/prometheus/client_golang/prometheus" mock "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -54,7 +55,7 @@ func TestExportWorker_IsSupported(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - r := NewExportWorker(nil, nil, nil, nil) + r := NewExportWorker(nil, nil, nil, nil, prometheus.DefaultRegisterer) got := r.IsSupported(context.Background(), tt.job) require.Equal(t, tt.want, got) }) @@ -68,7 +69,7 @@ func TestExportWorker_ProcessNoExportSettings(t *testing.T) { }, } - r := NewExportWorker(nil, nil, nil, nil) + r := NewExportWorker(nil, nil, nil, nil, prometheus.DefaultRegisterer) err := r.Process(context.Background(), nil, job, nil) require.EqualError(t, err, "missing export settings") } @@ -91,7 +92,7 @@ func TestExportWorker_ProcessWriteNotAllowed(t *testing.T) { }, }) - r := NewExportWorker(nil, nil, nil, nil) + r := NewExportWorker(nil, nil, nil, nil, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, nil) require.EqualError(t, err, "this repository is read only") } @@ -115,7 +116,7 @@ func TestExportWorker_ProcessBranchNotAllowedForLocal(t *testing.T) { }, }) - r := NewExportWorker(nil, nil, nil, nil) + r := NewExportWorker(nil, nil, nil, nil, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, nil) require.EqualError(t, err, "this repository does not support the branch workflow") } @@ -147,7 +148,7 @@ func TestExportWorker_ProcessFailedToCreateClients(t *testing.T) { return fn(repo, true) }) - r := NewExportWorker(mockClients, nil, nil, mockStageFn.Execute) + r := NewExportWorker(mockClients, nil, nil, mockStageFn.Execute, prometheus.DefaultRegisterer) mockProgress := jobs.NewMockJobProgressRecorder(t) err := r.Process(context.Background(), mockRepo, job, mockProgress) @@ -183,7 +184,7 @@ func TestExportWorker_ProcessNotReaderWriter(t *testing.T) { return fn(repo, true) }) - r := NewExportWorker(mockClients, nil, nil, mockStageFn.Execute) + r := NewExportWorker(mockClients, nil, nil, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "export job submitted targeting repository that is not a ReaderWriter") } @@ -219,7 +220,7 @@ func TestExportWorker_ProcessRepositoryResourcesError(t *testing.T) { mockStageFn.On("Execute", context.Background(), mockRepo, mock.Anything, mock.Anything).Return(func(ctx context.Context, repo repository.Repository, stageOpts repository.StageOptions, fn func(repository.Repository, bool) error) error { return fn(repo, true) }) - r := NewExportWorker(mockClients, mockRepoResources, nil, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, nil, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "create repository resource client: failed to create repository resources client") } @@ -270,7 +271,7 @@ func TestExportWorker_ProcessStageOptions(t *testing.T) { return fn(repo, true) }) - r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) } @@ -351,7 +352,7 @@ func TestExportWorker_ProcessStageOptionsWithBranch(t *testing.T) { return fn(repo, true) }) - r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) }) @@ -394,7 +395,7 @@ func TestExportWorker_ProcessExportFnError(t *testing.T) { return fn(repo, true) }) - r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "export failed") } @@ -422,7 +423,7 @@ func TestExportWorker_ProcessWrapWithStageFnError(t *testing.T) { mockStageFn := NewMockWrapWithStageFn(t) mockStageFn.On("Execute", mock.Anything, mockRepo, mock.Anything, mock.Anything).Return(errors.New("stage failed")) - r := NewExportWorker(nil, nil, nil, mockStageFn.Execute) + r := NewExportWorker(nil, nil, nil, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "stage failed") } @@ -448,7 +449,7 @@ func TestExportWorker_ProcessBranchNotAllowedForStageableRepositories(t *testing mockProgress := jobs.NewMockJobProgressRecorder(t) // No progress messages expected in current implementation - r := NewExportWorker(nil, nil, nil, nil) + r := NewExportWorker(nil, nil, nil, nil, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "this repository does not support the branch workflow") } @@ -499,7 +500,7 @@ func TestExportWorker_ProcessGitRepository(t *testing.T) { return fn(repo, true) }) - r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) } @@ -545,7 +546,7 @@ func TestExportWorker_ProcessGitRepositoryExportFnError(t *testing.T) { return fn(repo, true) }) - r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "export failed") } @@ -608,7 +609,7 @@ func TestExportWorker_RefURLsSetWithBranch(t *testing.T) { return fn(mockReaderWriter, true) }) - r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepoWithURLs, job, mockProgress) require.NoError(t, err) @@ -664,7 +665,7 @@ func TestExportWorker_RefURLsNotSetWithoutBranch(t *testing.T) { return fn(mockReaderWriter, true) }) - r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepoWithURLs, job, mockProgress) require.NoError(t, err) @@ -720,7 +721,7 @@ func TestExportWorker_RefURLsNotSetForNonURLRepository(t *testing.T) { return fn(mockReaderWriter, true) }) - r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute) + r := NewExportWorker(mockClients, mockRepoResources, mockExportFn.Execute, mockStageFn.Execute, prometheus.DefaultRegisterer) err := r.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) diff --git a/pkg/registry/apis/provisioning/jobs/move/worker.go b/pkg/registry/apis/provisioning/jobs/move/worker.go index 790630f101e..35c4e6b0342 100644 --- a/pkg/registry/apis/provisioning/jobs/move/worker.go +++ b/pkg/registry/apis/provisioning/jobs/move/worker.go @@ -14,19 +14,22 @@ import ( "github.com/grafana/grafana/apps/provisioning/pkg/safepath" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" + "github.com/prometheus/client_golang/prometheus" ) type Worker struct { syncWorker jobs.Worker wrapFn repository.WrapWithStageFn resourcesFactory resources.RepositoryResourcesFactory + registry prometheus.Registerer } -func NewWorker(syncWorker jobs.Worker, wrapFn repository.WrapWithStageFn, resourcesFactory resources.RepositoryResourcesFactory) *Worker { +func NewWorker(syncWorker jobs.Worker, wrapFn repository.WrapWithStageFn, resourcesFactory resources.RepositoryResourcesFactory, registry prometheus.Registerer) *Worker { return &Worker{ syncWorker: syncWorker, wrapFn: wrapFn, resourcesFactory: resourcesFactory, + registry: registry, } } diff --git a/pkg/registry/apis/provisioning/jobs/move/worker_test.go b/pkg/registry/apis/provisioning/jobs/move/worker_test.go index 3f9ea095d60..f8f82f20142 100644 --- a/pkg/registry/apis/provisioning/jobs/move/worker_test.go +++ b/pkg/registry/apis/provisioning/jobs/move/worker_test.go @@ -16,6 +16,7 @@ import ( "github.com/grafana/grafana/apps/provisioning/pkg/safepath" "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" + "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" ) @@ -75,7 +76,7 @@ func TestMoveWorker_IsSupported(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) result := worker.IsSupported(context.Background(), tt.job) require.Equal(t, tt.expected, result) }) @@ -89,7 +90,7 @@ func TestMoveWorker_ProcessMissingMoveSettings(t *testing.T) { }, } - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), nil, job, nil) require.EqualError(t, err, "missing move settings") } @@ -104,7 +105,7 @@ func TestMoveWorker_ProcessMissingTargetPath(t *testing.T) { }, } - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), nil, job, nil) require.EqualError(t, err, "target path is required for move operation") } @@ -120,7 +121,7 @@ func TestMoveWorker_ProcessInvalidTargetPath(t *testing.T) { }, } - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), nil, job, nil) require.EqualError(t, err, "target path must be a directory (should end with '/')") } @@ -152,7 +153,7 @@ func TestMoveWorker_ProcessNotReaderWriter(t *testing.T) { mockProgress.On("SetTotal", mock.Anything, 1).Return() mockProgress.On("StrictMaxErrors", 1).Return() - worker := NewWorker(nil, mockWrapFn.Execute, nil) + worker := NewWorker(nil, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "move files in repository: move job submitted targeting repository that is not a ReaderWriter") } @@ -176,7 +177,7 @@ func TestMoveWorker_ProcessWrapFnError(t *testing.T) { mockProgress.On("SetTotal", mock.Anything, 1).Return() mockProgress.On("StrictMaxErrors", 1).Return() - worker := NewWorker(nil, mockWrapFn.Execute, nil) + worker := NewWorker(nil, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "move files in repository: stage failed") } @@ -223,7 +224,7 @@ func TestMoveWorker_ProcessMoveFilesSuccess(t *testing.T) { return result.Path == "test/path2" && result.Action == repository.FileActionRenamed && result.Error == nil })).Return() - worker := NewWorker(nil, mockWrapFn.Execute, nil) + worker := NewWorker(nil, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) } @@ -262,7 +263,7 @@ func TestMoveWorker_ProcessMoveFilesWithError(t *testing.T) { })).Return() mockProgress.On("TooManyErrors").Return(errors.New("too many errors")) - worker := NewWorker(nil, mockWrapFn.Execute, nil) + worker := NewWorker(nil, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "move files in repository: too many errors") } @@ -307,7 +308,7 @@ func TestMoveWorker_ProcessWithSyncWorker(t *testing.T) { return syncJob.Spec.Pull != nil && !syncJob.Spec.Pull.Incremental }), mockProgress).Return(nil) - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, nil) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) } @@ -348,7 +349,7 @@ func TestMoveWorker_ProcessSyncWorkerError(t *testing.T) { syncError := errors.New("sync failed") mockSyncWorker.On("Process", mock.Anything, mockRepo, mock.Anything, mockProgress).Return(syncError) - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, nil) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "pull resources: sync failed") } @@ -429,7 +430,7 @@ func TestMoveWorker_moveFiles(t *testing.T) { } } - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) err := worker.moveFiles(context.Background(), mockRepo, mockProgress, opts, tt.paths...) if tt.expectedError != "" { @@ -485,7 +486,7 @@ func TestMoveWorker_constructTargetPath(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - worker := NewWorker(nil, nil, nil) + worker := NewWorker(nil, nil, nil, prometheus.DefaultRegisterer) result := worker.constructTargetPath(tt.jobTargetPath, tt.sourcePath) require.Equal(t, tt.expectedTarget, result) }) @@ -554,7 +555,7 @@ func TestMoveWorker_ProcessWithResourceReferences(t *testing.T) { return syncJob.Spec.Pull != nil && !syncJob.Spec.Pull.Incremental }), mockProgress).Return(nil) - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) } @@ -615,7 +616,7 @@ func TestMoveWorker_ProcessResourceReferencesError(t *testing.T) { return syncJob.Spec.Pull != nil && !syncJob.Spec.Pull.Incremental }), mockProgress).Return(nil) - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) // Should continue despite individual resource errors } @@ -655,7 +656,7 @@ func TestMoveWorker_ProcessResourcesFactoryError(t *testing.T) { factoryError := errors.New("failed to create resources client") mockResourcesFactory.On("Client", mock.Anything, mockRepo).Return(nil, factoryError) - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.EqualError(t, err, "move files in repository: create repository resources client: failed to create resources client") } @@ -773,7 +774,7 @@ func TestMoveWorker_resolveResourcesToPaths(t *testing.T) { } } - worker := NewWorker(nil, nil, mockResourcesFactory) + worker := NewWorker(nil, nil, mockResourcesFactory, prometheus.DefaultRegisterer) paths, err := worker.resolveResourcesToPaths(context.Background(), mockRepo, mockProgress, tt.resources) if tt.expectedError != "" { @@ -885,7 +886,7 @@ func TestMoveWorker_RefURLsSetWithRef(t *testing.T) { }, } - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepoWithURLs, job, mockProgress) require.NoError(t, err) @@ -942,7 +943,7 @@ func TestMoveWorker_RefURLsNotSetWithoutRef(t *testing.T) { }, } - worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(mockSyncWorker, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepoWithURLs, job, mockProgress) require.NoError(t, err) @@ -993,7 +994,7 @@ func TestMoveWorker_RefURLsNotSetForNonURLRepository(t *testing.T) { }, } - worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory) + worker := NewWorker(nil, mockWrapFn.Execute, mockResourcesFactory, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), mockRepo, job, mockProgress) require.NoError(t, err) diff --git a/pkg/registry/apis/provisioning/jobs/persistentstore.go b/pkg/registry/apis/provisioning/jobs/persistentstore.go index 5604e7bff0c..cd5b5c1561b 100644 --- a/pkg/registry/apis/provisioning/jobs/persistentstore.go +++ b/pkg/registry/apis/provisioning/jobs/persistentstore.go @@ -18,6 +18,7 @@ import ( provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1" client "github.com/grafana/grafana/apps/provisioning/pkg/generated/clientset/versioned/typed/provisioning/v0alpha1" "github.com/grafana/grafana/pkg/apimachinery/identity" + "github.com/prometheus/client_golang/prometheus" ) const ( @@ -72,18 +73,21 @@ type persistentStore struct { // expiry is the time after which a job is considered abandoned. // If a job is abandoned, it will have its claim cleaned up periodically. expiry time.Duration + + registry prometheus.Registerer } // NewJobStore creates a new job queue implementation using the API client. -func NewJobStore(provisioningClient client.ProvisioningV0alpha1Interface, expiry time.Duration) (*persistentStore, error) { +func NewJobStore(provisioningClient client.ProvisioningV0alpha1Interface, expiry time.Duration, registry prometheus.Registerer) (*persistentStore, error) { if expiry <= 0 { expiry = time.Second * 30 } return &persistentStore{ - client: provisioningClient, - clock: time.Now, - expiry: expiry, + client: provisioningClient, + clock: time.Now, + expiry: expiry, + registry: registry, }, nil } diff --git a/pkg/registry/apis/provisioning/jobs/sync/worker.go b/pkg/registry/apis/provisioning/jobs/sync/worker.go index 87ef447ff7f..a2e7f7feb69 100644 --- a/pkg/registry/apis/provisioning/jobs/sync/worker.go +++ b/pkg/registry/apis/provisioning/jobs/sync/worker.go @@ -10,6 +10,7 @@ import ( "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" + "github.com/prometheus/client_golang/prometheus" ) //go:generate mockery --name RepositoryPatchFn --structname MockRepositoryPatchFn --inpackage --filename repository_patch_fn_mock.go --with-expecter @@ -32,6 +33,9 @@ type SyncWorker struct { // Sync functions syncer Syncer + + // Registry for metrics + registry prometheus.Registerer } func NewSyncWorker( @@ -40,6 +44,7 @@ func NewSyncWorker( storageStatus dualwrite.Service, patchStatus RepositoryPatchFn, syncer Syncer, + registry prometheus.Registerer, ) *SyncWorker { return &SyncWorker{ clients: clients, @@ -47,6 +52,7 @@ func NewSyncWorker( patchStatus: patchStatus, storageStatus: storageStatus, syncer: syncer, + registry: registry, } } diff --git a/pkg/registry/apis/provisioning/jobs/sync/worker_test.go b/pkg/registry/apis/provisioning/jobs/sync/worker_test.go index 901c577aab8..cce78138148 100644 --- a/pkg/registry/apis/provisioning/jobs/sync/worker_test.go +++ b/pkg/registry/apis/provisioning/jobs/sync/worker_test.go @@ -10,6 +10,7 @@ import ( "github.com/grafana/grafana/pkg/registry/apis/provisioning/jobs" "github.com/grafana/grafana/pkg/registry/apis/provisioning/resources" "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" + "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -43,7 +44,7 @@ func TestSyncWorker_IsSupported(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - worker := NewSyncWorker(nil, nil, nil, nil, nil) + worker := NewSyncWorker(nil, nil, nil, nil, nil, prometheus.DefaultRegisterer) result := worker.IsSupported(context.Background(), tt.job) require.Equal(t, tt.expected, result) }) @@ -62,7 +63,7 @@ func TestSyncWorker_ProcessNotReaderWriter(t *testing.T) { }) fakeDualwrite := dualwrite.NewMockService(t) fakeDualwrite.On("ReadFromUnified", mock.Anything, mock.Anything).Return(true, nil).Twice() - worker := NewSyncWorker(nil, nil, fakeDualwrite, nil, nil) + worker := NewSyncWorker(nil, nil, fakeDualwrite, nil, nil, prometheus.DefaultRegisterer) err := worker.Process(context.Background(), repo, provisioning.Job{}, jobs.NewMockJobProgressRecorder(t)) require.EqualError(t, err, "sync job submitted for repository that does not support read-write -- this is a bug") } @@ -530,6 +531,7 @@ func TestSyncWorker_Process(t *testing.T) { dualwriteService, repositoryPatchFn.Execute, syncer, + prometheus.DefaultRegisterer, ) // Create test job diff --git a/pkg/registry/apis/provisioning/register.go b/pkg/registry/apis/provisioning/register.go index c0fdc90ca34..08a5c5cc034 100644 --- a/pkg/registry/apis/provisioning/register.go +++ b/pkg/registry/apis/provisioning/register.go @@ -117,6 +117,7 @@ type APIBuilder struct { extraWorkers []jobs.Worker restConfigGetter func(context.Context) (*clientrest.Config, error) + registry prometheus.Registerer } // NewAPIBuilder creates an API builder. @@ -139,6 +140,7 @@ func NewAPIBuilder( allowedTargets []provisioning.SyncTargetType, restConfigGetter func(context.Context) (*clientrest.Config, error), allowImageRendering bool, + registry prometheus.Registerer, newStandaloneClientFactoryFunc func(loopbackConfigProvider apiserver.RestConfigProvider) resources.ClientFactory, // optional, only used for standalone apiserver ) *APIBuilder { var clients resources.ClientFactory @@ -169,6 +171,7 @@ func NewAPIBuilder( allowedTargets: allowedTargets, restConfigGetter: restConfigGetter, allowImageRendering: allowImageRendering, + registry: registry, } for _, builder := range extraBuilders { @@ -254,6 +257,7 @@ func RegisterAPIService( allowedTargets, nil, // will use loopback instead cfg.ProvisioningAllowImageRendering, + reg, nil, ) apiregistration.RegisterAPI(builder) @@ -707,13 +711,13 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH b.client = c.ProvisioningV0alpha1() // Initialize the API client-based job store - b.jobs, err = jobs.NewJobStore(b.client, 30*time.Second) + b.jobs, err = jobs.NewJobStore(b.client, 30*time.Second, b.registry) if err != nil { return fmt.Errorf("create API client job store: %w", err) } - b.statusPatcher = appcontroller.NewRepositoryStatusPatcher(b.GetClient()) - b.healthChecker = controller.NewHealthChecker(b.statusPatcher) + b.statusPatcher = appcontroller.NewRepositoryStatusPatcher(b.GetClient(), b.registry) + b.healthChecker = controller.NewHealthChecker(b.statusPatcher, b.registry) // if running solely CRUD, skip the rest of the setup if b.onlyApiServer { @@ -739,6 +743,7 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH b.repositoryResources, export.ExportAll, stageIfPossible, + b.registry, ) syncer := sync.NewSyncer(sync.Compare, sync.FullSync, sync.IncrementalSync) @@ -748,6 +753,7 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH b.storageStatus, b.statusPatcher.Patch, syncer, + b.registry, ) signerFactory := signature.NewSignerFactory(b.clients) legacyResources := migrate.NewLegacyResourcesMigrator( @@ -779,8 +785,8 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH b.storageStatus, ) - deleteWorker := deletepkg.NewWorker(syncWorker, stageIfPossible, b.repositoryResources) - moveWorker := movepkg.NewWorker(syncWorker, stageIfPossible, b.repositoryResources) + deleteWorker := deletepkg.NewWorker(syncWorker, stageIfPossible, b.repositoryResources, b.registry) + moveWorker := movepkg.NewWorker(syncWorker, stageIfPossible, b.repositoryResources, b.registry) workers := []jobs.Worker{ deleteWorker, exportWorker, @@ -815,6 +821,7 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH 30*time.Second, // Lease renewal interval b.jobs, repoGetter, jobHistoryWriter, jobController.InsertNotifications(), + b.registry, workers..., ) if err != nil { @@ -837,6 +844,7 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH b.storageStatus, b.GetHealthChecker(), b.statusPatcher, + b.registry, ) if err != nil { return err diff --git a/pkg/registry/apis/provisioning/webhooks/register.go b/pkg/registry/apis/provisioning/webhooks/register.go index ce357188bad..6dc55a5b0de 100644 --- a/pkg/registry/apis/provisioning/webhooks/register.go +++ b/pkg/registry/apis/provisioning/webhooks/register.go @@ -18,6 +18,7 @@ import ( "github.com/grafana/grafana/pkg/services/rendering" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/unified/resource" + "github.com/prometheus/client_golang/prometheus" ) // WebhookExtraBuilder is a function that returns an ExtraBuilder. @@ -63,6 +64,7 @@ func ProvideWebhooksWithImages( renderer rendering.Service, blobstore resource.ResourceClient, configProvider apiserver.RestConfigProvider, + registry prometheus.Registerer, ) *WebhookExtraBuilder { urlProvider := func(_ string) string { return cfg.AppURL @@ -82,6 +84,7 @@ func ProvideWebhooksWithImages( isPublic, b, screenshotRenderer, + registry, ) evaluator := pullrequest.NewEvaluator(screenshotRenderer, parsers, urlProvider) @@ -98,7 +101,7 @@ func ProvideWebhooksWithImages( } } -func ProvideWebhooks(provisioningURL string) *WebhookExtraBuilder { +func ProvideWebhooks(provisioningURL string, registry prometheus.Registerer) *WebhookExtraBuilder { urlProvider := func(_ string) string { return provisioningURL } @@ -110,7 +113,7 @@ func ProvideWebhooks(provisioningURL string) *WebhookExtraBuilder { urlProvider: urlProvider, ExtraBuilder: func(b *provisioningapis.APIBuilder) provisioningapis.Extra { screenshotRenderer := pullrequest.NewNoOpRenderer() - webhook := NewWebhookConnector(isPublic, b, screenshotRenderer) + webhook := NewWebhookConnector(isPublic, b, screenshotRenderer, registry) return NewWebhookExtra(webhook) }, diff --git a/pkg/registry/apis/provisioning/webhooks/webhook.go b/pkg/registry/apis/provisioning/webhooks/webhook.go index 0997818d692..321f08fdafb 100644 --- a/pkg/registry/apis/provisioning/webhooks/webhook.go +++ b/pkg/registry/apis/provisioning/webhooks/webhook.go @@ -19,6 +19,7 @@ import ( "github.com/grafana/grafana/pkg/apimachinery/identity" provisioningapis "github.com/grafana/grafana/pkg/registry/apis/provisioning" "github.com/grafana/grafana/pkg/registry/apis/provisioning/webhooks/pullrequest" + "github.com/prometheus/client_golang/prometheus" ) type WebhookRepository interface { @@ -34,6 +35,7 @@ type webhookConnector struct { webhooksEnabled bool core *provisioningapis.APIBuilder renderer pullrequest.ScreenshotRenderer + registry prometheus.Registerer } func NewWebhookConnector( @@ -41,11 +43,13 @@ func NewWebhookConnector( // TODO: use interface for this core *provisioningapis.APIBuilder, renderer pullrequest.ScreenshotRenderer, + registry prometheus.Registerer, ) *webhookConnector { return &webhookConnector{ webhooksEnabled: webhooksEnabled, core: core, renderer: renderer, + registry: registry, } } diff --git a/pkg/server/wire_gen.go b/pkg/server/wire_gen.go index f0f9de53f93..eba4ca46baf 100644 --- a/pkg/server/wire_gen.go +++ b/pkg/server/wire_gen.go @@ -827,7 +827,7 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api userStorageAPIBuilder := userstorage.RegisterAPIService(featureToggles, apiserverService, registerer) apiBuilder := preferences.RegisterAPIService(cfg, featureToggles, sqlStore, prefService, starService, userService, apiserverService) legacyMigrator := legacy.ProvideLegacyMigrator(sqlStore, provisioningServiceImpl, libraryPanelService, dashboardPermissionsService, accessControl, featureToggles) - webhookExtraBuilder := webhooks.ProvideWebhooksWithImages(cfg, renderingService, resourceClient, eventualRestConfigProvider) + webhookExtraBuilder := webhooks.ProvideWebhooksWithImages(cfg, renderingService, resourceClient, eventualRestConfigProvider, registerer) v3 := extras.ProvideProvisioningExtraAPIs(webhookExtraBuilder) pullRequestWorker := pullrequest.ProvidePullRequestWorker(cfg, renderingService, resourceClient, eventualRestConfigProvider) v4 := extras.ProvideExtraWorkers(pullRequestWorker) @@ -1432,7 +1432,7 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac userStorageAPIBuilder := userstorage.RegisterAPIService(featureToggles, apiserverService, registerer) apiBuilder := preferences.RegisterAPIService(cfg, featureToggles, sqlStore, prefService, starService, userService, apiserverService) legacyMigrator := legacy.ProvideLegacyMigrator(sqlStore, provisioningServiceImpl, libraryPanelService, dashboardPermissionsService, accessControl, featureToggles) - webhookExtraBuilder := webhooks.ProvideWebhooksWithImages(cfg, renderingService, resourceClient, eventualRestConfigProvider) + webhookExtraBuilder := webhooks.ProvideWebhooksWithImages(cfg, renderingService, resourceClient, eventualRestConfigProvider, registerer) v3 := extras.ProvideProvisioningExtraAPIs(webhookExtraBuilder) pullRequestWorker := pullrequest.ProvidePullRequestWorker(cfg, renderingService, resourceClient, eventualRestConfigProvider) v4 := extras.ProvideExtraWorkers(pullRequestWorker)