From 7e8bbd2ec48891ee19bac0b2b0639614ae8d1002 Mon Sep 17 00:00:00 2001 From: Ryan McKinley Date: Tue, 9 Sep 2025 19:38:33 +0300 Subject: [PATCH] DualWrite: Avoid dynamic wrapper when mode5 is configured (#110823) --- .../dualwrite/dualwriter_mode1_test.go | 3 - pkg/storage/legacysql/dualwrite/filedb.go | 40 ----- .../legacysql/dualwrite/runtime_test.go | 6 +- pkg/storage/legacysql/dualwrite/service.go | 29 ++-- .../legacysql/dualwrite/service_test.go | 142 +++++++++++++----- 5 files changed, 127 insertions(+), 93 deletions(-) delete mode 100644 pkg/storage/legacysql/dualwrite/filedb.go diff --git a/pkg/storage/legacysql/dualwrite/dualwriter_mode1_test.go b/pkg/storage/legacysql/dualwrite/dualwriter_mode1_test.go index 0de214c4596..b864e3a360d 100644 --- a/pkg/storage/legacysql/dualwrite/dualwriter_mode1_test.go +++ b/pkg/storage/legacysql/dualwrite/dualwriter_mode1_test.go @@ -6,7 +6,6 @@ import ( "testing" "time" - "github.com/prometheus/client_golang/prometheus" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" "k8s.io/apimachinery/pkg/api/meta" @@ -27,8 +26,6 @@ var failingObj = &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() - func TestMode1_Create(t *testing.T) { type testCase struct { input runtime.Object diff --git a/pkg/storage/legacysql/dualwrite/filedb.go b/pkg/storage/legacysql/dualwrite/filedb.go deleted file mode 100644 index 88e5ad336e0..00000000000 --- a/pkg/storage/legacysql/dualwrite/filedb.go +++ /dev/null @@ -1,40 +0,0 @@ -package dualwrite - -import ( - "context" - "encoding/json" - "os" - "path/filepath" - - "github.com/grafana/grafana-app-sdk/logging" - "github.com/grafana/grafana/pkg/setting" -) - -// This format was used in early G12 provisioning config. It should be removed after 12.1 -// This migration will be called once, and will remove the file based option even if the input was invalid -func migrateFileDBTo(cfg *setting.Cfg, db *keyvalueDB) { - fpath := filepath.Join(cfg.DataPath, "dualwrite.json") - v, err := os.ReadFile(fpath) // nolint:gosec - if err != nil { - return // the file does not exist, so nothign required - } - logger := logging.DefaultLogger.With("logger", "dualwrite-migrator") - - old := make(map[string]StorageStatus) - err = json.Unmarshal(v, &old) - if err != nil { - logger.Warn("error loading dual write settings", "err", err) - } - - for _, v := range old { - err = db.set(context.Background(), v) - if err != nil { - logger.Warn("error migrating dual write value", "err", err) - } - } - - err = os.Remove(fpath) - if err != nil { - logger.Warn("error removing old dual write settings", "err", err) - } -} diff --git a/pkg/storage/legacysql/dualwrite/runtime_test.go b/pkg/storage/legacysql/dualwrite/runtime_test.go index 2e28ffe81b1..06393e39d09 100644 --- a/pkg/storage/legacysql/dualwrite/runtime_test.go +++ b/pkg/storage/legacysql/dualwrite/runtime_test.go @@ -77,7 +77,7 @@ func TestRuntime_Create(t *testing.T) { tt.setupStorageFn(us.Mock, tt.input) } - m := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, kvstore.NewFakeKVStore(), nil) + m := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), nil, kvstore.NewFakeKVStore(), nil) dw, err := m.NewStorage(kind, ls, us) require.NoError(t, err) @@ -149,7 +149,7 @@ func TestRuntime_Get(t *testing.T) { tt.setupStorageFn(us.Mock, name) } - m := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, kvstore.NewFakeKVStore(), nil) + m := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), nil, kvstore.NewFakeKVStore(), nil) dw, err := m.NewStorage(kind, ls, us) require.NoError(t, err) status, err := m.Status(context.Background(), kind) @@ -233,7 +233,7 @@ func TestRuntime_CreateWhileMigrating(t *testing.T) { } // Shared provider across all tests - dual := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, kvstore.NewFakeKVStore(), nil) + dual := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), nil, kvstore.NewFakeKVStore(), nil) for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { diff --git a/pkg/storage/legacysql/dualwrite/service.go b/pkg/storage/legacysql/dualwrite/service.go index 47414019f3a..689bfe56fa8 100644 --- a/pkg/storage/legacysql/dualwrite/service.go +++ b/pkg/storage/legacysql/dualwrite/service.go @@ -9,6 +9,7 @@ import ( "k8s.io/apimachinery/pkg/runtime/schema" "github.com/grafana/grafana-app-sdk/logging" + "github.com/grafana/grafana/pkg/apiserver/rest" "github.com/grafana/grafana/pkg/infra/kvstore" "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/setting" @@ -25,11 +26,26 @@ func ProvideService( features featuremgmt.FeatureToggles, reg prometheus.Registerer, kv kvstore.KVStore, - cfg *setting.Cfg) Service { + cfg *setting.Cfg, +) Service { enabled := features.IsEnabledGlobally(featuremgmt.FlagManagedDualWriter) || features.IsEnabledGlobally(featuremgmt.FlagProvisioning) // required for git provisioning - if !enabled && cfg != nil { - return &staticService{cfg} // fallback to using the dual write flags from cfg + + if cfg != nil { + // Avoid dynamic behavior when things are explicitly configured to mode5 + allMode5 := true + for _, gr := range []string{ + "dashboards.dashboard.grafana.app", + "folders.folder.grafana.app", + } { + if cfg.UnifiedStorage[gr].DualWriterMode != rest.Mode5 { + allMode5 = false + break + } + } + if allMode5 || !enabled { + return &staticService{cfg} // fallback to using the dual write flags from cfg + } } db := &keyvalueDB{ @@ -37,21 +53,14 @@ func ProvideService( logger: logging.DefaultLogger.With("logger", "dualwrite.kv"), } - // TODO: remove this after G12.1 - if cfg != nil { - migrateFileDBTo(cfg, db) - } - return &service{ db: db, - reg: reg, enabled: enabled, } } type service struct { db *keyvalueDB - reg prometheus.Registerer enabled bool } diff --git a/pkg/storage/legacysql/dualwrite/service_test.go b/pkg/storage/legacysql/dualwrite/service_test.go index 2f41a597ce0..5ac51584e08 100644 --- a/pkg/storage/legacysql/dualwrite/service_test.go +++ b/pkg/storage/legacysql/dualwrite/service_test.go @@ -8,51 +8,119 @@ import ( "github.com/stretchr/testify/require" "k8s.io/apimachinery/pkg/runtime/schema" + "github.com/grafana/grafana/pkg/apiserver/rest" "github.com/grafana/grafana/pkg/infra/kvstore" "github.com/grafana/grafana/pkg/services/featuremgmt" + "github.com/grafana/grafana/pkg/setting" ) func TestService(t *testing.T) { - ctx := context.Background() - mode := ProvideService(featuremgmt.WithFeatures(), nil, kvstore.NewFakeKVStore(), nil) + t.Run("dynamic", func(t *testing.T) { + ctx := context.Background() + mode := ProvideService(featuremgmt.WithFeatures(), nil, kvstore.NewFakeKVStore(), nil) - gr := schema.GroupResource{Group: "ggg", Resource: "rrr"} - status, err := mode.Status(ctx, gr) - require.NoError(t, err) - require.Equal(t, StorageStatus{ - Group: "ggg", - Resource: "rrr", - WriteLegacy: true, - WriteUnified: true, - ReadUnified: false, - Migrated: 0, - Migrating: 0, - Runtime: true, - UpdateKey: 1, - }, status, "should start with the right defaults") + gr := schema.GroupResource{Group: "ggg", Resource: "rrr"} + status, err := mode.Status(ctx, gr) + require.NoError(t, err) + require.Equal(t, StorageStatus{ + Group: "ggg", + Resource: "rrr", + WriteLegacy: true, + WriteUnified: true, + ReadUnified: false, + Migrated: 0, + Migrating: 0, + Runtime: true, + UpdateKey: 1, + }, status, "should start with the right defaults") - // Start migration - status, err = mode.StartMigration(ctx, gr, 1) - require.NoError(t, err) - require.Equal(t, status.UpdateKey, int64(2), "the key increased") - require.True(t, status.Migrating > 0, "migration is running") + // Start migration + status, err = mode.StartMigration(ctx, gr, 1) + require.NoError(t, err) + require.Equal(t, status.UpdateKey, int64(2), "the key increased") + require.True(t, status.Migrating > 0, "migration is running") - status.Migrated = time.Now().UnixMilli() - status.Migrating = 0 - status, err = mode.Update(ctx, status) - require.NoError(t, err) - require.Equal(t, status.UpdateKey, int64(3), "the key increased") - require.Equal(t, status.Migrating, int64(0), "done migrating") - require.True(t, status.Migrated > 0, "migration is running") + status.Migrated = time.Now().UnixMilli() + status.Migrating = 0 + status, err = mode.Update(ctx, status) + require.NoError(t, err) + require.Equal(t, status.UpdateKey, int64(3), "the key increased") + require.Equal(t, status.Migrating, int64(0), "done migrating") + require.True(t, status.Migrated > 0, "migration is running") - status.WriteUnified = false - status.ReadUnified = true - _, err = mode.Update(ctx, status) - require.Error(t, err) // must write unified if we read it + status.WriteUnified = false + status.ReadUnified = true + _, err = mode.Update(ctx, status) + require.Error(t, err) // must write unified if we read it - status.WriteUnified = false - status.ReadUnified = false - status.WriteLegacy = false - _, err = mode.Update(ctx, status) - require.Error(t, err) // must write something! + status.WriteUnified = false + status.ReadUnified = false + status.WriteLegacy = false + _, err = mode.Update(ctx, status) + require.Error(t, err) // must write something! + }) + + t.Run("static", func(t *testing.T) { + type testCase struct { + name string + flags featuremgmt.FeatureToggles + cfg setting.Cfg + + isStatic bool + foldersFromUnified bool + dashboardsFromUnified bool + } + + for _, tc := range []testCase{{ + name: "both mode5", + flags: featuremgmt.WithFeatures(featuremgmt.FlagProvisioning), + dashboardsFromUnified: true, + foldersFromUnified: true, + isStatic: true, + cfg: setting.Cfg{ + UnifiedStorage: map[string]setting.UnifiedStorageConfig{ + "dashboards.dashboard.grafana.app": { + DualWriterMode: rest.Mode5, + }, + "folders.folder.grafana.app": { + DualWriterMode: rest.Mode5, + }, + }, + }}, { + name: "dynamic", + flags: featuremgmt.WithFeatures(featuremgmt.FlagProvisioning), + isStatic: false, + cfg: setting.Cfg{ + UnifiedStorage: map[string]setting.UnifiedStorageConfig{ + "dashboards.dashboard.grafana.app": { + DualWriterMode: rest.Mode5, + }, + }, + }}, + } { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + svc := ProvideService(tc.flags, nil, kvstore.NewFakeKVStore(), &tc.cfg) + + _, isStatic := svc.(*staticService) + require.Equal(t, tc.isStatic, isStatic) + + if isStatic { + v, err := svc.ReadFromUnified(ctx, schema.GroupResource{ + Group: "dashboard.grafana.app", + Resource: "dashboards", + }) + require.NoError(t, err) + require.Equal(t, tc.dashboardsFromUnified, v, "XXX") + + v, err = svc.ReadFromUnified(ctx, schema.GroupResource{ + Group: "folder.grafana.app", + Resource: "folders", + }) + require.NoError(t, err) + require.Equal(t, tc.foldersFromUnified, v, "YYY") + } + }) + } + }) }