feat(unified-storage): provide delete function for bucket (#103825)
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user