Provisioning: Wire up prometheus (#111444)

This commit is contained in:
Stephanie Hingtgen
2025-09-22 09:54:50 -05:00
committed by GitHub
parent 04bc71fa6d
commit bd550d2f06
23 changed files with 169 additions and 105 deletions
+1 -1
View File
@@ -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
+6 -3
View File
@@ -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,
}
}
@@ -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 != "" {
+4 -3
View File
@@ -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(
+13 -9
View File
@@ -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
+8 -6
View File
@@ -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
}
@@ -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,
}
}
@@ -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)
@@ -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{
@@ -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
}
@@ -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,
}
}
@@ -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)
@@ -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,
}
}
@@ -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)
@@ -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,
}
}
@@ -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)
@@ -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
}
@@ -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,
}
}
@@ -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
+13 -5
View File
@@ -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
@@ -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)
},
@@ -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,
}
}
+2 -2
View File
@@ -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)