DualWrite: Avoid dynamic wrapper when mode5 is configured (#110823)

This commit is contained in:
Ryan McKinley
2025-09-09 19:38:33 +03:00
committed by GitHub
parent d26c6c112a
commit 7e8bbd2ec4
5 changed files with 127 additions and 93 deletions
@@ -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
-40
View File
@@ -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)
}
}
@@ -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) {
+19 -10
View File
@@ -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
}
+105 -37
View File
@@ -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")
}
})
}
})
}