refactor(unified-storage): set the GUID in the resource server (#105683)
This commit is contained in:
@@ -12,6 +12,10 @@ type WriteEvent struct {
|
||||
Key *resourcepb.ResourceKey // the request key
|
||||
PreviousRV int64 // only for Update+Delete
|
||||
|
||||
// GUID is optional and might be used when persisting an event.
|
||||
// It is always set by the resource server.
|
||||
GUID string
|
||||
|
||||
// The json payload (without resourceVersion)
|
||||
Value []byte
|
||||
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
"go.opentelemetry.io/otel/trace/noop"
|
||||
@@ -18,6 +19,7 @@ import (
|
||||
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
||||
|
||||
claims "github.com/grafana/authlib/types"
|
||||
|
||||
"github.com/grafana/grafana/pkg/apimachinery/utils"
|
||||
"github.com/grafana/grafana/pkg/storage/unified/resourcepb"
|
||||
)
|
||||
@@ -393,6 +395,7 @@ func (s *server) newEvent(ctx context.Context, user claims.AuthInfo, key *resour
|
||||
Value: value,
|
||||
Key: key,
|
||||
Object: obj,
|
||||
GUID: uuid.New().String(),
|
||||
}
|
||||
|
||||
if oldValue == nil {
|
||||
@@ -670,6 +673,7 @@ func (s *server) Delete(ctx context.Context, req *resourcepb.DeleteRequest) (*re
|
||||
Key: req.Key,
|
||||
Type: resourcepb.WatchEvent_DELETED,
|
||||
PreviousRV: latest.ResourceVersion,
|
||||
GUID: uuid.New().String(),
|
||||
}
|
||||
requester, ok := claims.AuthInfoFrom(ctx)
|
||||
if !ok {
|
||||
|
||||
@@ -10,7 +10,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/go-sql-driver/mysql"
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgconn"
|
||||
"github.com/lib/pq"
|
||||
"github.com/mattn/go-sqlite3"
|
||||
@@ -329,7 +328,6 @@ func (b *backend) create(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
ctx, span := b.tracer.Start(ctx, tracePrefix+"Create")
|
||||
defer span.End()
|
||||
|
||||
guid := uuid.New().String()
|
||||
folder := ""
|
||||
if event.Object != nil {
|
||||
folder = event.Object.GetFolder()
|
||||
@@ -341,12 +339,12 @@ func (b *backend) create(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
SQLTemplate: sqltemplate.New(b.dialect),
|
||||
WriteEvent: event,
|
||||
Folder: folder,
|
||||
GUID: guid,
|
||||
GUID: event.GUID,
|
||||
}); err != nil {
|
||||
if IsRowAlreadyExistsError(err) {
|
||||
return guid, resource.ErrResourceAlreadyExists
|
||||
return event.GUID, resource.ErrResourceAlreadyExists
|
||||
}
|
||||
return guid, fmt.Errorf("insert into resource: %w", err)
|
||||
return event.GUID, fmt.Errorf("insert into resource: %w", err)
|
||||
}
|
||||
|
||||
// 2. Insert into resource history
|
||||
@@ -355,9 +353,9 @@ func (b *backend) create(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
WriteEvent: event,
|
||||
Folder: folder,
|
||||
Generation: event.Object.GetGeneration(),
|
||||
GUID: guid,
|
||||
GUID: event.GUID,
|
||||
}); err != nil {
|
||||
return guid, fmt.Errorf("insert into resource history: %w", err)
|
||||
return event.GUID, fmt.Errorf("insert into resource history: %w", err)
|
||||
}
|
||||
_ = b.historyPruner.Add(pruningKey{
|
||||
namespace: event.Key.Namespace,
|
||||
@@ -368,7 +366,7 @@ func (b *backend) create(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
if b.simulatedNetworkLatency > 0 {
|
||||
time.Sleep(b.simulatedNetworkLatency)
|
||||
}
|
||||
return guid, nil
|
||||
return event.GUID, nil
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
@@ -418,7 +416,7 @@ func IsRowAlreadyExistsError(err error) bool {
|
||||
func (b *backend) update(ctx context.Context, event resource.WriteEvent) (int64, error) {
|
||||
ctx, span := b.tracer.Start(ctx, tracePrefix+"Update")
|
||||
defer span.End()
|
||||
guid := uuid.New().String()
|
||||
|
||||
folder := ""
|
||||
if event.Object != nil {
|
||||
folder = event.Object.GetFolder()
|
||||
@@ -431,10 +429,10 @@ func (b *backend) update(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
SQLTemplate: sqltemplate.New(b.dialect),
|
||||
WriteEvent: event,
|
||||
Folder: folder,
|
||||
GUID: guid,
|
||||
GUID: event.GUID,
|
||||
})
|
||||
if err != nil {
|
||||
return guid, fmt.Errorf("resource update: %w", err)
|
||||
return event.GUID, fmt.Errorf("resource update: %w", err)
|
||||
}
|
||||
|
||||
// 2. Insert into resource history
|
||||
@@ -442,10 +440,10 @@ func (b *backend) update(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
SQLTemplate: sqltemplate.New(b.dialect),
|
||||
WriteEvent: event,
|
||||
Folder: folder,
|
||||
GUID: guid,
|
||||
GUID: event.GUID,
|
||||
Generation: event.Object.GetGeneration(),
|
||||
}); err != nil {
|
||||
return guid, fmt.Errorf("insert into resource history: %w", err)
|
||||
return event.GUID, fmt.Errorf("insert into resource history: %w", err)
|
||||
}
|
||||
_ = b.historyPruner.Add(pruningKey{
|
||||
namespace: event.Key.Namespace,
|
||||
@@ -453,7 +451,7 @@ func (b *backend) update(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
resource: event.Key.Resource,
|
||||
name: event.Key.Name,
|
||||
})
|
||||
return guid, nil
|
||||
return event.GUID, nil
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
@@ -475,7 +473,7 @@ func (b *backend) update(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
func (b *backend) delete(ctx context.Context, event resource.WriteEvent) (int64, error) {
|
||||
ctx, span := b.tracer.Start(ctx, tracePrefix+"Delete")
|
||||
defer span.End()
|
||||
guid := uuid.New().String()
|
||||
|
||||
folder := ""
|
||||
if event.Object != nil {
|
||||
folder = event.Object.GetFolder()
|
||||
@@ -485,10 +483,10 @@ func (b *backend) delete(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
_, err := dbutil.Exec(ctx, tx, sqlResourceDelete, sqlResourceRequest{
|
||||
SQLTemplate: sqltemplate.New(b.dialect),
|
||||
WriteEvent: event,
|
||||
GUID: guid,
|
||||
GUID: event.GUID,
|
||||
})
|
||||
if err != nil {
|
||||
return guid, fmt.Errorf("delete resource: %w", err)
|
||||
return event.GUID, fmt.Errorf("delete resource: %w", err)
|
||||
}
|
||||
|
||||
// 2. Add event to resource history
|
||||
@@ -496,10 +494,10 @@ func (b *backend) delete(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
SQLTemplate: sqltemplate.New(b.dialect),
|
||||
WriteEvent: event,
|
||||
Folder: folder,
|
||||
GUID: guid,
|
||||
GUID: event.GUID,
|
||||
Generation: 0, // object does not exist
|
||||
}); err != nil {
|
||||
return guid, fmt.Errorf("insert into resource history: %w", err)
|
||||
return event.GUID, fmt.Errorf("insert into resource history: %w", err)
|
||||
}
|
||||
_ = b.historyPruner.Add(pruningKey{
|
||||
namespace: event.Key.Namespace,
|
||||
@@ -507,7 +505,7 @@ func (b *backend) delete(ctx context.Context, event resource.WriteEvent) (int64,
|
||||
resource: event.Key.Resource,
|
||||
name: event.Key.Name,
|
||||
})
|
||||
return guid, nil
|
||||
return event.GUID, nil
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
|
||||
@@ -226,7 +226,7 @@ func (b *backend) processBulk(ctx context.Context, setting resource.BulkSettings
|
||||
PreviousRV: -1, // Used for WATCH, but we want to skip watch events
|
||||
},
|
||||
Folder: req.Folder,
|
||||
GUID: uuid.NewString(),
|
||||
GUID: uuid.New().String(),
|
||||
ResourceVersion: rv.next(obj),
|
||||
}); err != nil {
|
||||
return rollbackWithError(fmt.Errorf("insert into resource history: %w", err))
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
claims "github.com/grafana/authlib/types"
|
||||
@@ -87,6 +88,7 @@ func RunTestSearchAndStorage(t *testing.T, ctx context.Context, backend resource
|
||||
Key: key,
|
||||
Value: value,
|
||||
Object: meta,
|
||||
GUID: uuid.New().String(),
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv, int64(0))
|
||||
|
||||
@@ -1106,6 +1106,7 @@ func writeEvent(ctx context.Context, store resource.StorageBackend, name string,
|
||||
return store.WriteEvent(ctx, resource.WriteEvent{
|
||||
Type: action,
|
||||
Value: options.Value,
|
||||
GUID: uuid.New().String(),
|
||||
Key: &resourcepb.ResourceKey{
|
||||
Namespace: options.Namespace,
|
||||
Group: options.Group,
|
||||
|
||||
Reference in New Issue
Block a user