diff --git a/pkg/storage/unified/resource/event.go b/pkg/storage/unified/resource/event.go index 443e79d56a4..4a1ceb06a21 100644 --- a/pkg/storage/unified/resource/event.go +++ b/pkg/storage/unified/resource/event.go @@ -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 diff --git a/pkg/storage/unified/resource/server.go b/pkg/storage/unified/resource/server.go index fc19fa9a161..c340c7a6e39 100644 --- a/pkg/storage/unified/resource/server.go +++ b/pkg/storage/unified/resource/server.go @@ -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 { diff --git a/pkg/storage/unified/sql/backend.go b/pkg/storage/unified/sql/backend.go index 720f4f24322..51d77335c4a 100644 --- a/pkg/storage/unified/sql/backend.go +++ b/pkg/storage/unified/sql/backend.go @@ -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 { diff --git a/pkg/storage/unified/sql/bulk.go b/pkg/storage/unified/sql/bulk.go index 401017d435f..dfb1aba9a3f 100644 --- a/pkg/storage/unified/sql/bulk.go +++ b/pkg/storage/unified/sql/bulk.go @@ -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)) diff --git a/pkg/storage/unified/testing/search_and_storage.go b/pkg/storage/unified/testing/search_and_storage.go index 015a499bd12..2bb1c70a3e7 100644 --- a/pkg/storage/unified/testing/search_and_storage.go +++ b/pkg/storage/unified/testing/search_and_storage.go @@ -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)) diff --git a/pkg/storage/unified/testing/storage_backend.go b/pkg/storage/unified/testing/storage_backend.go index 6295f5fdb90..c8ebb8b4692 100644 --- a/pkg/storage/unified/testing/storage_backend.go +++ b/pkg/storage/unified/testing/storage_backend.go @@ -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,