DualWrite: Cleanup and centralize the dual write creation (#90013)
This commit is contained in:
@@ -11,6 +11,7 @@ import (
|
||||
"k8s.io/apimachinery/pkg/api/meta"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||||
"k8s.io/apiserver/pkg/registry/rest"
|
||||
"k8s.io/klog/v2"
|
||||
)
|
||||
@@ -25,6 +26,9 @@ var (
|
||||
_ rest.SingularNameProvider = (DualWriter)(nil)
|
||||
)
|
||||
|
||||
// Function that will create a dual writer
|
||||
type DualWriteBuilder func(gr schema.GroupResource, legacy LegacyStorage, storage Storage) (Storage, error)
|
||||
|
||||
// Storage is a storage implementation that satisfies the same interfaces as genericregistry.Store.
|
||||
type Storage interface {
|
||||
rest.Storage
|
||||
@@ -152,7 +156,12 @@ func SetDualWritingMode(
|
||||
entity string,
|
||||
desiredMode DualWriterMode,
|
||||
reg prometheus.Registerer,
|
||||
) (DualWriter, error) {
|
||||
) (DualWriterMode, error) {
|
||||
// Mode0 means no DualWriter
|
||||
if desiredMode == Mode0 {
|
||||
return Mode0, nil
|
||||
}
|
||||
|
||||
toMode := map[string]DualWriterMode{
|
||||
// It is not possible to initialize a mode 0 dual writer. Mode 0 represents
|
||||
// writing to legacy storage without `unifiedStorage` enabled.
|
||||
@@ -166,7 +175,7 @@ func SetDualWritingMode(
|
||||
// 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")
|
||||
return Mode0, errors.New("failed to fetch current dual writing mode")
|
||||
}
|
||||
|
||||
currentMode, valid := toMode[m]
|
||||
@@ -182,7 +191,7 @@ func SetDualWritingMode(
|
||||
|
||||
err := kvs.Set(ctx, entity, fmt.Sprint(currentMode))
|
||||
if err != nil {
|
||||
return nil, errDualWriterSetCurrentMode
|
||||
return Mode0, errDualWriterSetCurrentMode
|
||||
}
|
||||
}
|
||||
|
||||
@@ -194,7 +203,7 @@ func SetDualWritingMode(
|
||||
|
||||
err := kvs.Set(ctx, entity, fmt.Sprint(currentMode))
|
||||
if err != nil {
|
||||
return nil, errDualWriterSetCurrentMode
|
||||
return Mode0, errDualWriterSetCurrentMode
|
||||
}
|
||||
}
|
||||
if (desiredMode == Mode1) && (currentMode == Mode2) {
|
||||
@@ -204,13 +213,13 @@ func SetDualWritingMode(
|
||||
|
||||
err := kvs.Set(ctx, entity, fmt.Sprint(currentMode))
|
||||
if err != nil {
|
||||
return nil, errDualWriterSetCurrentMode
|
||||
return Mode0, errDualWriterSetCurrentMode
|
||||
}
|
||||
}
|
||||
|
||||
// #TODO add support for other combinations of desired and current modes
|
||||
|
||||
return NewDualWriter(currentMode, legacy, storage, reg), nil
|
||||
return currentMode, nil
|
||||
}
|
||||
|
||||
var defaultConverter = runtime.UnstructuredConverter(runtime.DefaultUnstructuredConverter)
|
||||
|
||||
@@ -46,9 +46,9 @@ func TestSetDualWritingMode(t *testing.T) {
|
||||
kvStore := &fakeNamespacedKV{data: make(map[string]string), namespace: "storage.dualwriting." + tt.stackID}
|
||||
|
||||
p := prometheus.NewRegistry()
|
||||
dw, err := SetDualWritingMode(context.Background(), kvStore, ls, us, "playlist.grafana.app/v0alpha1", tt.desiredMode, p)
|
||||
dwMode, err := SetDualWritingMode(context.Background(), kvStore, ls, us, "playlist.grafana.app/v0alpha1", tt.desiredMode, p)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, tt.expectedMode, dw.Mode())
|
||||
assert.Equal(t, tt.expectedMode, dwMode)
|
||||
|
||||
// check kv store
|
||||
val, ok, err := kvStore.Get(context.Background(), "playlist.grafana.app/v0alpha1")
|
||||
|
||||
Reference in New Issue
Block a user