Storage: Consolidate error handling (#91167)
This commit is contained in:
@@ -17,8 +17,6 @@ import (
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
"go.opentelemetry.io/otel/trace/noop"
|
||||
"google.golang.org/protobuf/proto"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||||
)
|
||||
|
||||
const trace_prefix = "sql.resource."
|
||||
@@ -295,7 +293,7 @@ func (b *backend) delete(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
return newVersion, err
|
||||
}
|
||||
|
||||
func (b *backend) Read(ctx context.Context, req *resource.ReadRequest) (*resource.ReadResponse, error) {
|
||||
func (b *backend) ReadResource(ctx context.Context, req *resource.ReadRequest) *resource.ReadResponse {
|
||||
_, span := b.tracer.Start(ctx, trace_prefix+".Read")
|
||||
defer span.End()
|
||||
|
||||
@@ -315,23 +313,24 @@ func (b *backend) Read(ctx context.Context, req *resource.ReadRequest) (*resourc
|
||||
|
||||
res, err := dbutil.QueryRow(ctx, b.db, sr, readReq)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return nil, apierrors.NewNotFound(schema.GroupResource{
|
||||
Group: req.Key.Group,
|
||||
Resource: req.Key.Resource,
|
||||
}, req.Key.Name)
|
||||
return &resource.ReadResponse{
|
||||
Error: resource.NewNotFoundError(req.Key),
|
||||
}
|
||||
} else if err != nil {
|
||||
return nil, fmt.Errorf("get resource version: %w", err)
|
||||
return &resource.ReadResponse{Error: resource.AsErrorResult(err)}
|
||||
}
|
||||
|
||||
return &res.ReadResponse, nil
|
||||
return &res.ReadResponse
|
||||
}
|
||||
|
||||
func (b *backend) PrepareList(ctx context.Context, req *resource.ListRequest) (*resource.ListResponse, error) {
|
||||
func (b *backend) PrepareList(ctx context.Context, req *resource.ListRequest) *resource.ListResponse {
|
||||
_, span := b.tracer.Start(ctx, trace_prefix+"List")
|
||||
defer span.End()
|
||||
|
||||
if req.Options == nil || req.Options.Key.Group == "" || req.Options.Key.Resource == "" {
|
||||
return nil, fmt.Errorf("missing group or resource")
|
||||
return &resource.ListResponse{
|
||||
Error: resource.NewBadRequestError("missing group or resource"),
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: think about how to handler VersionMatch. We should be able to use latest for the first page (only).
|
||||
@@ -345,7 +344,7 @@ func (b *backend) PrepareList(ctx context.Context, req *resource.ListRequest) (*
|
||||
}
|
||||
|
||||
// listLatest fetches the resources from the resource table.
|
||||
func (b *backend) listLatest(ctx context.Context, req *resource.ListRequest) (*resource.ListResponse, error) {
|
||||
func (b *backend) listLatest(ctx context.Context, req *resource.ListRequest) *resource.ListResponse {
|
||||
out := &resource.ListResponse{
|
||||
ResourceVersion: 0,
|
||||
}
|
||||
@@ -387,19 +386,23 @@ func (b *backend) listLatest(ctx context.Context, req *resource.ListRequest) (*r
|
||||
|
||||
return nil
|
||||
})
|
||||
|
||||
return out, err
|
||||
if err != nil {
|
||||
out.Error = resource.AsErrorResult(err)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// listAtRevision fetches the resources from the resource_history table at a specific revision.
|
||||
func (b *backend) listAtRevision(ctx context.Context, req *resource.ListRequest) (*resource.ListResponse, error) {
|
||||
func (b *backend) listAtRevision(ctx context.Context, req *resource.ListRequest) *resource.ListResponse {
|
||||
// Get the RV
|
||||
rv := req.ResourceVersion
|
||||
offset := int64(0)
|
||||
if req.NextPageToken != "" {
|
||||
continueToken, err := GetContinueToken(req.NextPageToken)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("get continue token: %w", err)
|
||||
return &resource.ListResponse{
|
||||
Error: resource.AsErrorResult(fmt.Errorf("get continue token: %w", err)),
|
||||
}
|
||||
}
|
||||
rv = continueToken.ResourceVersion
|
||||
offset = continueToken.StartOffset
|
||||
@@ -443,8 +446,10 @@ func (b *backend) listAtRevision(ctx context.Context, req *resource.ListRequest)
|
||||
|
||||
return nil
|
||||
})
|
||||
|
||||
return out, err
|
||||
if err != nil {
|
||||
out.Error = resource.AsErrorResult(err)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func (b *backend) WatchWriteEvents(ctx context.Context) (<-chan *resource.WrittenEvent, error) {
|
||||
|
||||
@@ -88,14 +88,14 @@ func TestIntegrationBackendHappyPath(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("Read latest item 2", func(t *testing.T) {
|
||||
resp, err := store.Read(ctx, &resource.ReadRequest{Key: resourceKey("item2")})
|
||||
resp := store.ReadResource(ctx, &resource.ReadRequest{Key: resourceKey("item2")})
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(4), resp.ResourceVersion)
|
||||
require.Equal(t, "item2 MODIFIED", string(resp.Value))
|
||||
})
|
||||
|
||||
t.Run("Read early verion of item2", func(t *testing.T) {
|
||||
resp, err := store.Read(ctx, &resource.ReadRequest{
|
||||
resp := store.ReadResource(ctx, &resource.ReadRequest{
|
||||
Key: resourceKey("item2"),
|
||||
ResourceVersion: 3, // item2 was created at rv=2 and updated at rv=4
|
||||
})
|
||||
@@ -105,7 +105,7 @@ func TestIntegrationBackendHappyPath(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("PrepareList latest", func(t *testing.T) {
|
||||
resp, err := store.PrepareList(ctx, &resource.ListRequest{
|
||||
resp := store.PrepareList(ctx, &resource.ListRequest{
|
||||
Options: &resource.ListOptions{
|
||||
Key: &resource.ResourceKey{
|
||||
Namespace: "namespace",
|
||||
@@ -188,7 +188,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
_, _ = writeEvent(ctx, store, "item3", resource.WatchEvent_DELETED) // rv=7
|
||||
_, _ = writeEvent(ctx, store, "item6", resource.WatchEvent_ADDED) // rv=8
|
||||
t.Run("fetch all latest", func(t *testing.T) {
|
||||
res, err := store.PrepareList(ctx, &resource.ListRequest{
|
||||
res := store.PrepareList(ctx, &resource.ListRequest{
|
||||
Options: &resource.ListOptions{
|
||||
Key: &resource.ResourceKey{
|
||||
Group: "group",
|
||||
@@ -196,7 +196,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
},
|
||||
},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Nil(t, res.Error)
|
||||
require.Len(t, res.Items, 5)
|
||||
// should be sorted by resource version DESC
|
||||
require.Equal(t, "item6 ADDED", string(res.Items[0].Value))
|
||||
@@ -209,7 +209,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("list latest first page ", func(t *testing.T) {
|
||||
res, err := store.PrepareList(ctx, &resource.ListRequest{
|
||||
res := store.PrepareList(ctx, &resource.ListRequest{
|
||||
Limit: 3,
|
||||
Options: &resource.ListOptions{
|
||||
Key: &resource.ResourceKey{
|
||||
@@ -218,7 +218,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
},
|
||||
},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Nil(t, res.Error)
|
||||
require.Len(t, res.Items, 3)
|
||||
continueToken, err := sql.GetContinueToken(res.NextPageToken)
|
||||
require.NoError(t, err)
|
||||
@@ -230,7 +230,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("list at revision", func(t *testing.T) {
|
||||
res, err := store.PrepareList(ctx, &resource.ListRequest{
|
||||
res := store.PrepareList(ctx, &resource.ListRequest{
|
||||
ResourceVersion: 4,
|
||||
Options: &resource.ListOptions{
|
||||
Key: &resource.ResourceKey{
|
||||
@@ -239,7 +239,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
},
|
||||
},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Nil(t, res.Error)
|
||||
require.Len(t, res.Items, 4)
|
||||
require.Equal(t, "item4 ADDED", string(res.Items[0].Value))
|
||||
require.Equal(t, "item3 ADDED", string(res.Items[1].Value))
|
||||
@@ -249,7 +249,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("fetch first page at revision with limit", func(t *testing.T) {
|
||||
res, err := store.PrepareList(ctx, &resource.ListRequest{
|
||||
res := store.PrepareList(ctx, &resource.ListRequest{
|
||||
Limit: 3,
|
||||
ResourceVersion: 7,
|
||||
Options: &resource.ListOptions{
|
||||
@@ -259,7 +259,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
},
|
||||
},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Nil(t, res.Error)
|
||||
require.Len(t, res.Items, 3)
|
||||
t.Log(res.Items)
|
||||
require.Equal(t, "item2 MODIFIED", string(res.Items[0].Value))
|
||||
@@ -277,7 +277,7 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
ResourceVersion: 8,
|
||||
StartOffset: 2,
|
||||
}
|
||||
res, err := store.PrepareList(ctx, &resource.ListRequest{
|
||||
res := store.PrepareList(ctx, &resource.ListRequest{
|
||||
NextPageToken: continueToken.String(),
|
||||
Limit: 2,
|
||||
Options: &resource.ListOptions{
|
||||
@@ -287,12 +287,12 @@ func TestIntegrationBackendPrepareList(t *testing.T) {
|
||||
},
|
||||
},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Nil(t, res.Error)
|
||||
require.Len(t, res.Items, 2)
|
||||
require.Equal(t, "item5 ADDED", string(res.Items[0].Value))
|
||||
require.Equal(t, "item4 ADDED", string(res.Items[1].Value))
|
||||
|
||||
continueToken, err = sql.GetContinueToken(res.NextPageToken)
|
||||
continueToken, err := sql.GetContinueToken(res.NextPageToken)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(8), continueToken.ResourceVersion)
|
||||
require.Equal(t, int64(4), continueToken.StartOffset)
|
||||
|
||||
Reference in New Issue
Block a user