diff --git a/pkg/storage/unified/apistore/store.go b/pkg/storage/unified/apistore/store.go index a6c37c00dab..2ad474dc90f 100644 --- a/pkg/storage/unified/apistore/store.go +++ b/pkg/storage/unified/apistore/store.go @@ -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 }