Dashboard: Multi-version builder (#100305)

This commit is contained in:
Todd Treece
2025-02-21 06:50:29 -05:00
committed by GitHub
parent 7be1fd953a
commit 3992ac2ac1
31 changed files with 480 additions and 1621 deletions
+17 -25
View File
@@ -6,6 +6,7 @@
package apistore
import (
"bytes"
"context"
"errors"
"fmt"
@@ -47,7 +48,6 @@ var _ storage.Interface = (*Storage)(nil)
// Optional settings that apply to a single resource
type StorageOptions struct {
LargeObjectSupport LargeObjectSupport
InternalConversion func([]byte, runtime.Object) (runtime.Object, error)
RequireDeprecatedInternalID bool
}
@@ -148,9 +148,6 @@ func (s *Storage) Versioner() storage.Versioner {
}
func (s *Storage) convertToObject(data []byte, obj runtime.Object) (runtime.Object, error) {
if s.opts.InternalConversion != nil {
return s.opts.InternalConversion(data, obj)
}
obj, _, err := s.codec.Decode(data, nil, obj)
return obj, err
}
@@ -180,7 +177,7 @@ func (s *Storage) Create(ctx context.Context, key string, obj runtime.Object, ou
return resource.GetError(rsp.Error)
}
if err := copyModifiedObjectToDestination(obj, out); err != nil {
if _, err := s.convertToObject(req.Value, out); err != nil {
return err
}
@@ -294,7 +291,7 @@ func (s *Storage) Watch(ctx context.Context, key string, opts storage.ListOption
}
reporter := apierrors.NewClientErrorReporter(500, "WATCH", "")
decoder := newStreamDecoder(client, s.newFunc, predicate, s.codec, cancelWatch, s.opts.InternalConversion)
decoder := newStreamDecoder(client, s.newFunc, predicate, s.codec, cancelWatch)
return watch.NewStreamWatcher(decoder, reporter), nil
}
@@ -538,15 +535,22 @@ func (s *Storage) GuaranteedUpdate(
}
if unchanged {
if err := copyModifiedObjectToDestination(updatedObj, destination); err != nil {
var buf bytes.Buffer
if err = s.codec.Encode(updatedObj, &buf); err != nil {
return err
}
if _, err := s.convertToObject(buf.Bytes(), destination); err != nil {
return err
}
return nil
}
rv := int64(0)
var (
value []byte
rv int64
)
if created {
value, err := s.prepareObjectForStorage(ctx, updatedObj)
value, err = s.prepareObjectForStorage(ctx, updatedObj)
if err != nil {
return err
}
@@ -562,10 +566,11 @@ func (s *Storage) GuaranteedUpdate(
}
rv = rsp2.ResourceVersion
} else {
req.Value, err = s.prepareObjectForUpdate(ctx, updatedObj, existingObj)
value, err = s.prepareObjectForUpdate(ctx, updatedObj, existingObj)
if err != nil {
return err
}
req.Value = value
rsp2, err := s.store.Update(ctx, req)
if err != nil {
return resource.GetError(resource.AsErrorResult(err))
@@ -576,11 +581,11 @@ func (s *Storage) GuaranteedUpdate(
rv = rsp2.ResourceVersion
}
if err := s.versioner.UpdateObject(updatedObj, uint64(rv)); err != nil {
if _, err := s.convertToObject(value, destination); err != nil {
return err
}
if err := copyModifiedObjectToDestination(updatedObj, destination); err != nil {
if err := s.versioner.UpdateObject(destination, uint64(rv)); err != nil {
return err
}
@@ -634,16 +639,3 @@ func (s *Storage) validateMinimumResourceVersion(minimumResourceVersion string,
}
return nil
}
func copyModifiedObjectToDestination(updatedObj runtime.Object, destination runtime.Object) error {
u, err := conversion.EnforcePtr(updatedObj)
if err != nil {
return fmt.Errorf("unable to enforce updated object pointer: %w", err)
}
d, err := conversion.EnforcePtr(destination)
if err != nil {
return fmt.Errorf("unable to enforce destination pointer: %w", err)
}
d.Set(u)
return nil
}
+13 -19
View File
@@ -19,33 +19,27 @@ import (
)
type streamDecoder struct {
client resource.ResourceStore_WatchClient
newFunc func() runtime.Object
predicate storage.SelectionPredicate
codec runtime.Codec
cancelWatch context.CancelFunc
done sync.WaitGroup
internalConversion func([]byte, runtime.Object) (runtime.Object, error)
client resource.ResourceStore_WatchClient
newFunc func() runtime.Object
predicate storage.SelectionPredicate
codec runtime.Codec
cancelWatch context.CancelFunc
done sync.WaitGroup
}
func newStreamDecoder(client resource.ResourceStore_WatchClient, newFunc func() runtime.Object, predicate storage.SelectionPredicate, codec runtime.Codec, cancelWatch context.CancelFunc, internalConversion func([]byte, runtime.Object) (runtime.Object, error)) *streamDecoder {
func newStreamDecoder(client resource.ResourceStore_WatchClient, newFunc func() runtime.Object, predicate storage.SelectionPredicate, codec runtime.Codec, cancelWatch context.CancelFunc) *streamDecoder {
return &streamDecoder{
client: client,
newFunc: newFunc,
predicate: predicate,
codec: codec,
cancelWatch: cancelWatch,
internalConversion: internalConversion,
client: client,
newFunc: newFunc,
predicate: predicate,
codec: codec,
cancelWatch: cancelWatch,
}
}
func (d *streamDecoder) toObject(w *resource.WatchEvent_Resource) (runtime.Object, error) {
var obj runtime.Object
var err error
if d.internalConversion != nil {
obj, err = d.internalConversion(w.Value, d.newFunc())
} else {
obj, _, err = d.codec.Decode(w.Value, nil, d.newFunc())
}
obj, _, err = d.codec.Decode(w.Value, nil, d.newFunc())
if err == nil {
accessor, err := utils.MetaAccessor(obj)
if err != nil {
@@ -1,34 +0,0 @@
package apistore
import (
"testing"
"github.com/grafana/grafana/pkg/storage/unified/resource"
"github.com/stretchr/testify/require"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/runtime"
)
func TestStreamDecoder(t *testing.T) {
t.Run("toObject should handle internal conversion", func(t *testing.T) {
called := false
internalConversion := func(data []byte, obj runtime.Object) (runtime.Object, error) {
called = true
return obj, nil
}
decoder := &streamDecoder{
newFunc: func() runtime.Object { return &unstructured.Unstructured{} },
internalConversion: internalConversion,
}
event := &resource.WatchEvent_Resource{
Value: []byte("test"),
}
obj, err := decoder.toObject(event)
require.NoError(t, err)
require.NotNil(t, obj)
require.True(t, called, "internal conversion function should have been called")
})
}