feat(unified-storage): add tracing to apistore (#113714)

This commit is contained in:
Jean-Philippe Quéméner
2025-11-12 09:48:56 +01:00
committed by GitHub
parent 04902b0419
commit 76ab09a6a2
+65 -27
View File
@@ -18,6 +18,7 @@ import (
"time"
"github.com/bwmarrin/snowflake"
"go.opentelemetry.io/otel"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
metaV1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -46,7 +47,10 @@ const (
LargeObjectSupportDisabled = false
)
var _ storage.Interface = (*Storage)(nil)
var (
_ storage.Interface = (*Storage)(nil)
tracer = otel.Tracer("github.com/grafana/grafana/pkg/storage/unified/apistore")
)
type DefaultPermissionSetter = func(ctx context.Context, key *resourcepb.ResourceKey, id authtypes.AuthInfo, obj utils.GrafanaMetaAccessor) error
@@ -200,7 +204,9 @@ func (s *Storage) Versioner() storage.Versioner {
return s.versioner
}
func (s *Storage) convertToObject(data []byte, obj runtime.Object) (runtime.Object, error) {
func (s *Storage) convertToObject(ctx context.Context, data []byte, obj runtime.Object) (runtime.Object, error) {
_, span := tracer.Start(ctx, "apistore.Storage.convertToObject")
defer span.End()
obj, _, err := s.codec.Decode(data, nil, obj)
return obj, err
}
@@ -209,6 +215,8 @@ func (s *Storage) convertToObject(data []byte, obj runtime.Object) (runtime.Obje
// in seconds (0 means forever). If no error is returned and out is not nil, out will be
// set to the read value from database.
func (s *Storage) Create(ctx context.Context, key string, obj runtime.Object, out runtime.Object, ttl uint64) error {
ctx, span := tracer.Start(ctx, "apistore.Storage.Create")
defer span.End()
v, err := s.prepareObjectForStorage(ctx, obj)
if err != nil {
return s.handleManagedResourceRouting(ctx, err, resourcepb.WatchEvent_ADDED, key, obj, out)
@@ -238,7 +246,7 @@ func (s *Storage) Create(ctx context.Context, key string, obj runtime.Object, ou
return v.finish(ctx, err, s.opts.SecureValues)
}
if _, err := s.convertToObject(req.Value, out); err != nil {
if _, err := s.convertToObject(ctx, req.Value, out); err != nil {
return err
}
@@ -274,6 +282,8 @@ func (s *Storage) Delete(
_ runtime.Object,
opts storage.DeleteOptions,
) error {
ctx, span := tracer.Start(ctx, "apistore.Storage.Delete")
defer span.End()
info, ok := authtypes.AuthInfoFrom(ctx)
if !ok {
return errors.New("missing auth info")
@@ -336,6 +346,8 @@ func (s *Storage) Delete(
// This version is not yet passing the watch tests
func (s *Storage) Watch(ctx context.Context, key string, opts storage.ListOptions) (watch.Interface, error) {
ctx, span := tracer.Start(ctx, "apistore.Storage.Watch")
defer span.End()
k, err := s.getKey(key)
if err != nil {
return watch.NewEmptyWatch(), nil
@@ -379,6 +391,8 @@ func (s *Storage) Watch(ctx context.Context, key string, opts storage.ListOption
// The returned contents may be delayed, but it is guaranteed that they will
// match 'opts.ResourceVersion' according 'opts.ResourceVersionMatch'.
func (s *Storage) Get(ctx context.Context, key string, opts storage.GetOptions, objPtr runtime.Object) error {
ctx, span := tracer.Start(ctx, "apistore.Storage.Get")
defer span.End()
var err error
req := &resourcepb.ReadRequest{}
req.Key, err = s.getKey(key)
@@ -410,7 +424,7 @@ func (s *Storage) Get(ctx context.Context, key string, opts storage.GetOptions,
return resource.GetError(rsp.Error)
}
_, err = s.convertToObject(rsp.Value, objPtr)
_, err = s.convertToObject(ctx, rsp.Value, objPtr)
if err != nil {
return err
}
@@ -424,6 +438,8 @@ func (s *Storage) Get(ctx context.Context, key string, opts storage.GetOptions,
// The returned contents may be delayed, but it is guaranteed that they will
// match 'opts.ResourceVersion' according 'opts.ResourceVersionMatch'.
func (s *Storage) GetList(ctx context.Context, key string, opts storage.ListOptions, listObj runtime.Object) error {
ctx, span := tracer.Start(ctx, "apistore.Storage.GetList")
defer span.End()
k, err := s.getKey(key)
if err != nil {
return err
@@ -460,30 +476,11 @@ func (s *Storage) GetList(ctx context.Context, key string, opts storage.ListOpti
}
for _, item := range rsp.Items {
obj, err := s.convertToObject(item.Value, s.newFunc())
obj, shouldAppend, err := s.processItem(ctx, item, opts, predicate)
if err != nil {
return err
}
if err := s.versioner.UpdateObject(obj, uint64(item.ResourceVersion)); err != nil {
return err
}
if opts.ResourceVersionMatch == metaV1.ResourceVersionMatchExact {
currentVersion, err := s.versioner.ObjectResourceVersion(obj)
if err != nil {
return err
}
expectedRV, err := s.versioner.ParseResourceVersion(opts.ResourceVersion)
if err != nil {
return err
}
if currentVersion != expectedRV {
continue
}
}
ok, err := predicate.Matches(obj)
if err == nil && ok {
if shouldAppend {
v.Set(reflect.Append(v, reflect.ValueOf(obj).Elem()))
}
}
@@ -498,6 +495,45 @@ func (s *Storage) GetList(ctx context.Context, key string, opts storage.ListOpti
return nil
}
// processItem converts a raw item from the response, updates its version,
// checks resource version matching, and applies predicates.
// Returns the processed object, whether it should be appended to results, and any error.
func (s *Storage) processItem(ctx context.Context, item *resourcepb.ResourceWrapper, opts storage.ListOptions, predicate storage.SelectionPredicate) (runtime.Object, bool, error) {
_, span := tracer.Start(ctx, "apistore.Storage.processItem")
defer span.End()
obj, err := s.convertToObject(ctx, item.Value, s.newFunc())
if err != nil {
return nil, false, err
}
if err := s.versioner.UpdateObject(obj, uint64(item.ResourceVersion)); err != nil {
return nil, false, err
}
// Skip items that don't match the exact resource version if specified
if opts.ResourceVersionMatch == metaV1.ResourceVersionMatchExact {
currentVersion, err := s.versioner.ObjectResourceVersion(obj)
if err != nil {
return nil, false, err
}
expectedRV, err := s.versioner.ParseResourceVersion(opts.ResourceVersion)
if err != nil {
return nil, false, err
}
if currentVersion != expectedRV {
return nil, false, nil
}
}
// Apply predicate filtering
ok, err := predicate.Matches(obj)
if err != nil || !ok {
return nil, false, nil
}
return obj, true, nil
}
// GuaranteedUpdate keeps calling 'tryUpdate()' to update key 'key' (of type 'destination')
// retrying the update until success if there is index conflict.
// Note that object passed to tryUpdate may change across invocations of tryUpdate() if
@@ -521,6 +557,8 @@ func (s *Storage) GuaranteedUpdate(
tryUpdate storage.UpdateFunc,
cachedExistingObject runtime.Object,
) error {
ctx, span := tracer.Start(ctx, "apistore.Storage.GuaranteedUpdate")
defer span.End()
var (
res storage.ResponseMeta
updatedObj runtime.Object
@@ -586,7 +624,7 @@ func (s *Storage) GuaranteedUpdate(
}
existingBytes = readResponse.Value
existingObj, err = s.convertToObject(readResponse.Value, s.newFunc())
existingObj, err = s.convertToObject(ctx, readResponse.Value, s.newFunc())
if err != nil {
return err
}
@@ -649,7 +687,7 @@ func (s *Storage) GuaranteedUpdate(
rv = uint64(updateResponse.ResourceVersion)
}
if _, err := s.convertToObject(req.Value, destination); err != nil {
if _, err := s.convertToObject(ctx, req.Value, destination); err != nil {
return err
}