From 0ffc4c441b46756fb887ac8f08a6601cc9ebc4d8 Mon Sep 17 00:00:00 2001 From: Arati R <33031346+suntala@users.noreply.github.com> Date: Thu, 23 May 2024 00:12:46 +0200 Subject: [PATCH] Storage: Add mode reconciliation for modes 1 and 2 (#87919) * Add skeleton implementation for mode reconciliation between 1 and 2 * Track mode for each dual writer * Add test for setting dual writer * Include context when setting dual writing mode --------- Co-authored-by: Dan Cech --- .golangci.toml | 1 + pkg/apiserver/rest/dualwriter.go | 75 +++++++++++++++++++++++++- pkg/apiserver/rest/dualwriter_mode1.go | 5 ++ pkg/apiserver/rest/dualwriter_mode2.go | 5 ++ pkg/apiserver/rest/dualwriter_mode3.go | 5 ++ pkg/apiserver/rest/dualwriter_mode4.go | 5 ++ pkg/apiserver/rest/dualwriter_test.go | 61 +++++++++++++++++++++ pkg/registry/apis/playlist/register.go | 13 +++-- 8 files changed, 165 insertions(+), 5 deletions(-) create mode 100644 pkg/apiserver/rest/dualwriter_test.go diff --git a/.golangci.toml b/.golangci.toml index c32f1e7efc5..648b0dad93c 100644 --- a/.golangci.toml +++ b/.golangci.toml @@ -77,6 +77,7 @@ allow = [ "github.com/grafana/grafana/pkg/apiserver", "github.com/grafana/grafana/pkg/services/apiserver/utils", "github.com/grafana/grafana/pkg/services/featuremgmt", + "github.com/grafana/grafana/pkg/infra/kvstore", ] deny = [ { pkg = "github.com/grafana/grafana/pkg", desc = "apiserver is not allowed to import grafana core" } diff --git a/pkg/apiserver/rest/dualwriter.go b/pkg/apiserver/rest/dualwriter.go index b75398a0d6c..ff20cb1f4bd 100644 --- a/pkg/apiserver/rest/dualwriter.go +++ b/pkg/apiserver/rest/dualwriter.go @@ -2,10 +2,15 @@ package rest import ( "context" + "errors" + "fmt" + "github.com/grafana/grafana/pkg/infra/kvstore" + "github.com/grafana/grafana/pkg/services/featuremgmt" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apiserver/pkg/registry/rest" + "k8s.io/klog/v2" ) var ( @@ -69,12 +74,13 @@ type LegacyStorage interface { type DualWriter interface { Storage LegacyStorage + Mode() DualWriterMode } type DualWriterMode int const ( - Mode1 DualWriterMode = iota + Mode1 DualWriterMode = iota + 1 Mode2 Mode3 Mode4 @@ -117,3 +123,70 @@ func (u *updateWrapper) Preconditions() *metav1.Preconditions { func (u *updateWrapper) UpdatedObject(ctx context.Context, oldObj runtime.Object) (newObj runtime.Object, err error) { return u.updated, nil } + +func SetDualWritingMode( + ctx context.Context, + kvs *kvstore.NamespacedKVStore, + features featuremgmt.FeatureToggles, + entity string, + legacy LegacyStorage, + storage Storage, +) (DualWriter, error) { + toMode := map[string]DualWriterMode{ + "1": Mode1, + "2": Mode2, + "3": Mode3, + "4": Mode4, + } + errDualWriterSetCurrentMode := errors.New("failed to set current dual writing mode") + + // Use entity name as key + m, ok, err := kvs.Get(ctx, entity) + if err != nil { + return nil, errors.New("failed to fetch current dual writing mode") + } + + currentMode, valid := toMode[m] + + if !valid && ok { + // Only log if "ok" because initially all instances will have mode unset for playlists. + klog.Info("invalid dual writing mode for playlists mode:", m) + } + + if !valid || !ok { + // Default to mode 1 + currentMode = Mode1 + + err := kvs.Set(ctx, entity, fmt.Sprint(currentMode)) + if err != nil { + return nil, errDualWriterSetCurrentMode + } + } + + // Desired mode is 2 and current mode is 1 + if features.IsEnabledGlobally(featuremgmt.FlagDualWritePlaylistsMode2) && (currentMode == Mode1) { + // This is where we go through the different gates to allow the instance to migrate from mode 1 to mode 2. + // There are none between mode 1 and mode 2 + currentMode = Mode2 + + err := kvs.Set(ctx, entity, fmt.Sprint(currentMode)) + if err != nil { + return nil, errDualWriterSetCurrentMode + } + } + // #TODO enable this check when we have a flag/config for setting mode 1 as the desired mode + // if features.IsEnabledGlobally(featuremgmt.FlagDualWritePlaylistsMode1) && (currentMode == Mode2) { + // // This is where we go through the different gates to allow the instance to migrate from mode 2 to mode 1. + // // There are none between mode 1 and mode 2 + // currentMode = Mode1 + + // err := kvs.Set(ctx, entity, fmt.Sprint(currentMode)) + // if err != nil { + // return nil, errDualWriterSetCurrentMode + // } + // } + + // #TODO add support for other combinations of desired and current modes + + return NewDualWriter(currentMode, legacy, storage), nil +} diff --git a/pkg/apiserver/rest/dualwriter_mode1.go b/pkg/apiserver/rest/dualwriter_mode1.go index 3641039b77f..ede37691857 100644 --- a/pkg/apiserver/rest/dualwriter_mode1.go +++ b/pkg/apiserver/rest/dualwriter_mode1.go @@ -32,6 +32,11 @@ func NewDualWriterMode1(legacy LegacyStorage, storage Storage) *DualWriterMode1 return &DualWriterMode1{Legacy: legacy, Storage: storage, Log: klog.NewKlogr().WithName("DualWriterMode1"), dualWriterMetrics: metrics} } +// Mode returns the mode of the dual writer. +func (d *DualWriterMode1) Mode() DualWriterMode { + return Mode1 +} + // Create overrides the behavior of the generic DualWriter and writes only to LegacyStorage. func (d *DualWriterMode1) Create(ctx context.Context, obj runtime.Object, createValidation rest.ValidateObjectFunc, options *metav1.CreateOptions) (runtime.Object, error) { log := d.Log.WithValues("kind", options.Kind) diff --git a/pkg/apiserver/rest/dualwriter_mode2.go b/pkg/apiserver/rest/dualwriter_mode2.go index 84be2c08e97..413867f6ae1 100644 --- a/pkg/apiserver/rest/dualwriter_mode2.go +++ b/pkg/apiserver/rest/dualwriter_mode2.go @@ -30,6 +30,11 @@ func NewDualWriterMode2(legacy LegacyStorage, storage Storage) *DualWriterMode2 return &DualWriterMode2{Legacy: legacy, Storage: storage, Log: klog.NewKlogr().WithName("DualWriterMode2"), dualWriterMetrics: metrics} } +// Mode returns the mode of the dual writer. +func (d *DualWriterMode2) Mode() DualWriterMode { + return Mode2 +} + // Create overrides the behavior of the generic DualWriter and writes to LegacyStorage and Storage. func (d *DualWriterMode2) Create(ctx context.Context, obj runtime.Object, createValidation rest.ValidateObjectFunc, options *metav1.CreateOptions) (runtime.Object, error) { log := d.Log.WithValues("kind", options.Kind) diff --git a/pkg/apiserver/rest/dualwriter_mode3.go b/pkg/apiserver/rest/dualwriter_mode3.go index b0f2630143d..7295d99e93c 100644 --- a/pkg/apiserver/rest/dualwriter_mode3.go +++ b/pkg/apiserver/rest/dualwriter_mode3.go @@ -26,6 +26,11 @@ func NewDualWriterMode3(legacy LegacyStorage, storage Storage) *DualWriterMode3 return &DualWriterMode3{Legacy: legacy, Storage: storage, Log: klog.NewKlogr().WithName("DualWriterMode3"), dualWriterMetrics: metrics} } +// Mode returns the mode of the dual writer. +func (d *DualWriterMode3) Mode() DualWriterMode { + return Mode3 +} + // Create overrides the behavior of the generic DualWriter and writes to LegacyStorage and Storage. func (d *DualWriterMode3) Create(ctx context.Context, obj runtime.Object, createValidation rest.ValidateObjectFunc, options *metav1.CreateOptions) (runtime.Object, error) { log := klog.FromContext(ctx) diff --git a/pkg/apiserver/rest/dualwriter_mode4.go b/pkg/apiserver/rest/dualwriter_mode4.go index ef0847cfcaf..61269e91e63 100644 --- a/pkg/apiserver/rest/dualwriter_mode4.go +++ b/pkg/apiserver/rest/dualwriter_mode4.go @@ -25,6 +25,11 @@ func NewDualWriterMode4(legacy LegacyStorage, storage Storage) *DualWriterMode4 return &DualWriterMode4{Legacy: legacy, Storage: storage, Log: klog.NewKlogr().WithName("DualWriterMode4"), dualWriterMetrics: metrics} } +// Mode returns the mode of the dual writer. +func (d *DualWriterMode4) Mode() DualWriterMode { + return Mode4 +} + // #TODO remove all DualWriterMode4 methods once we remove the generic DualWriter implementation // Create overrides the behavior of the generic DualWriter and writes only to Storage. diff --git a/pkg/apiserver/rest/dualwriter_test.go b/pkg/apiserver/rest/dualwriter_test.go new file mode 100644 index 00000000000..14fc1421824 --- /dev/null +++ b/pkg/apiserver/rest/dualwriter_test.go @@ -0,0 +1,61 @@ +package rest + +import ( + "context" + "fmt" + "testing" + + "github.com/grafana/grafana/pkg/infra/kvstore" + "github.com/grafana/grafana/pkg/services/featuremgmt" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" +) + +func TestSetDualWritingMode(t *testing.T) { + type testCase struct { + name string + features []any + stackID string + expectedMode DualWriterMode + } + tests := + // #TODO add test cases for kv store failures. Requires adding support in kvstore test_utils.go + []testCase{ + { + name: "should return a mode 1 dual writer when no desired mode is set", + features: []any{}, + stackID: "stack-1", + expectedMode: Mode1, + }, + { + name: "should return a mode 2 dual writer when mode 2 is set as the desired mode", + features: []any{featuremgmt.FlagDualWritePlaylistsMode2}, + stackID: "stack-1", + expectedMode: Mode2, + }, + } + + for _, tt := range tests { + l := (LegacyStorage)(nil) + s := (Storage)(nil) + m := &mock.Mock{} + + ls := legacyStoreMock{m, l} + us := storageMock{m, s} + + f := featuremgmt.WithFeatures(tt.features...) + kvStore := kvstore.WithNamespace(kvstore.NewFakeKVStore(), 0, "storage.dualwriting."+tt.stackID) + + key := "playlist" + + dw, err := SetDualWritingMode(context.Background(), kvStore, f, key, ls, us) + assert.NoError(t, err) + assert.Equal(t, tt.expectedMode, dw.Mode()) + + // check kv store + val, ok, err := kvStore.Get(context.Background(), key) + assert.True(t, ok) + assert.NoError(t, err) + assert.Equal(t, val, fmt.Sprint(tt.expectedMode)) + } +} diff --git a/pkg/registry/apis/playlist/register.go b/pkg/registry/apis/playlist/register.go index 68623623e8a..98e52a22e7f 100644 --- a/pkg/registry/apis/playlist/register.go +++ b/pkg/registry/apis/playlist/register.go @@ -1,6 +1,7 @@ package playlist import ( + "context" "fmt" "time" @@ -17,6 +18,7 @@ import ( playlist "github.com/grafana/grafana/pkg/apis/playlist/v0alpha1" "github.com/grafana/grafana/pkg/apiserver/builder" grafanarest "github.com/grafana/grafana/pkg/apiserver/rest" + "github.com/grafana/grafana/pkg/infra/kvstore" "github.com/grafana/grafana/pkg/services/apiserver/endpoints/request" "github.com/grafana/grafana/pkg/services/apiserver/utils" "github.com/grafana/grafana/pkg/services/featuremgmt" @@ -32,18 +34,21 @@ type PlaylistAPIBuilder struct { namespacer request.NamespaceMapper gv schema.GroupVersion features featuremgmt.FeatureToggles + kvStore *kvstore.NamespacedKVStore } func RegisterAPIService(p playlistsvc.Service, apiregistration builder.APIRegistrar, cfg *setting.Cfg, features featuremgmt.FeatureToggles, + kvStore kvstore.KVStore, ) *PlaylistAPIBuilder { builder := &PlaylistAPIBuilder{ service: p, namespacer: request.GetNamespaceMapper(cfg), gv: playlist.PlaylistResourceInfo.GroupVersion(), features: features, + kvStore: kvstore.WithNamespace(kvStore, 0, "storage.dualwriting"), } apiregistration.RegisterAPI(builder) return builder @@ -123,11 +128,11 @@ func (b *PlaylistAPIBuilder) GetAPIGroupInfo( return nil, err } - mode := grafanarest.Mode1 - if b.features.IsEnabledGlobally(featuremgmt.FlagDualWritePlaylistsMode2) { - mode = grafanarest.Mode2 + dualWriter, err := grafanarest.SetDualWritingMode(context.Background(), b.kvStore, b.features, "playlist", legacyStore, store) + if err != nil { + return nil, err } - storage[resource.StoragePath()] = grafanarest.NewDualWriter(mode, legacyStore, store) + storage[resource.StoragePath()] = dualWriter } apiGroupInfo.VersionedResourcesStorageMap[playlist.VERSION] = storage