From 2e2b5942c81758e42e41bfbbc4607ece56cde658 Mon Sep 17 00:00:00 2001 From: Ryan McKinley Date: Fri, 21 Mar 2025 11:45:25 +0300 Subject: [PATCH] K8s/Unified: Consolidate generation logic in apistore client (#102260) --- pkg/apimachinery/utils/meta.go | 9 ++ pkg/apiserver/registry/generic/strategy.go | 18 +-- .../registry/generic/strategy_test.go | 34 +--- pkg/registry/apis/dashboard/register.go | 1 + pkg/registry/apis/folders/register.go | 3 +- pkg/storage/unified/apistore/prepare.go | 38 ++++- pkg/storage/unified/apistore/prepare_test.go | 153 +++++++++++++++++- pkg/storage/unified/apistore/store.go | 6 + pkg/tests/apis/playlist/playlist_test.go | 31 +++- 9 files changed, 237 insertions(+), 56 deletions(-) diff --git a/pkg/apimachinery/utils/meta.go b/pkg/apimachinery/utils/meta.go index 5ae3014c8c3..67076f2e96a 100644 --- a/pkg/apimachinery/utils/meta.go +++ b/pkg/apimachinery/utils/meta.go @@ -86,6 +86,7 @@ type GrafanaMetaAccessor interface { GetMessage() string SetMessage(msg string) SetAnnotation(key string, val string) + GetAnnotation(key string) string SetBlob(v *BlobInfo) GetBlob() *BlobInfo @@ -192,6 +193,14 @@ func (m *grafanaMetaAccessor) SetAnnotation(key string, val string) { m.obj.SetAnnotations(anno) } +func (m *grafanaMetaAccessor) GetAnnotation(key string) string { + anno := m.obj.GetAnnotations() + if anno != nil { + return anno[key] + } + return "" +} + func (m *grafanaMetaAccessor) get(key string) string { return m.obj.GetAnnotations()[key] } diff --git a/pkg/apiserver/registry/generic/strategy.go b/pkg/apiserver/registry/generic/strategy.go index 6cb5dd918fd..5293335f1f3 100644 --- a/pkg/apiserver/registry/generic/strategy.go +++ b/pkg/apiserver/registry/generic/strategy.go @@ -3,8 +3,6 @@ package generic import ( "context" - "github.com/grafana/grafana/pkg/apimachinery/utils" - apiequality "k8s.io/apimachinery/pkg/api/equality" "k8s.io/apimachinery/pkg/api/meta" "k8s.io/apimachinery/pkg/fields" "k8s.io/apimachinery/pkg/labels" @@ -14,6 +12,8 @@ import ( "k8s.io/apiserver/pkg/storage" "k8s.io/apiserver/pkg/storage/names" "sigs.k8s.io/structured-merge-diff/v4/fieldpath" + + "github.com/grafana/grafana/pkg/apimachinery/utils" ) type genericStrategy struct { @@ -81,20 +81,6 @@ func (g *genericStrategy) PrepareForUpdate(ctx context.Context, obj, old runtime } else { _ = newMeta.SetStatus(status) } - - spec, err := newMeta.GetSpec() - if err != nil { - return - } - - oldSpec, err := oldMeta.GetSpec() - if err != nil { - return - } - - if !apiequality.Semantic.DeepEqual(spec, oldSpec) { - newMeta.SetGeneration(oldMeta.GetGeneration() + 1) - } } func (g *genericStrategy) Validate(ctx context.Context, obj runtime.Object) field.ErrorList { diff --git a/pkg/apiserver/registry/generic/strategy_test.go b/pkg/apiserver/registry/generic/strategy_test.go index 969332305a2..906713d18fe 100644 --- a/pkg/apiserver/registry/generic/strategy_test.go +++ b/pkg/apiserver/registry/generic/strategy_test.go @@ -4,12 +4,13 @@ import ( "context" "testing" - "github.com/grafana/grafana/pkg/apiserver/registry/generic" "github.com/stretchr/testify/require" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apiserver/pkg/apis/example" + + "github.com/grafana/grafana/pkg/apiserver/registry/generic" ) func TestPrepareForUpdate(t *testing.T) { @@ -69,37 +70,6 @@ func TestPrepareForUpdate(t *testing.T) { }, }, }, - { - name: "increment generation if spec changes", - newObj: &example.Pod{ - ObjectMeta: metav1.ObjectMeta{ - Name: "test", - Namespace: "default", - Generation: 1, - }, - Spec: example.PodSpec{ - NodeSelector: map[string]string{"foo": "baz"}, - }, - Status: example.PodStatus{ - Phase: example.PodPhase("Running"), - }, - }, - oldObj: oldObj.DeepCopy(), - expectedGen: 2, - expectedObj: &example.Pod{ - ObjectMeta: metav1.ObjectMeta{ - Name: "test", - Namespace: "default", - Generation: 2, - }, - Spec: example.PodSpec{ - NodeSelector: map[string]string{"foo": "baz"}, - }, - Status: example.PodStatus{ - Phase: example.PodPhase("Running"), - }, - }, - }, } for _, tc := range testCases { diff --git a/pkg/registry/apis/dashboard/register.go b/pkg/registry/apis/dashboard/register.go index 101d36cd2da..c8a7b6f9f55 100644 --- a/pkg/registry/apis/dashboard/register.go +++ b/pkg/registry/apis/dashboard/register.go @@ -184,6 +184,7 @@ func (b *DashboardsAPIBuilder) Validate(ctx context.Context, a admission.Attribu func (b *DashboardsAPIBuilder) UpdateAPIGroupInfo(apiGroupInfo *genericapiserver.APIGroupInfo, opts builder.APIGroupOptions) error { storageOpts := apistore.StorageOptions{ + EnableFolderSupport: true, RequireDeprecatedInternalID: true, } diff --git a/pkg/registry/apis/folders/register.go b/pkg/registry/apis/folders/register.go index ca3d492815d..df8b2431f2f 100644 --- a/pkg/registry/apis/folders/register.go +++ b/pkg/registry/apis/folders/register.go @@ -6,7 +6,6 @@ import ( "fmt" "strings" - authtypes "github.com/grafana/authlib/types" "github.com/prometheus/client_golang/prometheus" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" @@ -18,6 +17,7 @@ import ( common "k8s.io/kube-openapi/pkg/common" "k8s.io/kube-openapi/pkg/spec3" + authtypes "github.com/grafana/authlib/types" "github.com/grafana/grafana/pkg/apimachinery/identity" "github.com/grafana/grafana/pkg/apimachinery/utils" "github.com/grafana/grafana/pkg/apis/folder/v0alpha1" @@ -156,6 +156,7 @@ func (b *FolderAPIBuilder) UpdateAPIGroupInfo(apiGroupInfo *genericapiserver.API } opts.StorageOptions(resourceInfo.GroupResource(), apistore.StorageOptions{ + EnableFolderSupport: true, RequireDeprecatedInternalID: true}) folderStore := &folderStorage{ diff --git a/pkg/storage/unified/apistore/prepare.go b/pkg/storage/unified/apistore/prepare.go index c9677fd0763..e6d9458ebfb 100644 --- a/pkg/storage/unified/apistore/prepare.go +++ b/pkg/storage/unified/apistore/prepare.go @@ -9,13 +9,14 @@ import ( "time" "github.com/google/uuid" + apiequality "k8s.io/apimachinery/pkg/api/equality" + apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" "k8s.io/apiserver/pkg/storage" "k8s.io/klog/v2" authtypes "github.com/grafana/authlib/types" - "github.com/grafana/grafana/pkg/apimachinery/utils" "github.com/grafana/grafana/pkg/storage/unified/resource" ) @@ -57,6 +58,9 @@ func (s *Storage) prepareObjectForStorage(ctx context.Context, newObject runtime if obj.GetUID() == "" { obj.SetUID(types.UID(uuid.NewString())) } + if obj.GetFolder() != "" && !s.opts.EnableFolderSupport { + return nil, apierrors.NewBadRequest(fmt.Sprintf("folders are not supported for: %s", s.gr.String())) + } if s.opts.RequireDeprecatedInternalID { // nolint:staticcheck @@ -77,6 +81,7 @@ func (s *Storage) prepareObjectForStorage(ctx context.Context, newObject runtime obj.SetUpdatedBy("") obj.SetUpdatedTimestamp(nil) obj.SetCreatedBy(info.GetUID()) + obj.SetGeneration(1) // the first time we write var buf bytes.Buffer if err = s.codec.Encode(newObject, &buf); err != nil { @@ -131,8 +136,35 @@ func (s *Storage) prepareObjectForUpdate(ctx context.Context, updateObject runti obj.SetDeprecatedInternalID(previousInternalID) // nolint:staticcheck } - obj.SetUpdatedBy(info.GetUID()) - obj.SetUpdatedTimestampMillis(time.Now().UnixMilli()) + // Check if we should bump the generation + changed := obj.GetFolder() != previous.GetFolder() + if changed { + if !s.opts.EnableFolderSupport { + return nil, apierrors.NewBadRequest(fmt.Sprintf("folders are not supported for: %s", s.gr.String())) + } + // TODO: check that we can move the folder? + } else if obj.GetDeletionTimestamp() != nil && previous.GetDeletionTimestamp() == nil { + changed = true // bump generation when deleted + } else { + spec, e1 := obj.GetSpec() + oldSpec, e2 := previous.GetSpec() + if e1 == nil && e2 == nil { + if !apiequality.Semantic.DeepEqual(spec, oldSpec) { + changed = true + } + } + } + + // Mark the resource as changed + if changed { + obj.SetGeneration(previous.GetGeneration() + 1) + obj.SetUpdatedBy(info.GetUID()) + obj.SetUpdatedTimestampMillis(time.Now().UnixMilli()) + } else { + obj.SetGeneration(previous.GetGeneration()) + obj.SetAnnotation(utils.AnnoKeyUpdatedBy, previous.GetAnnotation(utils.AnnoKeyUpdatedBy)) + obj.SetAnnotation(utils.AnnoKeyUpdatedTimestamp, previous.GetAnnotation(utils.AnnoKeyUpdatedTimestamp)) + } var buf bytes.Buffer if err = s.codec.Encode(updateObject, &buf); err != nil { diff --git a/pkg/storage/unified/apistore/prepare_test.go b/pkg/storage/unified/apistore/prepare_test.go index 1bb9c06cd04..1c17a9806c2 100644 --- a/pkg/storage/unified/apistore/prepare_test.go +++ b/pkg/storage/unified/apistore/prepare_test.go @@ -9,6 +9,8 @@ import ( "github.com/stretchr/testify/require" "golang.org/x/exp/rand" "k8s.io/apimachinery/pkg/api/apitesting" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/serializer" "k8s.io/apiserver/pkg/storage" @@ -30,11 +32,14 @@ func TestPrepareObjectForStorage(t *testing.T) { codec: apitesting.TestCodec(codecs, v0alpha1.DashboardResourceInfo.GroupVersion()), snowflake: node, opts: StorageOptions{ - LargeObjectSupport: nil, + EnableFolderSupport: true, + LargeObjectSupport: nil, }, } - ctx := authtypes.WithAuthInfo(context.Background(), &identity.StaticRequester{UserID: 1, UserUID: "user-uid", Type: authtypes.TypeUser}) + ctx := authtypes.WithAuthInfo(context.Background(), + &identity.StaticRequester{UserID: 1, UserUID: "user-uid", Type: authtypes.TypeUser}, + ) t.Run("Error getting auth info from context", func(t *testing.T) { _, err := s.prepareObjectForStorage(context.Background(), nil) @@ -81,7 +86,7 @@ func TestPrepareObjectForStorage(t *testing.T) { require.Empty(t, updatedTS) }) - t.Run("Should keep repo info", func(t *testing.T) { + t.Run("Should keep manager info", func(t *testing.T) { dashboard := v0alpha1.Dashboard{} dashboard.Name = "test-name" obj := dashboard.DeepCopyObject() @@ -117,6 +122,65 @@ func TestPrepareObjectForStorage(t *testing.T) { require.Equal(t, s.TimestampMillis, now.UnixMilli()) }) + t.Run("Update should manage incrementing generation and metadata", func(t *testing.T) { + dashboard := v0alpha1.Dashboard{} + dashboard.Name = "test-name" + obj := dashboard.DeepCopyObject() + meta, err := utils.MetaAccessor(obj) + meta.SetFolder("aaa") + require.NoError(t, err) + + encodedData, err := s.prepareObjectForStorage(ctx, obj) + require.NoError(t, err) + + insertedObject, _, err := s.codec.Decode(encodedData, nil, &v0alpha1.Dashboard{}) + require.NoError(t, err) + meta, err = utils.MetaAccessor(insertedObject) + require.NoError(t, err) + require.Equal(t, int64(1), meta.GetGeneration()) + require.Equal(t, "user:user-uid", meta.GetCreatedBy()) + require.Equal(t, "", meta.GetUpdatedBy()) // empty + ts, err := meta.GetUpdatedTimestamp() + require.NoError(t, err) + require.Nil(t, ts) + + // Change the user... and only update metadata + ctx = authtypes.WithAuthInfo(context.Background(), + &identity.StaticRequester{UserID: 1, UserUID: "user2", Type: authtypes.TypeUser}, + ) + + // Change the status... but generation is the same + updatedObject := insertedObject.DeepCopyObject() + meta, err = utils.MetaAccessor(updatedObject) + require.NoError(t, err) + err = meta.SetStatus(v0alpha1.DashboardStatus{ + Conversion: &v0alpha1.DashboardConversionStatus{ + Failed: true, + Error: "test", + }, + }) + require.NoError(t, err) + meta.SetGeneration(123) // will be removed + + // Update status without changing generation or update metadata + _, err = s.prepareObjectForUpdate(ctx, updatedObject, insertedObject) + require.NoError(t, err) + require.Equal(t, "", meta.GetUpdatedBy()) + require.Equal(t, int64(1), meta.GetGeneration()) + + // Change the folder -- the generation should increase and the updatedBy metadata + dashboard2 := &v0alpha1.Dashboard{ObjectMeta: v1.ObjectMeta{ + Name: dashboard.Name, + }} // TODO... deep copy, See: https://github.com/grafana/grafana/pull/102258 + meta2, err := utils.MetaAccessor(dashboard2) + require.NoError(t, err) + meta2.SetFolder("xyz") // will bump generation + _, err = s.prepareObjectForUpdate(ctx, dashboard2, updatedObject) + require.NoError(t, err) + require.Equal(t, "user:user2", meta2.GetUpdatedBy()) + require.Equal(t, int64(2), meta2.GetGeneration()) + }) + s.opts.RequireDeprecatedInternalID = true t.Run("Should generate internal id", func(t *testing.T) { dashboard := v0alpha1.Dashboard{} @@ -149,4 +213,87 @@ func TestPrepareObjectForStorage(t *testing.T) { require.NoError(t, err) require.Equal(t, meta.GetDeprecatedInternalID(), int64(1)) // nolint:staticcheck }) + + t.Run("calculate generation", func(t *testing.T) { + dash := &v0alpha1.Dashboard{ + ObjectMeta: v1.ObjectMeta{ + Name: "test", + }, + Spec: v0alpha1.DashboardSpec{ + Object: map[string]interface{}{ + "hello": "world", + }, + }, + } + out := getPreparedObject(t, ctx, s, dash, nil) + require.Equal(t, int64(1), out.GetGeneration()) + require.NotEmpty(t, out.GetAnnotation(utils.AnnoKeyCreatedBy)) + require.Equal(t, "", out.GetAnnotation(utils.AnnoKeyUpdatedBy)) + require.Equal(t, "", out.GetAnnotation(utils.AnnoKeyUpdatedTimestamp)) + + t.Run("increment when the spec changes", func(t *testing.T) { + b := dash.DeepCopy() + b.Spec.Object["x"] = "y" + out = getPreparedObject(t, ctx, s, b, dash) + require.Equal(t, int64(2), out.GetGeneration()) + require.NotEmpty(t, out.GetAnnotation(utils.AnnoKeyUpdatedBy)) + require.NotEmpty(t, out.GetAnnotation(utils.AnnoKeyUpdatedTimestamp)) + }) + + t.Run("increment when the folder changes", func(t *testing.T) { + b := dash.DeepCopy() + b.Annotations = map[string]string{ + utils.AnnoKeyFolder: "abc", + } + out = getPreparedObject(t, ctx, s, b, dash) + require.Equal(t, int64(2), out.GetGeneration()) + }) + + t.Run("increment when deleted", func(t *testing.T) { + now := v1.Now() + b := dash.DeepCopy() + b.DeletionTimestamp = &now + out = getPreparedObject(t, ctx, s, b, dash) + require.Equal(t, int64(2), out.GetGeneration()) + }) + + t.Run("keep when status, labels, or annotations change", func(t *testing.T) { + b := dash.DeepCopy() + b.Annotations = map[string]string{ + "x": "hello", + } + b.Labels = map[string]string{ + "a": "b", + } + b.Status = v0alpha1.DashboardStatus{ + Conversion: &v0alpha1.DashboardConversionStatus{ + Failed: true, + }, + } + out = getPreparedObject(t, ctx, s, b, dash) + require.Equal(t, int64(1), out.GetGeneration()) // still 1 + }) + }) +} + +func getPreparedObject(t *testing.T, ctx context.Context, s *Storage, obj runtime.Object, old runtime.Object) utils.GrafanaMetaAccessor { + t.Helper() + + var raw []byte + var err error + + if old == nil { + raw, err = s.prepareObjectForStorage(ctx, obj) + } else { + raw, err = s.prepareObjectForUpdate(ctx, obj, old) + } + require.NoError(t, err) + + out := &unstructured.Unstructured{} + err = out.UnmarshalJSON(raw) + require.NoError(t, err) + + meta, err := utils.MetaAccessor(out) + require.NoError(t, err) + return meta } diff --git a/pkg/storage/unified/apistore/store.go b/pkg/storage/unified/apistore/store.go index 966ba119c0a..4557c0a5d39 100644 --- a/pkg/storage/unified/apistore/store.go +++ b/pkg/storage/unified/apistore/store.go @@ -47,8 +47,14 @@ var _ storage.Interface = (*Storage)(nil) // Optional settings that apply to a single resource type StorageOptions struct { + // ????: should we constrain this to only dashboards for now? + // Not yet clear if this is a good general solution, or just a stop-gap LargeObjectSupport LargeObjectSupport + // Allow writing objects with metadata.annotations[grafana.app/folder] + EnableFolderSupport bool + + // Add internalID label when missing RequireDeprecatedInternalID bool } diff --git a/pkg/tests/apis/playlist/playlist_test.go b/pkg/tests/apis/playlist/playlist_test.go index a7b61abec50..a30aeba2fc2 100644 --- a/pkg/tests/apis/playlist/playlist_test.go +++ b/pkg/tests/apis/playlist/playlist_test.go @@ -10,11 +10,13 @@ import ( "testing" "github.com/stretchr/testify/require" + apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" + "github.com/grafana/grafana/pkg/apimachinery/utils" grafanarest "github.com/grafana/grafana/pkg/apiserver/rest" "github.com/grafana/grafana/pkg/services/apiserver/options" "github.com/grafana/grafana/pkg/services/featuremgmt" @@ -156,7 +158,7 @@ func TestIntegrationPlaylist(t *testing.T) { }) t.Run("with dual write (file, mode 5)", func(t *testing.T) { - doPlaylistTests(t, apis.NewK8sTestHelper(t, testinfra.GrafanaOpts{ + helper := doPlaylistTests(t, apis.NewK8sTestHelper(t, testinfra.GrafanaOpts{ AppModeProduction: true, DisableAnonymous: true, APIServerStorageType: "file", // write the files to disk @@ -169,6 +171,33 @@ func TestIntegrationPlaylist(t *testing.T) { featuremgmt.FlagKubernetesPlaylists, // Required so that legacy calls are also written }, })) + + client := helper.GetResourceClient(apis.ResourceClientArgs{ + User: helper.Org1.Editor, + GVR: gvr, + }) + + // Folder support needs to be enabled explicitly for this resource + t.Run("ensure writing folders is an error", func(t *testing.T) { + // Create works without folder + obj := helper.LoadYAMLOrJSONFile("testdata/playlist-generate.yaml") + out, err := client.Resource.Create(context.Background(), obj, metav1.CreateOptions{}) + require.NoError(t, err) + + meta, err := utils.MetaAccessor(out) + require.NoError(t, err) + require.Equal(t, int64(1), meta.GetGeneration()) + require.Equal(t, helper.Org1.Editor.Identity.GetUID(), meta.GetCreatedBy()) + require.Equal(t, "", meta.GetUpdatedBy()) + + meta, err = utils.MetaAccessor(obj) + require.NoError(t, err) + meta.SetFolder("FolderUID") + + _, err = client.Resource.Create(context.Background(), obj, metav1.CreateOptions{}) + require.Error(t, err) + require.True(t, apierrors.IsBadRequest(err)) + }) }) t.Run("with dual write (unified storage, mode 0)", func(t *testing.T) {