Storage: Propagate RV to server on update+delete (#111866)
This commit is contained in:
@@ -293,13 +293,6 @@ func (s *Storage) Delete(
|
||||
if err := preconditions.Check(key, out); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if preconditions.ResourceVersion != nil {
|
||||
cmd.ResourceVersion, err = strconv.ParseInt(*preconditions.ResourceVersion, 10, 64)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if preconditions.UID != nil {
|
||||
cmd.Uid = string(*preconditions.UID)
|
||||
}
|
||||
@@ -319,6 +312,10 @@ func (s *Storage) Delete(
|
||||
return s.handleManagedResourceRouting(ctx, err, resourcepb.WatchEvent_DELETED, key, out, out)
|
||||
}
|
||||
|
||||
cmd.ResourceVersion, err = meta.GetResourceVersionInt64()
|
||||
if err != nil {
|
||||
return resource.GetError(resource.AsErrorResult(err))
|
||||
}
|
||||
rsp, err := s.store.Delete(ctx, cmd)
|
||||
if err != nil {
|
||||
return resource.GetError(resource.AsErrorResult(err))
|
||||
@@ -611,41 +608,45 @@ func (s *Storage) GuaranteedUpdate(
|
||||
}
|
||||
continue
|
||||
}
|
||||
break
|
||||
}
|
||||
|
||||
v, err := s.prepareObjectForUpdate(ctx, updatedObj, existingObj)
|
||||
if err != nil {
|
||||
return s.handleManagedResourceRouting(ctx, err, resourcepb.WatchEvent_MODIFIED, key, updatedObj, destination)
|
||||
}
|
||||
|
||||
// Only update (for real) if the bytes have changed
|
||||
var rv uint64
|
||||
req.Value = v.raw.Bytes()
|
||||
if !bytes.Equal(req.Value, existingBytes) {
|
||||
updateResponse, err := s.store.Update(ctx, req)
|
||||
v, err := s.prepareObjectForUpdate(ctx, updatedObj, existingObj)
|
||||
if err != nil {
|
||||
err = resource.GetError(resource.AsErrorResult(err))
|
||||
} else if updateResponse.Error != nil {
|
||||
err = resource.GetError(updateResponse.Error)
|
||||
return s.handleManagedResourceRouting(ctx, err, resourcepb.WatchEvent_MODIFIED, key, updatedObj, destination)
|
||||
}
|
||||
|
||||
// Cleanup secure values
|
||||
if err = v.finish(ctx, err, s.opts.SecureValues); err != nil {
|
||||
// Only update (for real) if the bytes have changed
|
||||
var rv uint64
|
||||
req.Value = v.raw.Bytes()
|
||||
if !bytes.Equal(req.Value, existingBytes) {
|
||||
req.ResourceVersion = readResponse.ResourceVersion
|
||||
updateResponse, err := s.store.Update(ctx, req)
|
||||
if err != nil {
|
||||
err = resource.GetError(resource.AsErrorResult(err))
|
||||
} else if updateResponse.Error != nil {
|
||||
if attempt < MaxUpdateAttempts && updateResponse.Error.Code == http.StatusConflict {
|
||||
continue // try the read again
|
||||
}
|
||||
err = resource.GetError(updateResponse.Error)
|
||||
}
|
||||
|
||||
// Cleanup secure values
|
||||
if err = v.finish(ctx, err, s.opts.SecureValues); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
rv = uint64(updateResponse.ResourceVersion)
|
||||
}
|
||||
|
||||
if _, err := s.convertToObject(req.Value, destination); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
rv = uint64(updateResponse.ResourceVersion)
|
||||
}
|
||||
|
||||
if _, err := s.convertToObject(req.Value, destination); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if rv > 0 {
|
||||
if err := s.versioner.UpdateObject(destination, rv); err != nil {
|
||||
return err
|
||||
if rv > 0 {
|
||||
if err := s.versioner.UpdateObject(destination, rv); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
@@ -18,8 +18,7 @@ import (
|
||||
|
||||
// Package-level errors.
|
||||
var (
|
||||
ErrOptimisticLockingFailed = errors.New("optimistic locking failed")
|
||||
ErrNotImplementedYet = errors.New("not implemented yet")
|
||||
ErrNotImplementedYet = errors.New("not implemented yet")
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -31,6 +30,12 @@ var (
|
||||
Code: http.StatusConflict,
|
||||
},
|
||||
}
|
||||
|
||||
ErrOptimisticLockingFailed = resourcepb.ErrorResult{
|
||||
Code: http.StatusConflict,
|
||||
Reason: "optimistic locking failed",
|
||||
Message: "requested RV does not match saved RV",
|
||||
}
|
||||
)
|
||||
|
||||
func NewBadRequestError(msg string) *resourcepb.ErrorResult {
|
||||
|
||||
@@ -731,8 +731,12 @@ func (s *server) update(ctx context.Context, user claims.AuthInfo, req *resource
|
||||
return rsp, nil
|
||||
}
|
||||
|
||||
// TODO: once we know the client is always sending the RV, require ResourceVersion > 0
|
||||
// See: https://github.com/grafana/grafana/pull/111866
|
||||
if req.ResourceVersion > 0 && latest.ResourceVersion != req.ResourceVersion {
|
||||
return nil, ErrOptimisticLockingFailed
|
||||
return &resourcepb.UpdateResponse{
|
||||
Error: &ErrOptimisticLockingFailed,
|
||||
}, nil
|
||||
}
|
||||
|
||||
event, e := s.newEvent(ctx, user, req.Key, req.Value, latest.Value)
|
||||
@@ -796,7 +800,7 @@ func (s *server) delete(ctx context.Context, user claims.AuthInfo, req *resource
|
||||
return rsp, nil
|
||||
}
|
||||
if req.ResourceVersion > 0 && latest.ResourceVersion != req.ResourceVersion {
|
||||
rsp.Error = AsErrorResult(ErrOptimisticLockingFailed)
|
||||
rsp.Error = &ErrOptimisticLockingFailed
|
||||
return rsp, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -477,11 +477,12 @@ func TestSimpleServer(t *testing.T) {
|
||||
ResourceVersion: created.ResourceVersion})
|
||||
require.NoError(t, err)
|
||||
|
||||
_, err = server.Update(ctx, &resourcepb.UpdateRequest{
|
||||
rsp, _ := server.Update(ctx, &resourcepb.UpdateRequest{
|
||||
Key: key,
|
||||
Value: raw,
|
||||
ResourceVersion: created.ResourceVersion})
|
||||
require.ErrorIs(t, err, ErrOptimisticLockingFailed)
|
||||
require.Equal(t, rsp.Error.Code, ErrOptimisticLockingFailed.Code)
|
||||
require.Equal(t, rsp.Error.Message, ErrOptimisticLockingFailed.Message)
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -18,6 +18,7 @@ import (
|
||||
"go.opentelemetry.io/otel/trace/noop"
|
||||
"google.golang.org/protobuf/proto"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||||
|
||||
"github.com/grafana/grafana/pkg/util/sqlite"
|
||||
|
||||
@@ -405,15 +406,18 @@ func (b *backend) update(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
// Use rvManager.ExecWithRV instead of direct transaction
|
||||
rv, err := b.rvManager.ExecWithRV(ctx, event.Key, func(tx db.Tx) (string, error) {
|
||||
// 1. Update resource
|
||||
_, err := dbutil.Exec(ctx, tx, sqlResourceUpdate, sqlResourceRequest{
|
||||
res, err := dbutil.Exec(ctx, tx, sqlResourceUpdate, sqlResourceRequest{
|
||||
SQLTemplate: sqltemplate.New(b.dialect),
|
||||
WriteEvent: event,
|
||||
WriteEvent: event, // includes the RV
|
||||
Folder: folder,
|
||||
GUID: event.GUID,
|
||||
})
|
||||
if err != nil {
|
||||
return event.GUID, fmt.Errorf("resource update: %w", err)
|
||||
}
|
||||
if err = b.checkConflict(res, event.Key, event.PreviousRV); err != nil {
|
||||
return event.GUID, err
|
||||
}
|
||||
|
||||
// 2. Insert into resource history
|
||||
if _, err := dbutil.Exec(ctx, tx, sqlResourceHistoryInsert, sqlResourceRequest{
|
||||
@@ -460,7 +464,7 @@ func (b *backend) delete(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
}
|
||||
rv, err := b.rvManager.ExecWithRV(ctx, event.Key, func(tx db.Tx) (string, error) {
|
||||
// 1. delete from resource
|
||||
_, err := dbutil.Exec(ctx, tx, sqlResourceDelete, sqlResourceRequest{
|
||||
res, err := dbutil.Exec(ctx, tx, sqlResourceDelete, sqlResourceRequest{
|
||||
SQLTemplate: sqltemplate.New(b.dialect),
|
||||
WriteEvent: event,
|
||||
GUID: event.GUID,
|
||||
@@ -468,6 +472,9 @@ func (b *backend) delete(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
if err != nil {
|
||||
return event.GUID, fmt.Errorf("delete resource: %w", err)
|
||||
}
|
||||
if err = b.checkConflict(res, event.Key, event.PreviousRV); err != nil {
|
||||
return event.GUID, err
|
||||
}
|
||||
|
||||
// 2. Add event to resource history
|
||||
if _, err := dbutil.Exec(ctx, tx, sqlResourceHistoryInsert, sqlResourceRequest{
|
||||
@@ -504,6 +511,28 @@ func (b *backend) delete(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
return rv, nil
|
||||
}
|
||||
|
||||
func (b *backend) checkConflict(res db.Result, key *resourcepb.ResourceKey, rv int64) error {
|
||||
if rv == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
// The RV is part of the update request, and it may no longer be the most recent
|
||||
rows, err := res.RowsAffected()
|
||||
if err != nil {
|
||||
return fmt.Errorf("unable to verify RV: %w", err)
|
||||
}
|
||||
if rows == 1 {
|
||||
return nil // expected one result
|
||||
}
|
||||
if rows > 0 {
|
||||
return fmt.Errorf("multiple rows effected (%d)", rows)
|
||||
}
|
||||
return apierrors.NewConflict(schema.GroupResource{
|
||||
Group: key.Group,
|
||||
Resource: key.Resource,
|
||||
}, key.Name, fmt.Errorf("resource version does not match current value"))
|
||||
}
|
||||
|
||||
func (b *backend) ReadResource(ctx context.Context, req *resourcepb.ReadRequest) *resource.BackendReadResponse {
|
||||
_, span := b.tracer.Start(ctx, tracePrefix+".Read")
|
||||
defer span.End()
|
||||
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
"testing"
|
||||
|
||||
"github.com/DATA-DOG/go-sqlmock"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
||||
|
||||
|
||||
@@ -5,5 +5,6 @@ DELETE FROM {{ .Ident "resource" }}
|
||||
AND {{ .Ident "resource" }} = {{ .Arg .WriteEvent.Key.Resource }}
|
||||
{{ if .WriteEvent.Key.Name }}
|
||||
AND {{ .Ident "name" }} = {{ .Arg .WriteEvent.Key.Name }}
|
||||
AND {{ .Ident "resource_version" }} = {{ .Arg .WriteEvent.PreviousRV }}
|
||||
{{ end }}
|
||||
;
|
||||
|
||||
@@ -10,4 +10,5 @@ UPDATE {{ .Ident "resource" }}
|
||||
AND {{ .Ident "resource" }} = {{ .Arg .WriteEvent.Key.Resource }}
|
||||
AND {{ .Ident "namespace" }} = {{ .Arg .WriteEvent.Key.Namespace }}
|
||||
AND {{ .Ident "name" }} = {{ .Arg .WriteEvent.Key.Name }}
|
||||
AND {{ .Ident "resource_version" }} = {{ .Arg .WriteEvent.PreviousRV }}
|
||||
;
|
||||
|
||||
@@ -87,14 +87,23 @@ func TestIntegrationListIter(t *testing.T) {
|
||||
Group: item.group,
|
||||
Name: item.name,
|
||||
},
|
||||
Value: item.value,
|
||||
PreviousRV: 0,
|
||||
Value: item.value,
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to insert test data: %w", err)
|
||||
}
|
||||
_, err = dbutil.Exec(ctx, tx, sqlResourceUpdate, sqlResourceRequest{
|
||||
|
||||
if _, err = dbutil.Exec(ctx, tx, sqlResourceUpdateRV, sqlResourceUpdateRVRequest{
|
||||
SQLTemplate: sqltemplate.New(dialect),
|
||||
GUIDToRV: map[string]int64{
|
||||
item.guid: item.resourceVersion,
|
||||
},
|
||||
}); err != nil {
|
||||
return fmt.Errorf("failed to insert test data: %w", err)
|
||||
}
|
||||
|
||||
if _, err = dbutil.Exec(ctx, tx, sqlResourceUpdate, sqlResourceRequest{
|
||||
SQLTemplate: sqltemplate.New(dialect),
|
||||
GUID: item.guid,
|
||||
ResourceVersion: item.resourceVersion,
|
||||
@@ -110,8 +119,7 @@ func TestIntegrationListIter(t *testing.T) {
|
||||
PreviousRV: item.resourceVersion,
|
||||
Type: 1,
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
}); err != nil {
|
||||
return fmt.Errorf("failed to insert resource version: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -31,6 +31,21 @@ func TestUnifiedStorageQueries(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "with rv",
|
||||
Data: &sqlResourceRequest{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
WriteEvent: resource.WriteEvent{
|
||||
Key: &resourcepb.ResourceKey{
|
||||
Namespace: "nn",
|
||||
Group: "gg",
|
||||
Resource: "rr",
|
||||
Name: "name",
|
||||
},
|
||||
PreviousRV: 1234,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
sqlResourceInsert: {
|
||||
{
|
||||
@@ -63,6 +78,7 @@ func TestUnifiedStorageQueries(t *testing.T) {
|
||||
Resource: "rr",
|
||||
Name: "name",
|
||||
},
|
||||
PreviousRV: 1759304090100678,
|
||||
},
|
||||
Folder: "fldr",
|
||||
},
|
||||
|
||||
@@ -263,7 +263,7 @@ func (m *resourceVersionManager) execBatch(ctx context.Context, group, resource
|
||||
attribute.Int("operation_index", i),
|
||||
attribute.String("error", err.Error()),
|
||||
))
|
||||
return fmt.Errorf("failed to execute function: %w", err)
|
||||
return err
|
||||
}
|
||||
guids[i] = guid
|
||||
}
|
||||
|
||||
@@ -4,4 +4,5 @@ DELETE FROM `resource`
|
||||
AND `group` = 'gg'
|
||||
AND `resource` = 'rr'
|
||||
AND `name` = 'name'
|
||||
AND `resource_version` = 0
|
||||
;
|
||||
|
||||
+8
@@ -0,0 +1,8 @@
|
||||
DELETE FROM `resource`
|
||||
WHERE 1 = 1
|
||||
AND `namespace` = 'nn'
|
||||
AND `group` = 'gg'
|
||||
AND `resource` = 'rr'
|
||||
AND `name` = 'name'
|
||||
AND `resource_version` = 1234
|
||||
;
|
||||
@@ -10,4 +10,5 @@ UPDATE `resource`
|
||||
AND `resource` = 'rr'
|
||||
AND `namespace` = 'nn'
|
||||
AND `name` = 'name'
|
||||
AND `resource_version` = 1759304090100678
|
||||
;
|
||||
|
||||
@@ -4,4 +4,5 @@ DELETE FROM "resource"
|
||||
AND "group" = 'gg'
|
||||
AND "resource" = 'rr'
|
||||
AND "name" = 'name'
|
||||
AND "resource_version" = 0
|
||||
;
|
||||
|
||||
+8
@@ -0,0 +1,8 @@
|
||||
DELETE FROM "resource"
|
||||
WHERE 1 = 1
|
||||
AND "namespace" = 'nn'
|
||||
AND "group" = 'gg'
|
||||
AND "resource" = 'rr'
|
||||
AND "name" = 'name'
|
||||
AND "resource_version" = 1234
|
||||
;
|
||||
@@ -10,4 +10,5 @@ UPDATE "resource"
|
||||
AND "resource" = 'rr'
|
||||
AND "namespace" = 'nn'
|
||||
AND "name" = 'name'
|
||||
AND "resource_version" = 1759304090100678
|
||||
;
|
||||
|
||||
@@ -4,4 +4,5 @@ DELETE FROM "resource"
|
||||
AND "group" = 'gg'
|
||||
AND "resource" = 'rr'
|
||||
AND "name" = 'name'
|
||||
AND "resource_version" = 0
|
||||
;
|
||||
|
||||
+8
@@ -0,0 +1,8 @@
|
||||
DELETE FROM "resource"
|
||||
WHERE 1 = 1
|
||||
AND "namespace" = 'nn'
|
||||
AND "group" = 'gg'
|
||||
AND "resource" = 'rr'
|
||||
AND "name" = 'name'
|
||||
AND "resource_version" = 1234
|
||||
;
|
||||
@@ -10,4 +10,5 @@ UPDATE "resource"
|
||||
AND "resource" = 'rr'
|
||||
AND "namespace" = 'nn'
|
||||
AND "name" = 'name'
|
||||
AND "resource_version" = 1759304090100678
|
||||
;
|
||||
|
||||
@@ -124,13 +124,13 @@ func runTestIntegrationBackendHappyPath(t *testing.T, backend resource.StorageBa
|
||||
})
|
||||
|
||||
t.Run("Update item2", func(t *testing.T) {
|
||||
rv4, err = writeEvent(ctx, backend, "item2", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rv4, err = writeEvent(ctx, backend, "item2", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rv2))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv4, rv3)
|
||||
})
|
||||
|
||||
t.Run("Delete item1", func(t *testing.T) {
|
||||
rv5, err = writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_DELETED, WithNamespace(ns))
|
||||
rv5, err = writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_DELETED, WithNamespaceAndRV(ns, rv1))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv5, rv4)
|
||||
})
|
||||
@@ -352,10 +352,10 @@ func runTestIntegrationBackendList(t *testing.T, backend resource.StorageBackend
|
||||
rv5, err := writeEvent(ctx, backend, "item5", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv5, rv4)
|
||||
rv6, err := writeEvent(ctx, backend, "item2", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rv6, err := writeEvent(ctx, backend, "item2", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rv2))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv6, rv5)
|
||||
rv7, err := writeEvent(ctx, backend, "item3", resourcepb.WatchEvent_DELETED, WithNamespace(ns))
|
||||
rv7, err := writeEvent(ctx, backend, "item3", resourcepb.WatchEvent_DELETED, WithNamespaceAndRV(ns, rv3))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv7, rv6)
|
||||
rv8, err := writeEvent(ctx, backend, "item6", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
@@ -490,10 +490,10 @@ func runTestIntegrationBackendListModifiedSince(t *testing.T, backend resource.S
|
||||
ns := nsPrefix + "-history-ns"
|
||||
rvCreated, _ := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
require.Greater(t, rvCreated, int64(0))
|
||||
rvUpdated, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rvUpdated, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rvCreated))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvUpdated, rvCreated)
|
||||
rvDeleted, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_DELETED, WithNamespace(ns))
|
||||
rvDeleted, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_DELETED, WithNamespaceAndRV(ns, rvUpdated))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvDeleted, rvUpdated)
|
||||
|
||||
@@ -610,19 +610,19 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage
|
||||
require.Greater(t, rv1, int64(0))
|
||||
|
||||
// add 5 events for item1 - should be saved to history
|
||||
rvHistory1, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rvHistory1, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rv1))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvHistory1, rv1)
|
||||
rvHistory2, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rvHistory2, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rvHistory1))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvHistory2, rvHistory1)
|
||||
rvHistory3, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rvHistory3, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rvHistory2))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvHistory3, rvHistory2)
|
||||
rvHistory4, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rvHistory4, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rvHistory3))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvHistory4, rvHistory3)
|
||||
rvHistory5, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rvHistory5, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rvHistory4))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvHistory5, rvHistory4)
|
||||
|
||||
@@ -804,8 +804,9 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage
|
||||
resourceVersions = append(resourceVersions, initialRV)
|
||||
|
||||
// Create 9 more versions with modifications
|
||||
rv := initialRV
|
||||
for i := 0; i < 9; i++ {
|
||||
rv, err := writeEvent(ctx, backend, "paged-item", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns2))
|
||||
rv, err = writeEvent(ctx, backend, "paged-item", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns2, rv))
|
||||
require.NoError(t, err)
|
||||
resourceVersions = append(resourceVersions, rv)
|
||||
}
|
||||
@@ -907,7 +908,7 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage
|
||||
// Create a resource and delete it
|
||||
rv, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
require.NoError(t, err)
|
||||
rvDeleted, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_DELETED, WithNamespace(ns))
|
||||
rvDeleted, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_DELETED, WithNamespaceAndRV(ns, rv))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvDeleted, rv)
|
||||
|
||||
@@ -932,7 +933,7 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage
|
||||
// Create a resource and delete it
|
||||
rv, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
require.NoError(t, err)
|
||||
rvDeleted, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_DELETED, WithNamespace(ns))
|
||||
rvDeleted, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_DELETED, WithNamespaceAndRV(ns, rv))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvDeleted, rv)
|
||||
|
||||
@@ -940,7 +941,7 @@ func runTestIntegrationBackendListHistory(t *testing.T, backend resource.Storage
|
||||
rv1, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv1, rvDeleted)
|
||||
rv2, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_MODIFIED, WithNamespace(ns))
|
||||
rv2, err := writeEvent(ctx, backend, "deleted-item", resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, rv1))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv2, rv1)
|
||||
|
||||
@@ -983,8 +984,8 @@ func runTestIntegrationBackendListHistoryErrorReporting(t *testing.T, backend re
|
||||
|
||||
const events = 500
|
||||
prevRv := origRv
|
||||
for i := 0; i < events; i++ {
|
||||
rv, err := writeEvent(ctx, backend, name, resourcepb.WatchEvent_MODIFIED, WithNamespace(ns), WithGroup(group), WithResource(resourceName))
|
||||
for range events {
|
||||
rv, err := writeEvent(ctx, backend, name, resourcepb.WatchEvent_MODIFIED, WithNamespaceAndRV(ns, prevRv), WithGroup(group), WithResource(resourceName))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv, prevRv)
|
||||
prevRv = rv
|
||||
@@ -1131,6 +1132,14 @@ func runTestIntegrationBackendCreateNewResource(t *testing.T, backend resource.S
|
||||
// WriteEventOption is a function that modifies WriteEventOptions
|
||||
type WriteEventOption func(*WriteEventOptions)
|
||||
|
||||
// WithNamespace sets the namespace for the write event
|
||||
func WithNamespaceAndRV(namespace string, rv int64) WriteEventOption {
|
||||
return func(o *WriteEventOptions) {
|
||||
o.Namespace = namespace
|
||||
o.PreviousRV = rv
|
||||
}
|
||||
}
|
||||
|
||||
// WithNamespace sets the namespace for the write event
|
||||
func WithNamespace(namespace string) WriteEventOption {
|
||||
return func(o *WriteEventOptions) {
|
||||
@@ -1180,11 +1189,12 @@ func WithValue(value string) WriteEventOption {
|
||||
}
|
||||
|
||||
type WriteEventOptions struct {
|
||||
Namespace string
|
||||
Group string
|
||||
Resource string
|
||||
Folder string
|
||||
Value []byte
|
||||
Namespace string
|
||||
Group string
|
||||
Resource string
|
||||
Folder string
|
||||
Value []byte
|
||||
PreviousRV int64
|
||||
}
|
||||
|
||||
func writeEvent(ctx context.Context, store resource.StorageBackend, name string, action resourcepb.WatchEvent_Type, opts ...WriteEventOption) (int64, error) {
|
||||
@@ -1236,6 +1246,7 @@ func writeEvent(ctx context.Context, store resource.StorageBackend, name string,
|
||||
Resource: options.Resource,
|
||||
Name: name,
|
||||
},
|
||||
PreviousRV: options.PreviousRV,
|
||||
}
|
||||
switch action {
|
||||
case resourcepb.WatchEvent_DELETED:
|
||||
@@ -1285,18 +1296,15 @@ func runTestIntegrationBackendTrash(t *testing.T, backend resource.StorageBacken
|
||||
rv1, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv1, int64(0))
|
||||
rvDelete1, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_DELETED, WithNamespace(ns))
|
||||
rvDelete1, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_DELETED, WithNamespaceAndRV(ns, rv1))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvDelete1, rv1)
|
||||
rvDelete2, err := writeEvent(ctx, backend, "item1", resourcepb.WatchEvent_DELETED, WithNamespace(ns))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvDelete2, rvDelete1)
|
||||
|
||||
// item2 deleted and recreated, should not be returned in trash
|
||||
rv2, err := writeEvent(ctx, backend, "item2", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv2, int64(0))
|
||||
rvDelete3, err := writeEvent(ctx, backend, "item2", resourcepb.WatchEvent_DELETED, WithNamespace(ns))
|
||||
rvDelete3, err := writeEvent(ctx, backend, "item2", resourcepb.WatchEvent_DELETED, WithNamespaceAndRV(ns, rv2))
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rvDelete3, rv2)
|
||||
rv3, err := writeEvent(ctx, backend, "item2", resourcepb.WatchEvent_ADDED, WithNamespace(ns))
|
||||
@@ -1325,10 +1333,10 @@ func runTestIntegrationBackendTrash(t *testing.T, backend resource.StorageBacken
|
||||
},
|
||||
},
|
||||
},
|
||||
expectedVersions: []int64{rvDelete2},
|
||||
expectedVersions: []int64{rvDelete1},
|
||||
expectedValues: []string{"item1 DELETED"},
|
||||
minExpectedHeadRV: rvDelete2,
|
||||
expectedContinueRV: rvDelete2,
|
||||
minExpectedHeadRV: rvDelete1,
|
||||
expectedContinueRV: rvDelete1,
|
||||
expectedSortAsc: false,
|
||||
},
|
||||
{
|
||||
|
||||
Reference in New Issue
Block a user