diff --git a/pkg/apiserver/rest/dualwriter.go b/pkg/apiserver/rest/dualwriter.go index 9f7e29492bb..072d26ee846 100644 --- a/pkg/apiserver/rest/dualwriter.go +++ b/pkg/apiserver/rest/dualwriter.go @@ -102,6 +102,7 @@ func SetDualWritingMode( ctx context.Context, kvs NamespacedKVStore, cfg *SyncerConfig, + metrics *DualWriterMetrics, ) (DualWriterMode, error) { if cfg == nil { return Mode0, errors.New("syncer config is nil") @@ -167,7 +168,7 @@ func SetDualWritingMode( // Before running the sync, set the syncer config to the current mode, as we have to run the syncer // once in the current active mode before we can upgrade. cfg.Mode = currentMode - syncOk, err := runDataSyncer(ctx, cfg) + syncOk, err := runDataSyncer(ctx, cfg, metrics) // Once we are done with running the syncer, we can change the mode back on the config to the desired one. cfg.Mode = cfgModeTmp if err != nil { diff --git a/pkg/apiserver/rest/dualwriter_syncer.go b/pkg/apiserver/rest/dualwriter_syncer.go index fb064370fc6..4931e4f73eb 100644 --- a/pkg/apiserver/rest/dualwriter_syncer.go +++ b/pkg/apiserver/rest/dualwriter_syncer.go @@ -6,7 +6,6 @@ import ( "math/rand" "time" - "github.com/prometheus/client_golang/prometheus" apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" metainternalversion "k8s.io/apimachinery/pkg/apis/meta/internalversion" @@ -40,8 +39,6 @@ type SyncerConfig struct { DataSyncerInterval time.Duration DataSyncerRecordsLimit int - - Reg prometheus.Registerer } func (s *SyncerConfig) Validate() error { @@ -69,15 +66,12 @@ func (s *SyncerConfig) Validate() error { if s.DataSyncerRecordsLimit == 0 { s.DataSyncerRecordsLimit = 1000 } - if s.Reg == nil { - s.Reg = prometheus.DefaultRegisterer - } return nil } // StartPeriodicDataSyncer starts a background job that will execute the DataSyncer, syncing the data // from the hosted grafana backend into the unified storage backend. This is run in the grafana instance. -func StartPeriodicDataSyncer(ctx context.Context, cfg *SyncerConfig) error { +func StartPeriodicDataSyncer(ctx context.Context, cfg *SyncerConfig, metrics *DualWriterMetrics) error { if err := cfg.Validate(); err != nil { return fmt.Errorf("invalid syncer config: %w", err) } @@ -95,14 +89,14 @@ func StartPeriodicDataSyncer(ctx context.Context, cfg *SyncerConfig) error { time.Sleep(time.Second * time.Duration(jitterSeconds)) // run it immediately - syncOK, err := runDataSyncer(ctx, cfg) + syncOK, err := runDataSyncer(ctx, cfg, metrics) log.Info("data syncer finished", "syncOK", syncOK, "error", err) ticker := time.NewTicker(cfg.DataSyncerInterval) for { select { case <-ticker.C: - syncOK, err = runDataSyncer(ctx, cfg) + syncOK, err = runDataSyncer(ctx, cfg, metrics) log.Info("data syncer finished", "syncOK", syncOK, ", error", err) case <-ctx.Done(): return @@ -114,7 +108,7 @@ func StartPeriodicDataSyncer(ctx context.Context, cfg *SyncerConfig) error { // runDataSyncer will ensure that data between legacy storage and unified storage are in sync. // The sync implementation depends on the DualWriter mode -func runDataSyncer(ctx context.Context, cfg *SyncerConfig) (bool, error) { +func runDataSyncer(ctx context.Context, cfg *SyncerConfig, metrics *DualWriterMetrics) (bool, error) { if err := cfg.Validate(); err != nil { return false, fmt.Errorf("invalid syncer config: %w", err) } @@ -126,19 +120,17 @@ func runDataSyncer(ctx context.Context, cfg *SyncerConfig) (bool, error) { // implementation depends on the current DualWriter mode switch cfg.Mode { case Mode1, Mode2: - return legacyToUnifiedStorageDataSyncer(ctx, cfg) + return legacyToUnifiedStorageDataSyncer(ctx, cfg, metrics) default: klog.Info("data syncer not implemented for mode:", cfg.Mode) return false, nil } } -func legacyToUnifiedStorageDataSyncer(ctx context.Context, cfg *SyncerConfig) (bool, error) { +func legacyToUnifiedStorageDataSyncer(ctx context.Context, cfg *SyncerConfig, metrics *DualWriterMetrics) (bool, error) { if err := cfg.Validate(); err != nil { return false, fmt.Errorf("invalid syncer config: %w", err) } - metrics := &dualWriterMetrics{} - metrics.init(cfg.Reg) log := klog.NewKlogr().WithName("legacyToUnifiedStorageDataSyncer").WithValues("mode", cfg.Mode, "resource", cfg.Kind) diff --git a/pkg/apiserver/rest/dualwriter_syncer_test.go b/pkg/apiserver/rest/dualwriter_syncer_test.go index c1e0a216e9f..9e2cb8a21b6 100644 --- a/pkg/apiserver/rest/dualwriter_syncer_test.go +++ b/pkg/apiserver/rest/dualwriter_syncer_test.go @@ -201,13 +201,12 @@ func TestLegacyToUnifiedStorage_DataSyncer(t *testing.T) { LegacyStorage: ls, Storage: us, Kind: "test.kind", - Reg: p, ServerLockService: &fakeServerLock{}, RequestInfo: &request.RequestInfo{}, DataSyncerRecordsLimit: 1000, DataSyncerInterval: time.Hour, - }) + }, NewDualWriterMetrics(nil)) if tt.wantErr { assert.Error(t, err) return @@ -241,13 +240,12 @@ func TestLegacyToUnifiedStorage_DataSyncer(t *testing.T) { LegacyStorage: ls, Storage: us, Kind: "test.kind", - Reg: p, ServerLockService: &fakeServerLock{}, RequestInfo: &request.RequestInfo{}, DataSyncerRecordsLimit: 1000, DataSyncerInterval: time.Hour, - }) + }, NewDualWriterMetrics(nil)) if tt.wantErr { assert.Error(t, err) return diff --git a/pkg/apiserver/rest/dualwriter_test.go b/pkg/apiserver/rest/dualwriter_test.go index b618382f3f6..011abcf26f3 100644 --- a/pkg/apiserver/rest/dualwriter_test.go +++ b/pkg/apiserver/rest/dualwriter_test.go @@ -105,11 +105,10 @@ func TestSetDualWritingMode(t *testing.T) { SkipDataSync: tt.skipDataSync, ServerLockService: serverLockSvc, RequestInfo: &request.RequestInfo{}, - Reg: p, DataSyncerRecordsLimit: 1000, DataSyncerInterval: time.Hour, - }) + }, NewDualWriterMetrics(nil)) require.NoError(t, err) require.Equal(t, tt.expectedMode, dwMode) diff --git a/pkg/apiserver/rest/metrics.go b/pkg/apiserver/rest/metrics.go index 9380e948f75..74892fa4e47 100644 --- a/pkg/apiserver/rest/metrics.go +++ b/pkg/apiserver/rest/metrics.go @@ -6,47 +6,40 @@ import ( "time" "github.com/prometheus/client_golang/prometheus" - "k8s.io/klog/v2" + "github.com/prometheus/client_golang/prometheus/promauto" ) -type dualWriterMetrics struct { - syncer *prometheus.HistogramVec +type DualWriterMetrics struct { + // DualWriterSyncerDuration is a metric summary for dual writer sync duration per mode + syncer *prometheus.HistogramVec + // DualWriterDataSyncerOutcome is a metric summary for dual writer data syncer outcome comparison between the 2 stores per mode syncerOutcome *prometheus.HistogramVec } -// DualWriterSyncerDuration is a metric summary for dual writer sync duration per mode -var DualWriterSyncerDuration = prometheus.NewHistogramVec(prometheus.HistogramOpts{ - Name: "dual_writer_data_syncer_duration_seconds", - Help: "Histogram for the runtime of dual writer data syncer duration per mode", - Namespace: "grafana", - NativeHistogramBucketFactor: 1.1, -}, []string{"is_error", "mode", "resource"}) +func NewDualWriterMetrics(reg prometheus.Registerer) *DualWriterMetrics { + return &DualWriterMetrics{ + syncer: promauto.With(reg).NewHistogramVec(prometheus.HistogramOpts{ + Name: "dual_writer_data_syncer_duration_seconds", + Help: "Histogram for the runtime of dual writer data syncer duration per mode", + Namespace: "grafana", + NativeHistogramBucketFactor: 1.1, + }, []string{"is_error", "mode", "resource"}), -// DualWriterDataSyncerOutcome is a metric summary for dual writer data syncer outcome comparison between the 2 stores per mode -var DualWriterDataSyncerOutcome = prometheus.NewHistogramVec(prometheus.HistogramOpts{ - Name: "dual_writer_data_syncer_outcome", - Help: "Histogram for the runtime of dual writer data syncer outcome comparison between the 2 stores per mode", - Namespace: "grafana", - NativeHistogramBucketFactor: 1.1, -}, []string{"mode", "resource"}) - -func (m *dualWriterMetrics) init(reg prometheus.Registerer) { - log := klog.NewKlogr() - m.syncer = DualWriterSyncerDuration - m.syncerOutcome = DualWriterDataSyncerOutcome - errSyncer := reg.Register(m.syncer) - errSyncerOutcome := reg.Register(m.syncerOutcome) - if errSyncer != nil || errSyncerOutcome != nil { - log.Info("cloud migration metrics already registered") + syncerOutcome: promauto.With(reg).NewHistogramVec(prometheus.HistogramOpts{ + Name: "dual_writer_data_syncer_outcome", + Help: "Histogram for the runtime of dual writer data syncer outcome comparison between the 2 stores per mode", + Namespace: "grafana", + NativeHistogramBucketFactor: 1.1, + }, []string{"mode", "resource"}), } } -func (m *dualWriterMetrics) recordDataSyncerDuration(isError bool, mode DualWriterMode, resource string, startFrom time.Time) { +func (m *DualWriterMetrics) recordDataSyncerDuration(isError bool, mode DualWriterMode, resource string, startFrom time.Time) { duration := time.Since(startFrom).Seconds() m.syncer.WithLabelValues(strconv.FormatBool(isError), fmt.Sprintf("%d", mode), resource).Observe(duration) } -func (m *dualWriterMetrics) recordDataSyncerOutcome(mode DualWriterMode, resource string, synced bool) { +func (m *DualWriterMetrics) recordDataSyncerOutcome(mode DualWriterMode, resource string, synced bool) { var observeValue float64 if !synced { observeValue = 1 diff --git a/pkg/apiserver/rest/storage_mocks_test.go b/pkg/apiserver/rest/storage_mocks_test.go index a85a372c5e1..a047211fe61 100644 --- a/pkg/apiserver/rest/storage_mocks_test.go +++ b/pkg/apiserver/rest/storage_mocks_test.go @@ -5,7 +5,6 @@ import ( "errors" "time" - "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/mock" metainternalversion "k8s.io/apimachinery/pkg/apis/meta/internalversion" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -21,8 +20,6 @@ var anotherObj = &example.Pod{TypeMeta: metav1.TypeMeta{Kind: "foo"}, ObjectMeta var exampleList = &example.PodList{TypeMeta: metav1.TypeMeta{Kind: "foo"}, ListMeta: metav1.ListMeta{}, Items: []example.Pod{*exampleObj}} var anotherList = &example.PodList{Items: []example.Pod{*anotherObj}} -var p = prometheus.NewRegistry() - type storageMock struct { *mock.Mock Storage diff --git a/pkg/services/apiserver/builder/helper.go b/pkg/services/apiserver/builder/helper.go index c33b5a062e7..6cbd65b0d9f 100644 --- a/pkg/services/apiserver/builder/helper.go +++ b/pkg/services/apiserver/builder/helper.go @@ -284,6 +284,7 @@ func InstallAPIs( // support the legacy storage type. var dualWrite grafanarest.DualWriteBuilder metrics := newBuilderMetrics(reg) + dualWriterMetrics := grafanarest.NewDualWriterMetrics(reg) // nolint:staticcheck if storageOpts.StorageType != options.StorageTypeLegacy { @@ -337,11 +338,10 @@ func InstallAPIs( ServerLockService: serverLock, DataSyncerInterval: dataSyncerInterval, DataSyncerRecordsLimit: dataSyncerRecordsLimit, - Reg: reg, } // This also sets the currentMode on the syncer config. - currentMode, err := grafanarest.SetDualWritingMode(ctx, kvStore, syncerCfg) + currentMode, err := grafanarest.SetDualWritingMode(ctx, kvStore, syncerCfg, dualWriterMetrics) if err != nil { return nil, err } @@ -359,7 +359,7 @@ func InstallAPIs( if dualWriterPeriodicDataSyncJobEnabled { // The mode might have changed in SetDualWritingMode, so apply current mode first. syncerCfg.Mode = currentMode - if err := grafanarest.StartPeriodicDataSyncer(ctx, syncerCfg); err != nil { + if err := grafanarest.StartPeriodicDataSyncer(ctx, syncerCfg, dualWriterMetrics); err != nil { return nil, err } }