diff --git a/pkg/storage/unified/resource/cdk_bucket.go b/pkg/storage/unified/resource/cdk_bucket.go index 7a789c26eb1..e148326f8a9 100644 --- a/pkg/storage/unified/resource/cdk_bucket.go +++ b/pkg/storage/unified/resource/cdk_bucket.go @@ -19,10 +19,12 @@ type CDKBucket interface { ListPage(context.Context, []byte, int, *blob.ListOptions) ([]*blob.ListObject, []byte, error) WriteAll(context.Context, string, []byte, *blob.WriterOptions) error ReadAll(context.Context, string) ([]byte, error) + Delete(context.Context, string) error SignedURL(context.Context, string, *blob.SignedURLOptions) (string, error) } var _ CDKBucket = (*blob.Bucket)(nil) +var _ CDKBucket = (*InstrumentedBucket)(nil) const ( cdkBucketOperationLabel = "operation" @@ -159,6 +161,29 @@ func (b *InstrumentedBucket) WriteAll(ctx context.Context, key string, p []byte, return err } +func (b *InstrumentedBucket) Delete(ctx context.Context, key string) error { + ctx, span := b.tracer.Start(ctx, "InstrumentedBucket/Delete") + defer span.End() + start := time.Now() + err := b.bucket.Delete(ctx, key) + end := time.Since(start).Seconds() + labels := prometheus.Labels{ + cdkBucketOperationLabel: "Delete", + } + if err != nil { + labels[cdkBucketStatusLabel] = cdkBucketStatusError + b.requests.With(labels).Inc() + b.latency.With(labels).Observe(end) + span.RecordError(err) + span.SetStatus(codes.Error, err.Error()) + return err + } + labels[cdkBucketStatusLabel] = cdkBucketStatusSuccess + b.requests.With(labels).Inc() + b.latency.With(labels).Observe(end) + return nil +} + func (b *InstrumentedBucket) SignedURL(ctx context.Context, key string, opts *blob.SignedURLOptions) (string, error) { ctx, span := b.tracer.Start(ctx, "InstrumentedBucket/SignedURL") defer span.End() diff --git a/pkg/storage/unified/resource/cdk_bucket_test.go b/pkg/storage/unified/resource/cdk_bucket_test.go index 71949d94f80..9a577a23864 100644 --- a/pkg/storage/unified/resource/cdk_bucket_test.go +++ b/pkg/storage/unified/resource/cdk_bucket_test.go @@ -12,6 +12,8 @@ import ( "gocloud.dev/blob" ) +var _ CDKBucket = (*fakeCDKBucket)(nil) + type fakeCDKBucket struct { attributesFunc func(ctx context.Context, key string) (*blob.Attributes, error) writeAllFunc func(ctx context.Context, key string, p []byte, opts *blob.WriterOptions) error @@ -19,6 +21,7 @@ type fakeCDKBucket struct { signedURLFunc func(ctx context.Context, key string, opts *blob.SignedURLOptions) (string, error) listFunc func(opts *blob.ListOptions) *blob.ListIterator listPageFunc func(ctx context.Context, pageToken []byte, pageSize int, opts *blob.ListOptions) ([]*blob.ListObject, []byte, error) + deleteFunc func(ctx context.Context, key string) error } func (f *fakeCDKBucket) Attributes(ctx context.Context, key string) (*blob.Attributes, error) { @@ -63,6 +66,13 @@ func (f *fakeCDKBucket) ListPage(ctx context.Context, pageToken []byte, pageSize return nil, nil, nil } +func (f *fakeCDKBucket) Delete(ctx context.Context, key string) error { + if f.deleteFunc != nil { + return f.deleteFunc(ctx, key) + } + return nil +} + func TestInstrumentedBucket(t *testing.T) { operations := []struct { name string @@ -108,6 +118,24 @@ func TestInstrumentedBucket(t *testing.T) { return err }, }, + { + name: "Delete", + operation: "Delete", + setup: func(fakeBucket *fakeCDKBucket, success bool) { + if success { + fakeBucket.deleteFunc = func(ctx context.Context, key string) error { + return nil + } + } else { + fakeBucket.deleteFunc = func(ctx context.Context, key string) error { + return fmt.Errorf("some error") + } + } + }, + call: func(instrumentedBucket *InstrumentedBucket) error { + return instrumentedBucket.Delete(context.Background(), "key") + }, + }, { name: "ReadAll", operation: "ReadAll",