DualWrite: Manage values from KV store (not file) (#106772)

This commit is contained in:
Ryan McKinley
2025-06-18 10:37:44 +03:00
committed by GitHub
parent 3cb62e370b
commit 945bc53b4c
8 changed files with 143 additions and 111 deletions
+6 -6
View File
@@ -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()
+4 -2
View File
@@ -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,
+21 -80
View File
@@ -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
}
@@ -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))
}
@@ -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) {
+30 -19
View File
@@ -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)
}
@@ -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)
+17
View File
@@ -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