diff --git a/pkg/registry/apis/dashboard/search_test.go b/pkg/registry/apis/dashboard/search_test.go index f6dfdb03892..31203b57efd 100644 --- a/pkg/registry/apis/dashboard/search_test.go +++ b/pkg/registry/apis/dashboard/search_test.go @@ -35,7 +35,7 @@ func TestSearchFallback(t *testing.T) { "dashboards.dashboard.grafana.app": {DualWriterMode: rest.Mode0}, }, } - dual := dualwrite.ProvideService(featuremgmt.WithFeatures(), nil, cfg) + dual := dualwrite.ProvideStaticServiceForTests(cfg) searchHandler := NewSearchHandler(tracing.NewNoopTracerService(), dual, mockLegacyClient, mockClient, nil) rr := httptest.NewRecorder() @@ -62,7 +62,7 @@ func TestSearchFallback(t *testing.T) { "dashboards.dashboard.grafana.app": {DualWriterMode: rest.Mode1}, }, } - dual := dualwrite.ProvideService(featuremgmt.WithFeatures(), nil, cfg) + dual := dualwrite.ProvideStaticServiceForTests(cfg) searchHandler := NewSearchHandler(tracing.NewNoopTracerService(), dual, mockLegacyClient, mockClient, nil) rr := httptest.NewRecorder() @@ -89,7 +89,7 @@ func TestSearchFallback(t *testing.T) { "dashboards.dashboard.grafana.app": {DualWriterMode: rest.Mode2}, }, } - dual := dualwrite.ProvideService(featuremgmt.WithFeatures(), nil, cfg) + dual := dualwrite.ProvideStaticServiceForTests(cfg) searchHandler := NewSearchHandler(tracing.NewNoopTracerService(), dual, mockLegacyClient, mockClient, nil) rr := httptest.NewRecorder() @@ -116,7 +116,7 @@ func TestSearchFallback(t *testing.T) { "dashboards.dashboard.grafana.app": {DualWriterMode: rest.Mode3}, }, } - dual := dualwrite.ProvideService(featuremgmt.WithFeatures(), nil, cfg) + dual := dualwrite.ProvideStaticServiceForTests(cfg) searchHandler := NewSearchHandler(tracing.NewNoopTracerService(), dual, mockLegacyClient, mockClient, nil) rr := httptest.NewRecorder() @@ -143,7 +143,7 @@ func TestSearchFallback(t *testing.T) { "dashboards.dashboard.grafana.app": {DualWriterMode: rest.Mode4}, }, } - dual := dualwrite.ProvideService(featuremgmt.WithFeatures(), nil, cfg) + dual := dualwrite.ProvideStaticServiceForTests(cfg) searchHandler := NewSearchHandler(tracing.NewNoopTracerService(), dual, mockLegacyClient, mockClient, nil) rr := httptest.NewRecorder() @@ -170,7 +170,7 @@ func TestSearchFallback(t *testing.T) { "dashboards.dashboard.grafana.app": {DualWriterMode: rest.Mode5}, }, } - dual := dualwrite.ProvideService(featuremgmt.WithFeatures(), nil, cfg) + dual := dualwrite.ProvideStaticServiceForTests(cfg) searchHandler := NewSearchHandler(tracing.NewNoopTracerService(), dual, mockLegacyClient, mockClient, nil) rr := httptest.NewRecorder() diff --git a/pkg/services/apiserver/service.go b/pkg/services/apiserver/service.go index 2ac8391ebe5..d33d5eccf5e 100644 --- a/pkg/services/apiserver/service.go +++ b/pkg/services/apiserver/service.go @@ -6,6 +6,7 @@ import ( "net/http" "path" + "github.com/prometheus/client_golang/prometheus" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/runtime/serializer" @@ -47,7 +48,6 @@ import ( "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" "github.com/grafana/grafana/pkg/storage/unified/apistore" "github.com/grafana/grafana/pkg/storage/unified/resource" - "github.com/prometheus/client_golang/prometheus" ) var ( @@ -325,7 +325,9 @@ func (s *service) start(ctx context.Context) error { // Install the API group+version err = builder.InstallAPIs(s.scheme, s.codecs, server, serverConfig.RESTOptionsGetter, builders, o.StorageOptions, // Required for the dual writer initialization - s.metrics, request.GetNamespaceMapper(s.cfg), kvstore.WithNamespace(s.kvStore, 0, "storage.dualwriting"), + s.metrics, + request.GetNamespaceMapper(s.cfg), + kvstore.WithNamespace(s.kvStore, 0, "storage.dualwriting"), // NOTE: will be removed and replaced with the dual writer utility s.serverLockService, s.storageStatus, optsregister, diff --git a/pkg/storage/legacysql/dualwrite/filedb.go b/pkg/storage/legacysql/dualwrite/filedb.go index fba0ab212b7..88e5ad336e0 100644 --- a/pkg/storage/legacysql/dualwrite/filedb.go +++ b/pkg/storage/legacysql/dualwrite/filedb.go @@ -4,96 +4,37 @@ import ( "context" "encoding/json" "os" - "sync" - - "k8s.io/apimachinery/pkg/runtime/schema" + "path/filepath" "github.com/grafana/grafana-app-sdk/logging" + "github.com/grafana/grafana/pkg/setting" ) -// Simple file implementation -- useful while testing and not yet sure about the SQL structure! -// When a path exists, read/write it from disk; otherwise it is held in memory -type fileDB struct { - path string - changed int64 - db map[string]StorageStatus - mu sync.RWMutex - logger logging.Logger -} - -// File implementation while testing -- values are saved in the data directory -func newFileDB(path string) *fileDB { - return &fileDB{ - db: make(map[string]StorageStatus), - path: path, - logger: logging.DefaultLogger.With("logger", "fileDB"), +// 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") -func (m *fileDB) Get(ctx context.Context, gr schema.GroupResource) (StorageStatus, bool, error) { - m.mu.RLock() - defer m.mu.RUnlock() + old := make(map[string]StorageStatus) + err = json.Unmarshal(v, &old) + if err != nil { + logger.Warn("error loading dual write settings", "err", err) + } - info, err := os.Stat(m.path) - if err == nil && info.ModTime().UnixMilli() != m.changed { - v, err := os.ReadFile(m.path) - if err == nil { - err = json.Unmarshal(v, &m.db) - m.changed = info.ModTime().UnixMilli() - } + for _, v := range old { + err = db.set(context.Background(), v) if err != nil { - m.logger.Warn("error reading filedb", "err", err) - } - - changed := false - for k, v := range m.db { - // Must write to unified if we are reading unified - if v.ReadUnified && !v.WriteUnified { - v.WriteUnified = true - m.db[k] = v - changed = true - } - - // Make sure we are writing something! - if !v.WriteLegacy && !v.WriteUnified { - v.WriteLegacy = true - m.db[k] = v - changed = true - } - } - if changed { - err = m.save() - m.logger.Warn("error saving changes filedb", "err", err) + logger.Warn("error migrating dual write value", "err", err) } } - v, ok := m.db[gr.String()] - return v, ok, nil -} - -func (m *fileDB) Set(ctx context.Context, status StorageStatus) error { - m.mu.Lock() - defer m.mu.Unlock() - - gr := schema.GroupResource{ - Group: status.Group, - Resource: status.Resource, + err = os.Remove(fpath) + if err != nil { + logger.Warn("error removing old dual write settings", "err", err) } - m.db[gr.String()] = status - - return m.save() -} - -func (m *fileDB) save() error { - if m.path != "" { - data, err := json.MarshalIndent(m.db, "", " ") - if err != nil { - return err - } - err = os.WriteFile(m.path, data, 0644) - if err != nil { - return err - } - } - return nil } diff --git a/pkg/storage/legacysql/dualwrite/kvstore.go b/pkg/storage/legacysql/dualwrite/kvstore.go new file mode 100644 index 00000000000..3b3e4f2f344 --- /dev/null +++ b/pkg/storage/legacysql/dualwrite/kvstore.go @@ -0,0 +1,59 @@ +package dualwrite + +import ( + "context" + "encoding/json" + + "k8s.io/apimachinery/pkg/runtime/schema" + + "github.com/grafana/grafana-app-sdk/logging" + "github.com/grafana/grafana/pkg/infra/kvstore" +) + +type keyvalueDB struct { + db kvstore.KVStore + logger logging.Logger +} + +// The setting is for all orgs +const globalKVOrgID = 0 + +// NOTE: this will replace any usage of "storage.dualwriting" and that will be removed +const globalKVNamespace = "unified.dualwrite" + +func (m *keyvalueDB) get(ctx context.Context, gr schema.GroupResource) (status StorageStatus, ok bool, err error) { + val, ok, err := m.db.Get(ctx, globalKVOrgID, globalKVNamespace, gr.String()) + if err != nil { + return status, false, err + } + + save := !ok + if ok { + err = json.Unmarshal([]byte(val), &status) + if err != nil { + m.logger.Warn("error reading filedb", "err", err) + save = true + } + } + + if status.validate() || save { + err = m.set(ctx, status) // will be the default values + } + return status, ok, err +} + +func (m *keyvalueDB) set(ctx context.Context, status StorageStatus) error { + gr := schema.GroupResource{ + Group: status.Group, + Resource: status.Resource, + } + + _ = status.validate() + + data, err := json.Marshal(status) + if err != nil { + return err + } + + return m.db.Set(ctx, globalKVOrgID, globalKVNamespace, gr.String(), string(data)) +} diff --git a/pkg/storage/legacysql/dualwrite/runtime_test.go b/pkg/storage/legacysql/dualwrite/runtime_test.go index a4b5af5223c..2e28ffe81b1 100644 --- a/pkg/storage/legacysql/dualwrite/runtime_test.go +++ b/pkg/storage/legacysql/dualwrite/runtime_test.go @@ -12,6 +12,7 @@ import ( "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" ) @@ -76,7 +77,7 @@ func TestRuntime_Create(t *testing.T) { tt.setupStorageFn(us.Mock, tt.input) } - m := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, nil) + m := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, kvstore.NewFakeKVStore(), nil) dw, err := m.NewStorage(kind, ls, us) require.NoError(t, err) @@ -148,7 +149,7 @@ func TestRuntime_Get(t *testing.T) { tt.setupStorageFn(us.Mock, name) } - m := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, nil) + m := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, kvstore.NewFakeKVStore(), nil) dw, err := m.NewStorage(kind, ls, us) require.NoError(t, err) status, err := m.Status(context.Background(), kind) @@ -232,7 +233,7 @@ func TestRuntime_CreateWhileMigrating(t *testing.T) { } // Shared provider across all tests - dual := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, nil) + dual := ProvideService(featuremgmt.WithFeatures(featuremgmt.FlagManagedDualWriter), p, 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 3c00f276e70..47414019f3a 100644 --- a/pkg/storage/legacysql/dualwrite/service.go +++ b/pkg/storage/legacysql/dualwrite/service.go @@ -3,47 +3,58 @@ package dualwrite import ( "context" "fmt" - "path/filepath" "time" "github.com/prometheus/client_golang/prometheus" "k8s.io/apimachinery/pkg/runtime/schema" + "github.com/grafana/grafana-app-sdk/logging" + "github.com/grafana/grafana/pkg/infra/kvstore" "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/setting" ) -func ProvideService(features featuremgmt.FeatureToggles, reg prometheus.Registerer, cfg *setting.Cfg) Service { +func ProvideStaticServiceForTests(cfg *setting.Cfg) Service { + if cfg == nil { + cfg = &setting.Cfg{} + } + return &staticService{cfg} +} + +func ProvideService( + features featuremgmt.FeatureToggles, + reg prometheus.Registerer, + kv kvstore.KVStore, + 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 } - path := "" // storage path + db := &keyvalueDB{ + db: kv, + logger: logging.DefaultLogger.With("logger", "dualwrite.kv"), + } + + // TODO: remove this after G12.1 if cfg != nil { - path = filepath.Join(cfg.DataPath, "dualwrite.json") + migrateFileDBTo(cfg, db) } return &service{ - db: newFileDB(path), + db: db, reg: reg, enabled: enabled, } } type service struct { - db statusStorage + db *keyvalueDB reg prometheus.Registerer enabled bool } -// The storage interface has zero business logic and simply writes values to a database -type statusStorage interface { - Get(ctx context.Context, gr schema.GroupResource) (StorageStatus, bool, error) - Set(ctx context.Context, status StorageStatus) error -} - // Hardcoded list of resources that should be controlled by the database (eventually everything?) func (m *service) ShouldManage(gr schema.GroupResource) bool { if !m.enabled { @@ -59,13 +70,13 @@ func (m *service) ShouldManage(gr schema.GroupResource) bool { } func (m *service) ReadFromUnified(ctx context.Context, gr schema.GroupResource) (bool, error) { - v, ok, err := m.db.Get(ctx, gr) + v, ok, err := m.db.get(ctx, gr) return ok && v.ReadUnified, err } // Status implements Service. func (m *service) Status(ctx context.Context, gr schema.GroupResource) (StorageStatus, error) { - v, found, err := m.db.Get(ctx, gr) + v, found, err := m.db.get(ctx, gr) if err != nil { return v, err } @@ -81,7 +92,7 @@ func (m *service) Status(ctx context.Context, gr schema.GroupResource) (StorageS Runtime: true, // need to explicitly ask for not runtime UpdateKey: 1, } - err := m.db.Set(ctx, v) + err := m.db.set(ctx, v) return v, err } return v, nil @@ -90,7 +101,7 @@ func (m *service) Status(ctx context.Context, gr schema.GroupResource) (StorageS // StartMigration implements Service. func (m *service) StartMigration(ctx context.Context, gr schema.GroupResource, key int64) (StorageStatus, error) { now := time.Now().UnixMilli() - v, ok, err := m.db.Get(ctx, gr) + v, ok, err := m.db.get(ctx, gr) if err != nil { return v, err } @@ -120,13 +131,13 @@ func (m *service) StartMigration(ctx context.Context, gr schema.GroupResource, k UpdateKey: 1, } } - err = m.db.Set(ctx, v) + err = m.db.set(ctx, v) return v, err } // FinishMigration implements Service. func (m *service) Update(ctx context.Context, status StorageStatus) (StorageStatus, error) { - v, ok, err := m.db.Get(ctx, schema.GroupResource{Group: status.Group, Resource: status.Resource}) + v, ok, err := m.db.get(ctx, schema.GroupResource{Group: status.Group, Resource: status.Resource}) if err != nil { return v, err } @@ -151,5 +162,5 @@ func (m *service) Update(ctx context.Context, status StorageStatus) (StorageStat return v, fmt.Errorf("must write either legacy or unified") } status.UpdateKey++ - return status, m.db.Set(ctx, status) + return status, m.db.set(ctx, status) } diff --git a/pkg/storage/legacysql/dualwrite/service_test.go b/pkg/storage/legacysql/dualwrite/service_test.go index e2b5d9037c6..2f41a597ce0 100644 --- a/pkg/storage/legacysql/dualwrite/service_test.go +++ b/pkg/storage/legacysql/dualwrite/service_test.go @@ -8,12 +8,13 @@ import ( "github.com/stretchr/testify/require" "k8s.io/apimachinery/pkg/runtime/schema" + "github.com/grafana/grafana/pkg/infra/kvstore" "github.com/grafana/grafana/pkg/services/featuremgmt" ) func TestService(t *testing.T) { ctx := context.Background() - mode := ProvideService(featuremgmt.WithFeatures(), nil, nil) + mode := ProvideService(featuremgmt.WithFeatures(), nil, kvstore.NewFakeKVStore(), nil) gr := schema.GroupResource{Group: "ggg", Resource: "rrr"} status, err := mode.Status(ctx, gr) diff --git a/pkg/storage/legacysql/dualwrite/types.go b/pkg/storage/legacysql/dualwrite/types.go index ea3346f2c34..bd1f2fa4415 100644 --- a/pkg/storage/legacysql/dualwrite/types.go +++ b/pkg/storage/legacysql/dualwrite/types.go @@ -32,6 +32,23 @@ type StorageStatus struct { UpdateKey int64 `json:"update_key" xorm:"update_key"` } +func (status *StorageStatus) validate() bool { + changed := false + + // Must write to unified if we are reading unified + if status.ReadUnified && !status.WriteUnified { + status.WriteUnified = true + changed = true + } + + // Make sure we are writing somewhere + if !status.WriteLegacy && !status.WriteUnified { + status.WriteLegacy = true + changed = true + } + return changed +} + // Service is a service for managing the dual write storage // //go:generate mockery --name Service --structname MockService --inpackage --filename service_mock.go --with-expecter