From 46c38fdbb7030537be90587d9cba2c4cfdc081d0 Mon Sep 17 00:00:00 2001 From: Dana Axinte <53751979+dana-axinte@users.noreply.github.com> Date: Fri, 4 Jul 2025 13:13:48 +0100 Subject: [PATCH] SecretsManager: Introduce worker and secret async service (#107614) SecretsManager: Introduce worker and secret aysnc service Co-authored-by: PoorlyDefinedBehaviour Co-authored-by: Matheus Macabu Co-authored-by: Michael Mandrus --- .../apis/secret/service/secure_value.go | 268 +++++++++++++++++ .../apis/secret/testutils/testutils.go | 202 +++++++++++++ pkg/registry/apis/secret/tracectx/carrier.go | 57 ++++ .../apis/secret/tracectx/carrier_test.go | 88 ++++++ pkg/registry/apis/secret/worker/metrics.go | 44 +++ pkg/registry/apis/secret/worker/worker.go | 274 ++++++++++++++++++ .../apis/secret/worker/worker_test.go | 246 ++++++++++++++++ .../backgroundsvcs/background_services.go | 3 + pkg/server/wire.go | 9 +- pkg/server/wire_gen.go | 79 ++++- pkg/storage/secret/metadata/outbox_store.go | 4 +- 11 files changed, 1264 insertions(+), 10 deletions(-) create mode 100644 pkg/registry/apis/secret/service/secure_value.go create mode 100644 pkg/registry/apis/secret/testutils/testutils.go create mode 100644 pkg/registry/apis/secret/tracectx/carrier.go create mode 100644 pkg/registry/apis/secret/tracectx/carrier_test.go create mode 100644 pkg/registry/apis/secret/worker/metrics.go create mode 100644 pkg/registry/apis/secret/worker/worker.go create mode 100644 pkg/registry/apis/secret/worker/worker_test.go diff --git a/pkg/registry/apis/secret/service/secure_value.go b/pkg/registry/apis/secret/service/secure_value.go new file mode 100644 index 00000000000..e667ad4c4be --- /dev/null +++ b/pkg/registry/apis/secret/service/secure_value.go @@ -0,0 +1,268 @@ +package service + +import ( + "context" + "fmt" + + claims "github.com/grafana/authlib/types" + "github.com/grafana/grafana/pkg/apimachinery/utils" + secretv0alpha1 "github.com/grafana/grafana/pkg/apis/secret/v0alpha1" + "github.com/grafana/grafana/pkg/registry/apis/secret/contracts" + "github.com/grafana/grafana/pkg/registry/apis/secret/tracectx" + "github.com/grafana/grafana/pkg/registry/apis/secret/xkube" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" +) + +type SecureValueService struct { + tracer trace.Tracer + accessClient claims.AccessClient + database contracts.Database + secureValueMetadataStorage contracts.SecureValueMetadataStorage + outboxQueue contracts.OutboxQueue + encryptionManager contracts.EncryptionManager +} + +func ProvideSecureValueService( + tracer trace.Tracer, + accessClient claims.AccessClient, + database contracts.Database, + secureValueMetadataStorage contracts.SecureValueMetadataStorage, + outboxQueue contracts.OutboxQueue, + encryptionManager contracts.EncryptionManager, +) *SecureValueService { + return &SecureValueService{ + tracer: tracer, + accessClient: accessClient, + database: database, + secureValueMetadataStorage: secureValueMetadataStorage, + outboxQueue: outboxQueue, + encryptionManager: encryptionManager, + } +} + +func (s *SecureValueService) Create(ctx context.Context, sv *secretv0alpha1.SecureValue, actorUID string) (*secretv0alpha1.SecureValue, error) { + ctx, span := s.tracer.Start(ctx, "SecureValueService.Create", trace.WithAttributes( + attribute.String("name", sv.GetName()), + attribute.String("namespace", sv.GetNamespace()), + attribute.String("actor", actorUID), + )) + defer span.End() + + sv.Status = secretv0alpha1.SecureValueStatus{Phase: secretv0alpha1.SecureValuePhasePending, Message: "Creating secure value"} + + var out *secretv0alpha1.SecureValue + + encryptedSecret, err := s.encryptionManager.Encrypt(ctx, sv.Namespace, []byte(sv.Spec.Value.DangerouslyExposeAndConsumeValue())) + if err != nil { + return nil, fmt.Errorf("encrypting secure value secret: %w", err) + } + + // Specifically here so that the spans from the worker are not inside the transaction. + requestID := tracectx.HexEncodeTraceFromContext(ctx) + + if err := s.database.Transaction(ctx, func(ctx context.Context) error { + createdSecureValue, err := s.secureValueMetadataStorage.Create(ctx, sv, actorUID) + if err != nil { + return fmt.Errorf("failed to create securevalue: %w", err) + } + out = createdSecureValue + + if _, err := s.outboxQueue.Append(ctx, contracts.AppendOutboxMessage{ + RequestID: requestID, + Type: contracts.CreateSecretOutboxMessage, + Name: sv.Name, + Namespace: sv.Namespace, + EncryptedSecret: string(encryptedSecret), + KeeperName: sv.Spec.Keeper, + }); err != nil { + return fmt.Errorf("failed to append message to create secure value to outbox queue: %w", err) + } + + return nil + }); err != nil { + return out, err + } + + return out, nil +} + +func (s *SecureValueService) Read(ctx context.Context, namespace xkube.Namespace, name string) (*secretv0alpha1.SecureValue, error) { + ctx, span := s.tracer.Start(ctx, "SecureValueService.Read", trace.WithAttributes( + attribute.String("name", name), + attribute.String("namespace", namespace.String()), + )) + defer span.End() + + return s.secureValueMetadataStorage.Read(ctx, namespace, name, contracts.ReadOpts{ForUpdate: false}) +} + +func (s *SecureValueService) List(ctx context.Context, namespace xkube.Namespace) (*secretv0alpha1.SecureValueList, error) { + ctx, span := s.tracer.Start(ctx, "SecureValueService.List", trace.WithAttributes( + attribute.String("namespace", namespace.String()), + )) + defer span.End() + + user, ok := claims.AuthInfoFrom(ctx) + if !ok { + return nil, fmt.Errorf("missing auth info in context") + } + + hasPermissionFor, err := s.accessClient.Compile(ctx, user, claims.ListRequest{ + Group: secretv0alpha1.GROUP, + Resource: secretv0alpha1.SecureValuesResourceInfo.GetName(), + Namespace: namespace.String(), + Verb: utils.VerbGet, // Why not VerbList? + }) + if err != nil { + return nil, fmt.Errorf("failed to compile checker: %w", err) + } + + secureValuesMetadata, err := s.secureValueMetadataStorage.List(ctx, namespace) + if err != nil { + return nil, fmt.Errorf("fetching secure values from storage: %+w", err) + } + + out := make([]secretv0alpha1.SecureValue, 0) + + for _, metadata := range secureValuesMetadata { + // Check whether the user has permission to access this specific SecureValue in the namespace. + if !hasPermissionFor(metadata.Name, "") { + continue + } + + out = append(out, metadata) + } + + return &secretv0alpha1.SecureValueList{ + Items: out, + }, nil +} + +func (s *SecureValueService) Update(ctx context.Context, newSecureValue *secretv0alpha1.SecureValue, actorUID string) (*secretv0alpha1.SecureValue, bool, error) { + ctx, span := s.tracer.Start(ctx, "SecureValueService.Create", trace.WithAttributes( + attribute.String("name", newSecureValue.GetName()), + attribute.String("namespace", newSecureValue.GetNamespace()), + attribute.String("actor", actorUID), + )) + defer span.End() + + // True when the effects of an update can be seen immediately. + // Never true in this case since updating a secure value is async. + const updateIsSync = false + + var ( + out *secretv0alpha1.SecureValue + encryptedSecret string + ) + + if newSecureValue.Spec.Value != "" { + buffer, err := s.encryptionManager.Encrypt(ctx, newSecureValue.Namespace, []byte(newSecureValue.Spec.Value.DangerouslyExposeAndConsumeValue())) + if err != nil { + return nil, false, fmt.Errorf("encrypting secure value secret: %w", err) + } + encryptedSecret = string(buffer) + } + + // Especifically here so that the spans from the worker are not inside the transaction. + requestID := tracectx.HexEncodeTraceFromContext(ctx) + + if err := s.database.Transaction(ctx, func(ctx context.Context) error { + sv, err := s.secureValueMetadataStorage.Read(ctx, xkube.Namespace(newSecureValue.Namespace), newSecureValue.Name, contracts.ReadOpts{ForUpdate: true}) + if err != nil { + return fmt.Errorf("fetching secure value: %+w", err) + } + + if sv.Status.Phase == secretv0alpha1.SecureValuePhasePending { + return contracts.ErrSecureValueOperationInProgress + } + + // Succeed immediately if the value is not going to be updated + if encryptedSecret == "" { + newSecureValue.Status = secretv0alpha1.SecureValueStatus{Phase: secretv0alpha1.SecureValuePhaseSucceeded} + } else { + newSecureValue.Status = secretv0alpha1.SecureValueStatus{ + Message: "Updating secure value", + Phase: secretv0alpha1.SecureValuePhasePending, + } + } + + // Current implementation replaces everything passed in the spec, so it is not a PATCH. Do we want/need to support that? + updatedSecureValue, err := s.secureValueMetadataStorage.Update(ctx, newSecureValue, actorUID) + if err != nil { + return fmt.Errorf("failed to update secure value: %w", err) + } + out = updatedSecureValue + + // Only the value needs to be updated asynchronously by the outbox worker + if encryptedSecret != "" { + if _, err := s.outboxQueue.Append(ctx, contracts.AppendOutboxMessage{ + RequestID: requestID, + Type: contracts.UpdateSecretOutboxMessage, + Name: newSecureValue.Name, + Namespace: newSecureValue.Namespace, + EncryptedSecret: encryptedSecret, + KeeperName: newSecureValue.Spec.Keeper, + ExternalID: &updatedSecureValue.Status.ExternalID, + }); err != nil { + return fmt.Errorf("failed to append message to update secure value to outbox queue: %w", err) + } + } + + return nil + }); err != nil { + return out, updateIsSync, err + } + + return out, updateIsSync, nil +} + +func (s *SecureValueService) Delete(ctx context.Context, namespace xkube.Namespace, name string) (*secretv0alpha1.SecureValue, error) { + ctx, span := s.tracer.Start(ctx, "SecureValueService.Delete", trace.WithAttributes( + attribute.String("name", name), + attribute.String("namespace", namespace.String()), + )) + defer span.End() + + // Set inside of the transaction callback + var out *secretv0alpha1.SecureValue + + // Especifically here so that the spans from the worker are not inside the transaction. + requestID := tracectx.HexEncodeTraceFromContext(ctx) + + if err := s.database.Transaction(ctx, func(ctx context.Context) error { + sv, err := s.secureValueMetadataStorage.Read(ctx, namespace, name, contracts.ReadOpts{ForUpdate: true}) + if err != nil { + return fmt.Errorf("fetching secure value: %+w", err) + } + + if sv.Status.Phase == secretv0alpha1.SecureValuePhasePending { + return contracts.ErrSecureValueOperationInProgress + } + + sv.Status = secretv0alpha1.SecureValueStatus{Phase: secretv0alpha1.SecureValuePhasePending, Message: "Deleting secure value"} + + if err := s.secureValueMetadataStorage.SetStatus(ctx, namespace, name, sv.Status); err != nil { + return fmt.Errorf("setting secure value status phase: %+w", err) + } + + if _, err := s.outboxQueue.Append(ctx, contracts.AppendOutboxMessage{ + RequestID: requestID, + Type: contracts.DeleteSecretOutboxMessage, + Name: name, + Namespace: namespace.String(), + KeeperName: sv.Spec.Keeper, + ExternalID: &sv.Status.ExternalID, + }); err != nil { + return fmt.Errorf("appending delete secure value message to outbox queue: %+w", err) + } + + out = sv + + return nil + }); err != nil { + return out, err + } + + return out, nil +} diff --git a/pkg/registry/apis/secret/testutils/testutils.go b/pkg/registry/apis/secret/testutils/testutils.go new file mode 100644 index 00000000000..40e571aca63 --- /dev/null +++ b/pkg/registry/apis/secret/testutils/testutils.go @@ -0,0 +1,202 @@ +package testutils + +import ( + "context" + "testing" + "time" + + secretv0alpha1 "github.com/grafana/grafana/pkg/apis/secret/v0alpha1" + encryptionstorage "github.com/grafana/grafana/pkg/storage/secret/encryption" + "go.opentelemetry.io/otel/trace/noop" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + "github.com/grafana/grafana/pkg/infra/usagestats" + "github.com/grafana/grafana/pkg/registry/apis/secret/contracts" + "github.com/grafana/grafana/pkg/registry/apis/secret/encryption" + "github.com/grafana/grafana/pkg/registry/apis/secret/encryption/manager" + "github.com/grafana/grafana/pkg/registry/apis/secret/secretkeeper/sqlkeeper" + "github.com/grafana/grafana/pkg/registry/apis/secret/service" + "github.com/grafana/grafana/pkg/registry/apis/secret/worker" + "github.com/grafana/grafana/pkg/registry/apis/secret/xkube" + "github.com/grafana/grafana/pkg/services/accesscontrol" + "github.com/grafana/grafana/pkg/services/accesscontrol/actest" + "github.com/grafana/grafana/pkg/services/featuremgmt" + "github.com/grafana/grafana/pkg/services/sqlstore" + "github.com/grafana/grafana/pkg/setting" + "github.com/grafana/grafana/pkg/storage/secret/database" + "github.com/grafana/grafana/pkg/storage/secret/metadata" + "github.com/grafana/grafana/pkg/storage/secret/migrator" + "github.com/stretchr/testify/require" +) + +type setupConfig struct { + workerCfg worker.Config + keeperService contracts.KeeperService +} + +func defaultSetupCfg() setupConfig { + return setupConfig{ + workerCfg: worker.Config{ + BatchSize: 10, + ReceiveTimeout: 1 * time.Second, + PollingInterval: time.Millisecond, + MaxMessageProcessingAttempts: 5, + }, + } +} + +func WithWorkerConfig(cfg worker.Config) func(*setupConfig) { + return func(setupCfg *setupConfig) { + setupCfg.workerCfg = cfg + } +} + +func WithKeeperService(keeperService contracts.KeeperService) func(*setupConfig) { + return func(setupCfg *setupConfig) { + setupCfg.keeperService = keeperService + } +} + +func Setup(t *testing.T, opts ...func(*setupConfig)) Sut { + setupCfg := defaultSetupCfg() + for _, opt := range opts { + opt(&setupCfg) + } + + tracer := noop.NewTracerProvider().Tracer("test") + testDB := sqlstore.NewTestStore(t, sqlstore.WithMigrator(migrator.New())) + + database := database.ProvideDatabase(testDB, tracer) + + outboxQueue := metadata.ProvideOutboxQueue(database, tracer, nil) + + features := featuremgmt.WithFeatures(featuremgmt.FlagGrafanaAPIServerWithExperimentalAPIs, featuremgmt.FlagSecretsManagementAppPlatform) + + keeperMetadataStorage, err := metadata.ProvideKeeperMetadataStorage(database, tracer, features, nil) + require.NoError(t, err) + + secureValueMetadataStorage, err := metadata.ProvideSecureValueMetadataStorage(database, tracer, features, nil) + require.NoError(t, err) + + // Initialize access client + access control + accessControl := &actest.FakeAccessControl{ExpectedEvaluate: true} + accessClient := accesscontrol.NewLegacyAccessClient(accessControl) + + defaultKey := "SdlklWklckeLS" + cfg := &setting.Cfg{ + SecretsManagement: setting.SecretsManagerSettings{ + SecretKey: defaultKey, + EncryptionProvider: "secretKey.v1", + }, + } + store, err := encryptionstorage.ProvideDataKeyStorage(database, tracer, features, nil) + require.NoError(t, err) + + usageStats := &usagestats.UsageStatsMock{T: t} + + encryptionManager, err := manager.ProvideEncryptionManager( + tracer, + store, + cfg, + usageStats, + encryption.ProvideThirdPartyProviderMap(), + ) + require.NoError(t, err) + + // Initialize encrypted value storage with a fake db + encValueStore, err := encryptionstorage.ProvideEncryptedValueStorage(database, tracer, features) + require.NoError(t, err) + + sqlKeeper := sqlkeeper.NewSQLKeeper(tracer, encryptionManager, encValueStore, nil) + + var keeperService contracts.KeeperService = newKeeperServiceWrapper(sqlKeeper) + + if setupCfg.keeperService != nil { + keeperService = setupCfg.keeperService + } + + secureValueService := service.ProvideSecureValueService(tracer, accessClient, database, secureValueMetadataStorage, outboxQueue, encryptionManager) + + worker, err := worker.NewWorker( + setupCfg.workerCfg, + tracer, + database, + outboxQueue, + secureValueMetadataStorage, + keeperMetadataStorage, + keeperService, + encryptionManager, + features, + nil, // metrics + ) + require.NoError(t, err) + + return Sut{Worker: worker, SecureValueService: secureValueService, SecureValueMetadataStorage: secureValueMetadataStorage, OutboxQueue: outboxQueue, Database: database} +} + +type Sut struct { + Worker *worker.Worker + SecureValueService *service.SecureValueService + SecureValueMetadataStorage contracts.SecureValueMetadataStorage + OutboxQueue contracts.OutboxQueue + Database *database.Database +} + +type CreateSvConfig struct { + Sv *secretv0alpha1.SecureValue +} + +func CreateSvWithSv(sv *secretv0alpha1.SecureValue) func(*CreateSvConfig) { + return func(cfg *CreateSvConfig) { + cfg.Sv = sv + } +} + +func (s *Sut) CreateSv(ctx context.Context, opts ...func(*CreateSvConfig)) (*secretv0alpha1.SecureValue, error) { + cfg := CreateSvConfig{ + Sv: &secretv0alpha1.SecureValue{ + ObjectMeta: metav1.ObjectMeta{ + Name: "sv1", + Namespace: "ns1", + }, + Spec: secretv0alpha1.SecureValueSpec{ + Description: "desc1", + Value: secretv0alpha1.NewExposedSecureValue("v1"), + }, + Status: secretv0alpha1.SecureValueStatus{ + Phase: secretv0alpha1.SecureValuePhasePending, + }, + }, + } + for _, opt := range opts { + opt(&cfg) + } + + createdSv, err := s.SecureValueService.Create(ctx, cfg.Sv, "actor") + if err != nil { + return nil, err + } + return createdSv, nil +} + +func (s *Sut) UpdateSv(ctx context.Context, sv *secretv0alpha1.SecureValue) (*secretv0alpha1.SecureValue, error) { + newSv, _, err := s.SecureValueService.Update(ctx, sv, "actor") + return newSv, err +} + +func (s *Sut) DeleteSv(ctx context.Context, namespace, name string) (*secretv0alpha1.SecureValue, error) { + sv, err := s.SecureValueService.Delete(ctx, xkube.Namespace(namespace), name) + return sv, err +} + +type keeperServiceWrapper struct { + keeper contracts.Keeper +} + +func newKeeperServiceWrapper(keeper contracts.Keeper) *keeperServiceWrapper { + return &keeperServiceWrapper{keeper: keeper} +} + +func (wrapper *keeperServiceWrapper) KeeperForConfig(cfg secretv0alpha1.KeeperConfig) (contracts.Keeper, error) { + return wrapper.keeper, nil +} diff --git a/pkg/registry/apis/secret/tracectx/carrier.go b/pkg/registry/apis/secret/tracectx/carrier.go new file mode 100644 index 00000000000..a8153c42d1e --- /dev/null +++ b/pkg/registry/apis/secret/tracectx/carrier.go @@ -0,0 +1,57 @@ +package tracectx + +import ( + "context" + "encoding/hex" + "fmt" + "strings" + + "go.opentelemetry.io/otel/propagation" +) + +const ( + kvSeparator = "=" + pairSeparator = "#" +) + +func HexEncodeTraceFromContext(ctx context.Context) string { + carrier := propagation.MapCarrier(make(map[string]string)) + + propagation.TraceContext{}.Inject(ctx, carrier) + + // no trace in context + if len(carrier) == 0 { + return "" + } + + pairs := make([]string, 0, len(carrier)) + for k, v := range carrier { + pairs = append(pairs, k+kvSeparator+v) + } + + return hex.EncodeToString([]byte(strings.Join(pairs, pairSeparator))) +} + +func HexDecodeTraceIntoContext(ctx context.Context, encoded string) (context.Context, error) { + if encoded == "" { + return ctx, nil + } + + decoded, err := hex.DecodeString(encoded) + if err != nil { + return nil, err + } + + pairs := strings.Split(string(decoded), pairSeparator) + + carrier := make(propagation.MapCarrier, len(pairs)) + for _, pair := range pairs { + kv := strings.SplitN(pair, kvSeparator, 2) + if len(kv) != 2 || kv[0] == "" || kv[1] == "" { + return nil, fmt.Errorf("invalid key-value pair: %s", pair) + } + carrier[kv[0]] = kv[1] + } + + return propagation.TraceContext{}.Extract(ctx, carrier), nil +} diff --git a/pkg/registry/apis/secret/tracectx/carrier_test.go b/pkg/registry/apis/secret/tracectx/carrier_test.go new file mode 100644 index 00000000000..043428faeb2 --- /dev/null +++ b/pkg/registry/apis/secret/tracectx/carrier_test.go @@ -0,0 +1,88 @@ +package tracectx + +import ( + "context" + "encoding/hex" + "testing" + + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/propagation" + "go.opentelemetry.io/otel/trace" +) + +func TestHexEncodeTraceFromContext(t *testing.T) { + t.Run("when no trace is present in context, it returns empty string", func(t *testing.T) { + ctx := context.Background() + + encoded := HexEncodeTraceFromContext(ctx) + require.Empty(t, encoded) + }) + + t.Run("when trace is present in context, it returns hex-encoded string", func(t *testing.T) { + carrier := propagation.MapCarrier{ + "traceparent": "00-446e31681d64f9dcefd947c95ef321d0-009e2f3d8ded1892-01", + "tracestate": "first=abc1234,second=xyz7890", + } + ctx := propagation.TraceContext{}.Extract(context.Background(), carrier) + + encoded := HexEncodeTraceFromContext(ctx) + require.NotEmpty(t, encoded) + + traceCtx, err := HexDecodeTraceIntoContext(context.Background(), encoded) + require.NoError(t, err) + + span := trace.SpanFromContext(traceCtx) + require.True(t, span.SpanContext().IsValid()) + + carrier = propagation.MapCarrier(make(map[string]string)) + propagation.TraceContext{}.Inject(traceCtx, carrier) + require.Contains(t, carrier, "traceparent") + require.Contains(t, carrier, "tracestate") + }) +} + +func TestHexDecodeTraceIntoContext(t *testing.T) { + t.Run("when encoded string is empty, it returns original context", func(t *testing.T) { + ctx := context.Background() + + result, err := HexDecodeTraceIntoContext(ctx, "") + require.NoError(t, err) + require.Equal(t, ctx, result) + }) + + t.Run("when encoded string is valid hex, it returns context with trace", func(t *testing.T) { + encoded := hex.EncodeToString([]byte("traceparent=00-446e31681d64f9dcefd947c95ef321d0-009e2f3d8ded1892-01#tracestate=first=abc1234,second=xyz7890")) + + ctx, err := HexDecodeTraceIntoContext(context.Background(), encoded) + require.NoError(t, err) + + span := trace.SpanFromContext(ctx) + require.True(t, span.SpanContext().IsValid()) + }) + + t.Run("when encoded string has invalid hex encoding, it returns an error", func(t *testing.T) { + invalidHex := "invalid-hex-zzz" + + result, err := HexDecodeTraceIntoContext(context.Background(), invalidHex) + require.Error(t, err) + require.Nil(t, result) + }) + + t.Run("when decoded string has invalid key-value pair format, it returns an error", func(t *testing.T) { + // missing key + encoded := hex.EncodeToString([]byte("00-446e31681d64f9dcefd947c95ef321d0-009e2f3d8ded1892-01")) + + result, err := HexDecodeTraceIntoContext(context.Background(), encoded) + require.Error(t, err) + require.Nil(t, result) + }) + + t.Run("when decoded string has key without value, it returns error", func(t *testing.T) { + // missing value + encoded := hex.EncodeToString([]byte("traceparent=")) + + result, err := HexDecodeTraceIntoContext(context.Background(), encoded) + require.Error(t, err) + require.Nil(t, result) + }) +} diff --git a/pkg/registry/apis/secret/worker/metrics.go b/pkg/registry/apis/secret/worker/metrics.go new file mode 100644 index 00000000000..a6f690f7163 --- /dev/null +++ b/pkg/registry/apis/secret/worker/metrics.go @@ -0,0 +1,44 @@ +package worker + +import ( + "github.com/prometheus/client_golang/prometheus" +) + +const ( + namespace = "grafana_secrets_manager" + subsystem = "outbox_worker" +) + +// OutboxMetrics is a struct that contains all the metrics for an implementation of the secrets service. +type OutboxMetrics struct { + OutboxMessageProcessingDuration *prometheus.HistogramVec +} + +func newOutboxMetrics() *OutboxMetrics { + return &OutboxMetrics{ + OutboxMessageProcessingDuration: prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Namespace: namespace, + Subsystem: subsystem, + Name: "message_processing_duration_seconds", + Help: "Duration of outbox message processing", + Buckets: prometheus.DefBuckets, + }, []string{"message_type", "keeper_type"}), + } +} + +// NewOutboxMetrics creates a new SecretsMetrics struct containing registered metrics +func NewOutboxMetrics(reg prometheus.Registerer) *OutboxMetrics { + m := newOutboxMetrics() + + if reg != nil { + reg.MustRegister( + m.OutboxMessageProcessingDuration, + ) + } + + return m +} + +func NewTestMetrics() *OutboxMetrics { + return newOutboxMetrics() +} diff --git a/pkg/registry/apis/secret/worker/worker.go b/pkg/registry/apis/secret/worker/worker.go new file mode 100644 index 00000000000..35474fdae47 --- /dev/null +++ b/pkg/registry/apis/secret/worker/worker.go @@ -0,0 +1,274 @@ +package worker + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/grafana/grafana-app-sdk/logging" + secretv0alpha1 "github.com/grafana/grafana/pkg/apis/secret/v0alpha1" + "github.com/grafana/grafana/pkg/registry" + "github.com/grafana/grafana/pkg/registry/apis/secret/contracts" + "github.com/grafana/grafana/pkg/registry/apis/secret/tracectx" + "github.com/grafana/grafana/pkg/registry/apis/secret/xkube" + "github.com/grafana/grafana/pkg/services/featuremgmt" + "github.com/prometheus/client_golang/prometheus" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" +) + +// Consumes and processes messages from the secure value outbox queue +type Worker struct { + config Config + tracer trace.Tracer + database contracts.Database + outboxQueue contracts.OutboxQueue + secureValueMetadataStorage contracts.SecureValueMetadataStorage + keeperMetadataStorage contracts.KeeperMetadataStorage + keeperService contracts.KeeperService + encryptionManager contracts.EncryptionManager + metrics *OutboxMetrics + enabled bool +} + +// DefaultConfig for the secure value outbox worker. +var DefaultConfig = Config{ + BatchSize: 20, + ReceiveTimeout: 5 * time.Second, + PollingInterval: 100 * time.Millisecond, + MaxMessageProcessingAttempts: 10, +} + +// ProvideWorkerConfig used for wire. +func ProvideWorkerConfig() Config { + return DefaultConfig +} + +type Config struct { + // The max number of messages to fetch from the outbox queue in a batch + BatchSize uint + // How long to wait for a request to fetch messages from the outbox queue + ReceiveTimeout time.Duration + // How often to poll the outbox queue for new messages + PollingInterval time.Duration + // How many tries to try to process a message before marking the operation as failed + MaxMessageProcessingAttempts uint +} + +func NewWorker( + config Config, + tracer trace.Tracer, + database contracts.Database, + outboxQueue contracts.OutboxQueue, + secureValueMetadataStorage contracts.SecureValueMetadataStorage, + keeperMetadataStorage contracts.KeeperMetadataStorage, + keeperService contracts.KeeperService, + encryptionManager contracts.EncryptionManager, + features featuremgmt.FeatureToggles, + reg prometheus.Registerer, +) (*Worker, error) { + if config.BatchSize == 0 { + return nil, fmt.Errorf("config.BatchSize is required") + } + if config.ReceiveTimeout == 0 { + return nil, fmt.Errorf("config.ReceiveTimeout is required") + } + if config.PollingInterval == 0 { + return nil, fmt.Errorf("config.PollingInterval is required") + } + if config.MaxMessageProcessingAttempts == 0 { + return nil, fmt.Errorf("config.MaxMessageProcessingAttempts is required") + } + + // Require both features to be enabled for the worker to run. + enabled := features.IsEnabledGlobally(featuremgmt.FlagGrafanaAPIServerWithExperimentalAPIs) && features.IsEnabledGlobally(featuremgmt.FlagSecretsManagementAppPlatform) + + return &Worker{ + config: config, + tracer: tracer, + database: database, + outboxQueue: outboxQueue, + secureValueMetadataStorage: secureValueMetadataStorage, + keeperMetadataStorage: keeperMetadataStorage, + keeperService: keeperService, + encryptionManager: encryptionManager, + metrics: NewOutboxMetrics(reg), + enabled: enabled, + }, nil +} + +// Ensure that Worker implements the BackgroundService interface, so we can start it as a background service. +var _ registry.BackgroundService = (*Worker)(nil) + +// Run is the main method to drive the worker +func (w *Worker) Run(ctx context.Context) error { + if !w.enabled { + return nil + } + + logging.FromContext(ctx).Debug("starting worker control loop") + + t := time.NewTicker(w.config.PollingInterval) + defer t.Stop() + + for { + select { + // If the context was canceled + case <-ctx.Done(): + // return the reason it was canceled + return ctx.Err() + + // Otherwise try to receive messages + case <-t.C: + if ctx.Err() != nil { + return ctx.Err() + } + + if err := w.ReceiveAndProcessMessages(ctx); err != nil { + logging.FromContext(ctx).Error("receiving outbox messages", "err", err.Error()) + } + } + } +} + +// TODO: don't rollback every message when a single error happens +func (w *Worker) ReceiveAndProcessMessages(ctx context.Context) error { + messageIDs := make([]int64, 0) + + txErr := w.database.Transaction(ctx, func(ctx context.Context) error { + timeoutCtx, cancel := context.WithTimeout(ctx, w.config.ReceiveTimeout) + messages, err := w.outboxQueue.ReceiveN(timeoutCtx, w.config.BatchSize) + cancel() + if err != nil { + return err + } + + for _, message := range messages { + messageIDs = append(messageIDs, message.MessageID) + if err := w.processMessage(ctx, message); err != nil { + return fmt.Errorf("processing message: %+v %w", message, err) + } + } + return nil + }) + + // This call is made outside the transaction to make sure the receive count is updated on rollbacks. + incrementErr := w.outboxQueue.IncrementReceiveCount(ctx, messageIDs) + if incrementErr != nil { + incrementErr = fmt.Errorf("incrementing receive count for outbox message: %w", incrementErr) + } + + return errors.Join(txErr, incrementErr) +} + +func (w *Worker) processMessage(ctx context.Context, message contracts.OutboxMessage) error { + start := time.Now() + keeperType := "unknown" + defer func() { + w.metrics.OutboxMessageProcessingDuration.WithLabelValues(string(message.Type), keeperType).Observe(time.Since(start).Seconds()) + }() + logging.FromContext(ctx).Debug("processing message", "type", message.Type, "name", message.Name, "namespace", message.Namespace, "receiveCount", message.ReceiveCount) + + opts := []trace.SpanStartOption{} + // If there's no request ID in the message, start a new root span and log an error. + ctx, err := tracectx.HexDecodeTraceIntoContext(ctx, message.RequestID) + if err != nil { + opts = append(opts, trace.WithNewRoot()) + logging.FromContext(ctx).Error("decoding trace context from message", "err", err.Error(), "message.requestID", message.RequestID) + } + + opts = append(opts, trace.WithAttributes( + attribute.String("message.requestID", message.RequestID), + attribute.Int64("message.id", message.MessageID), + attribute.String("message.type", string(message.Type)), + attribute.String("message.namespace", message.Namespace), + attribute.String("message.secureValue.name", message.Name), + attribute.Int("message.receive.count", message.ReceiveCount), + )) + + ctx, span := w.tracer.Start(ctx, "Worker.ProcessMessage", opts...) + defer span.End() + + if message.ReceiveCount >= int(w.config.MaxMessageProcessingAttempts) { + if err := w.secureValueMetadataStorage.SetStatus(ctx, xkube.Namespace(message.Namespace), message.Name, secretv0alpha1.SecureValueStatus{Phase: secretv0alpha1.SecureValuePhaseFailed, Message: fmt.Sprintf("Reached max number of attempts to complete operation: %s", message.Type)}); err != nil { + return fmt.Errorf("setting secret metadata status to Succeeded: message=%+v", message) + } + if err := w.outboxQueue.Delete(ctx, message.MessageID); err != nil { + return fmt.Errorf("deleting message from outbox queue: %w", err) + } + return nil + } + + keeperCfg, err := w.keeperMetadataStorage.GetKeeperConfig(ctx, message.Namespace, message.KeeperName, contracts.ReadOpts{ForUpdate: true}) + if err != nil { + return fmt.Errorf("fetching keeper config: namespace=%+v keeperName=%+v %w", message.Namespace, message.KeeperName, err) + } + keeperType = string(keeperCfg.Type()) + + keeper, err := w.keeperService.KeeperForConfig(keeperCfg) + if err != nil { + return fmt.Errorf("getting keeper for config: namespace=%+v keeperName=%+v %w", message.Namespace, message.KeeperName, err) + } + logging.FromContext(ctx).Debug("retrieved keeper", "namespace", message.Namespace, "keeperName", message.KeeperName, "type", keeperCfg.Type()) + + switch message.Type { + case contracts.CreateSecretOutboxMessage: + rawSecret, err := w.encryptionManager.Decrypt(ctx, message.Namespace, []byte(message.EncryptedSecret)) + if err != nil { + return fmt.Errorf("decrypting secure value secret: %w", err) + } + + externalID, err := keeper.Store(ctx, keeperCfg, message.Namespace, string(rawSecret)) + if err != nil { + return fmt.Errorf("storing secret: message=%+v %w", message, err) + } + + if err := w.secureValueMetadataStorage.SetExternalID(ctx, xkube.Namespace(message.Namespace), message.Name, externalID); err != nil { + return fmt.Errorf("setting secret metadata externalID: externalID=%+v message=%+v %w", externalID, message, err) + } + + // Setting the status to Succeeded must be the last action + // since it acts as a fence to clients. + if err := w.secureValueMetadataStorage.SetStatus(ctx, xkube.Namespace(message.Namespace), message.Name, secretv0alpha1.SecureValueStatus{Phase: secretv0alpha1.SecureValuePhaseSucceeded}); err != nil { + return fmt.Errorf("setting secret metadata status to Succeeded: message=%+v %w", message, err) + } + + case contracts.UpdateSecretOutboxMessage: + rawSecret, err := w.encryptionManager.Decrypt(ctx, message.Namespace, []byte(message.EncryptedSecret)) + if err != nil { + return fmt.Errorf("decrypting secure value secret: %w", err) + } + + if err := keeper.Update(ctx, keeperCfg, message.Namespace, contracts.ExternalID(*message.ExternalID), string(rawSecret)); err != nil { + return fmt.Errorf("calling keeper to update secret: %w", err) + } + + // Setting the status to Succeeded must be the last action + // since it acts as a fence to clients. + if err := w.secureValueMetadataStorage.SetStatus(ctx, xkube.Namespace(message.Namespace), message.Name, secretv0alpha1.SecureValueStatus{Phase: secretv0alpha1.SecureValuePhaseSucceeded}); err != nil { + return fmt.Errorf("setting secret metadata status to Succeeded: message=%+v", message) + } + + case contracts.DeleteSecretOutboxMessage: + if err := keeper.Delete(ctx, keeperCfg, message.Namespace, contracts.ExternalID(*message.ExternalID)); err != nil { + return fmt.Errorf("calling keeper to delete secret: %w", err) + } + if err := w.secureValueMetadataStorage.Delete(ctx, xkube.Namespace(message.Namespace), message.Name); err != nil { + return fmt.Errorf("deleting secure value metadata: %+w", err) + } + + default: + return fmt.Errorf("unhandled message type: %s", message.Type) + } + + // Delete the message from the queue after completing all operations because + // if the message is deleted first, the response may be lost, + // resulting in an error, but since the message was actually deleted + // the worker would never retry. + if err := w.outboxQueue.Delete(ctx, message.MessageID); err != nil { + return fmt.Errorf("deleting message from outbox queue: %w", err) + } + + return nil +} diff --git a/pkg/registry/apis/secret/worker/worker_test.go b/pkg/registry/apis/secret/worker/worker_test.go new file mode 100644 index 00000000000..0d20a668e52 --- /dev/null +++ b/pkg/registry/apis/secret/worker/worker_test.go @@ -0,0 +1,246 @@ +package worker_test + +import ( + "context" + "fmt" + "testing" + "time" + + secretv0alpha1 "github.com/grafana/grafana/pkg/apis/secret/v0alpha1" + "github.com/grafana/grafana/pkg/registry/apis/secret/contracts" + "github.com/grafana/grafana/pkg/registry/apis/secret/testutils" + "github.com/grafana/grafana/pkg/registry/apis/secret/worker" + "github.com/grafana/grafana/pkg/registry/apis/secret/xkube" + "github.com/stretchr/testify/require" +) + +type fakeKeeperService struct { + keeperForConfigFunc func(cfg secretv0alpha1.KeeperConfig) (contracts.Keeper, error) +} + +func newFakeKeeperService(keeperForConfigFunc func(cfg secretv0alpha1.KeeperConfig) (contracts.Keeper, error)) *fakeKeeperService { + return &fakeKeeperService{keeperForConfigFunc: keeperForConfigFunc} +} + +func (s *fakeKeeperService) KeeperForConfig(cfg secretv0alpha1.KeeperConfig) (contracts.Keeper, error) { + return s.keeperForConfigFunc(cfg) +} + +func TestProcessMessage(t *testing.T) { + t.Parallel() + + t.Run("secure value metadata status is set to Failed when processing a message fails too many times", func(t *testing.T) { + t.Parallel() + + // Given a worker that will attempt to process a message N times + workerCfg := worker.Config{ + BatchSize: 10, + ReceiveTimeout: 1 * time.Second, + PollingInterval: time.Millisecond, + MaxMessageProcessingAttempts: 2, + } + + // And an error that keeps happening + keeperService := newFakeKeeperService(func(cfg secretv0alpha1.KeeperConfig) (contracts.Keeper, error) { + return nil, fmt.Errorf("oops") + }) + + sut := testutils.Setup(t, testutils.WithWorkerConfig(workerCfg), testutils.WithKeeperService(keeperService)) + ctx := context.Background() + + // Queue a create secure value operation + sv, err := sut.CreateSv(ctx) + require.NoError(t, err) + + for range workerCfg.MaxMessageProcessingAttempts + 1 { + // The secure value status should be Pending while the worker is trying to process the message + sv, err = sut.SecureValueMetadataStorage.Read(ctx, xkube.Namespace(sv.Namespace), sv.Name, contracts.ReadOpts{}) + require.NoError(t, err) + require.Equal(t, secretv0alpha1.SecureValuePhasePending, sv.Status.Phase) + + // Worker tries to process messages + _ = sut.Worker.ReceiveAndProcessMessages(ctx) + } + + // After the worker fails to process a message too many times, + // the secure value status is changed to Failed + sv, err = sut.SecureValueMetadataStorage.Read(ctx, xkube.Namespace(sv.Namespace), sv.Name, contracts.ReadOpts{}) + require.NoError(t, err) + require.Equal(t, secretv0alpha1.SecureValuePhaseFailed, sv.Status.Phase) + + messages, err := sut.OutboxQueue.ReceiveN(ctx, 100) + require.NoError(t, err) + require.Empty(t, messages) + }) + + t.Run("create sv: secure value metadata status is set to Succeeded when message is processed successfully", func(t *testing.T) { + t.Parallel() + + sut := testutils.Setup(t) + ctx := context.Background() + + // Queue a create secure value operation + sv, err := sut.CreateSv(ctx) + require.NoError(t, err) + + // Worker receives and processes the message + require.NoError(t, sut.Worker.ReceiveAndProcessMessages(ctx)) + + // and sets the secure value status to Succeeded + sv, err = sut.SecureValueMetadataStorage.Read(ctx, xkube.Namespace(sv.Namespace), sv.Name, contracts.ReadOpts{}) + require.NoError(t, err) + require.Equal(t, secretv0alpha1.SecureValuePhaseSucceeded, sv.Status.Phase) + + messages, err := sut.OutboxQueue.ReceiveN(ctx, 100) + require.NoError(t, err) + require.Empty(t, messages) + }) + + t.Run("update sv: secure value metadata status is set to Succeeded when message is processed successfully", func(t *testing.T) { + t.Parallel() + + sut := testutils.Setup(t) + ctx := context.Background() + + // Queue a create secure value operation + sv, err := sut.CreateSv(ctx) + require.NoError(t, err) + + // Worker receives and processes the message + require.NoError(t, sut.Worker.ReceiveAndProcessMessages(ctx)) + + // and sets the secure value status to Succeeded + sv, err = sut.SecureValueMetadataStorage.Read(ctx, xkube.Namespace(sv.Namespace), sv.Name, contracts.ReadOpts{}) + require.NoError(t, err) + require.Equal(t, secretv0alpha1.SecureValuePhaseSucceeded, sv.Status.Phase) + + sv.Spec.Description = "desc2" + sv.Spec.Value = secretv0alpha1.NewExposedSecureValue("v2") + + // Queue an update operation + sv, err = sut.UpdateSv(ctx, sv) + require.NoError(t, err) + require.Equal(t, secretv0alpha1.SecureValuePhasePending, sv.Status.Phase) + + // Worker receives and processes the message + require.NoError(t, sut.Worker.ReceiveAndProcessMessages(ctx)) + updatedSv, err := sut.SecureValueMetadataStorage.Read(ctx, xkube.Namespace(sv.Namespace), sv.Name, contracts.ReadOpts{}) + require.NoError(t, err) + require.Equal(t, secretv0alpha1.SecureValuePhaseSucceeded, updatedSv.Status.Phase) + require.Equal(t, sv.Spec.Description, updatedSv.Spec.Description) + + messages, err := sut.OutboxQueue.ReceiveN(ctx, 100) + require.NoError(t, err) + require.Empty(t, messages) + }) + + t.Run("delete sv: secure value metadata is deleted", func(t *testing.T) { + t.Parallel() + + sut := testutils.Setup(t) + ctx := context.Background() + + // Queue a create secure value operation + sv, err := sut.CreateSv(ctx) + require.NoError(t, err) + + // Worker receives and processes the message + require.NoError(t, sut.Worker.ReceiveAndProcessMessages(ctx)) + + // and sets the secure value status to Succeeded + sv, err = sut.SecureValueMetadataStorage.Read(ctx, xkube.Namespace(sv.Namespace), sv.Name, contracts.ReadOpts{}) + require.NoError(t, err) + require.Equal(t, secretv0alpha1.SecureValuePhaseSucceeded, sv.Status.Phase) + + // Queue a delete operation + updatedSv, err := sut.DeleteSv(ctx, sv.Namespace, sv.Name) + require.NoError(t, err) + require.Equal(t, secretv0alpha1.SecureValuePhasePending, updatedSv.Status.Phase) + + // Worker receives and processes the message + require.NoError(t, sut.Worker.ReceiveAndProcessMessages(ctx)) + + // The secure value has been deleted + _, err = sut.SecureValueMetadataStorage.Read(ctx, xkube.Namespace(sv.Namespace), sv.Name, contracts.ReadOpts{}) + require.ErrorIs(t, err, contracts.ErrSecureValueNotFound) + + messages, err := sut.OutboxQueue.ReceiveN(ctx, 100) + require.NoError(t, err) + require.Empty(t, messages) + }) + + t.Run("when creating a secure value, the secret is encrypted before it is added to the outbox queue", func(t *testing.T) { + t.Parallel() + + sut := testutils.Setup(t) + ctx := context.Background() + + // Queue a create secure value operation + var secret string + _, err := sut.CreateSv(ctx, func(cfg *testutils.CreateSvConfig) { + secret = string(cfg.Sv.Spec.Value) + }) + require.NoError(t, err) + + messages, err := sut.OutboxQueue.ReceiveN(ctx, 100) + require.NoError(t, err) + require.Equal(t, 1, len(messages)) + + encryptedSecret := messages[0].EncryptedSecret + require.NotEmpty(t, secret) + require.NotEmpty(t, encryptedSecret) + require.NotEqual(t, secret, encryptedSecret) + }) + + t.Run("when updating a secure value, the secret is encrypted before it is added to the outbox queue", func(t *testing.T) { + t.Parallel() + + sut := testutils.Setup(t) + ctx := context.Background() + + // Queue a create secure value operation + sv, err := sut.CreateSv(ctx) + require.NoError(t, err) + sv.Spec.Value = secretv0alpha1.NewExposedSecureValue("v2") + + require.NoError(t, sut.Worker.ReceiveAndProcessMessages(ctx)) + + newValue := "v2" + sv.Spec.Value = secretv0alpha1.NewExposedSecureValue(newValue) + + // Queue an update secure value operation + _, err = sut.UpdateSv(ctx, sv) + require.NoError(t, err) + + messages, err := sut.OutboxQueue.ReceiveN(ctx, 100) + require.NoError(t, err) + require.Equal(t, 1, len(messages)) + + encryptedSecret := messages[0].EncryptedSecret + require.NotEmpty(t, encryptedSecret) + require.NotEqual(t, newValue, encryptedSecret) + }) + + t.Run("when deleting a secure value, no value is added to the outbox message", func(t *testing.T) { + t.Parallel() + + sut := testutils.Setup(t) + ctx := context.Background() + + // Queue a create secure value operation + sv, err := sut.CreateSv(ctx) + require.NoError(t, err) + sv.Spec.Value = secretv0alpha1.NewExposedSecureValue("v2") + + require.NoError(t, sut.Worker.ReceiveAndProcessMessages(ctx)) + + // Queue a delete secure value operation + _, err = sut.DeleteSv(ctx, sv.Namespace, sv.Name) + require.NoError(t, err) + + messages, err := sut.OutboxQueue.ReceiveN(ctx, 100) + require.NoError(t, err) + require.Equal(t, 1, len(messages)) + require.Empty(t, messages[0].EncryptedSecret) + }) +} diff --git a/pkg/registry/backgroundsvcs/background_services.go b/pkg/registry/backgroundsvcs/background_services.go index f6ae88f9dd1..d8a9010bfd2 100644 --- a/pkg/registry/backgroundsvcs/background_services.go +++ b/pkg/registry/backgroundsvcs/background_services.go @@ -9,6 +9,7 @@ import ( "github.com/grafana/grafana/pkg/infra/usagestats/statscollector" "github.com/grafana/grafana/pkg/registry" apiregistry "github.com/grafana/grafana/pkg/registry/apis" + secretworker "github.com/grafana/grafana/pkg/registry/apis/secret/worker" appregistry "github.com/grafana/grafana/pkg/registry/apps" "github.com/grafana/grafana/pkg/services/accesscontrol/dualwrite" "github.com/grafana/grafana/pkg/services/anonymous/anonimpl" @@ -70,6 +71,7 @@ func ProvideBackgroundServiceRegistry( appRegistry *appregistry.Service, pluginDashboardUpdater *plugindashboardsservice.DashboardUpdater, dashboardServiceImpl *service.DashboardServiceImpl, + secretManagerWorker *secretworker.Worker, // Need to make sure these are initialized, is there a better place to put them? _ dashboardsnapshots.Service, _ serviceaccounts.Service, @@ -117,6 +119,7 @@ func ProvideBackgroundServiceRegistry( appRegistry, pluginDashboardUpdater, dashboardServiceImpl, + secretManagerWorker, ) } diff --git a/pkg/server/wire.go b/pkg/server/wire.go index 42cf7f6f499..efb82d04bc0 100644 --- a/pkg/server/wire.go +++ b/pkg/server/wire.go @@ -45,6 +45,8 @@ import ( secretdecrypt "github.com/grafana/grafana/pkg/registry/apis/secret/decrypt" gsmEncryption "github.com/grafana/grafana/pkg/registry/apis/secret/encryption" encryptionManager "github.com/grafana/grafana/pkg/registry/apis/secret/encryption/manager" + secretsecurevalueservice "github.com/grafana/grafana/pkg/registry/apis/secret/service" + secretworker "github.com/grafana/grafana/pkg/registry/apis/secret/worker" appregistry "github.com/grafana/grafana/pkg/registry/apps" "github.com/grafana/grafana/pkg/services/accesscontrol" "github.com/grafana/grafana/pkg/services/accesscontrol/acimpl" @@ -425,16 +427,19 @@ var wireBasicSet = wire.NewSet( secretmetadata.ProvideSecureValueMetadataStorage, secretmetadata.ProvideKeeperMetadataStorage, secretmetadata.ProvideDecryptStorage, + secretdecrypt.ProvideDecryptAuthorizer, + secretdecrypt.ProvideDecryptAllowList, secretencryption.ProvideDataKeyStorage, secretencryption.ProvideEncryptedValueStorage, secretmetadata.ProvideOutboxQueue, + secretsecurevalueservice.ProvideSecureValueService, secretmigrator.NewWithEngine, secretdatabase.ProvideDatabase, wire.Bind(new(secretcontracts.Database), new(*secretdatabase.Database)), encryptionManager.ProvideEncryptionManager, gsmEncryption.ProvideThirdPartyProviderMap, - secretdecrypt.ProvideDecryptAuthorizer, - secretdecrypt.ProvideDecryptAllowList, + secretworker.ProvideWorkerConfig, + secretworker.NewWorker, // Unified storage resource.ProvideStorageMetrics, resource.ProvideIndexMetrics, diff --git a/pkg/server/wire_gen.go b/pkg/server/wire_gen.go index d52743dc2da..166c20a19dc 100644 --- a/pkg/server/wire_gen.go +++ b/pkg/server/wire_gen.go @@ -61,8 +61,11 @@ import ( "github.com/grafana/grafana/pkg/registry/apis/secret" "github.com/grafana/grafana/pkg/registry/apis/secret/contracts" "github.com/grafana/grafana/pkg/registry/apis/secret/decrypt" - encryption3 "github.com/grafana/grafana/pkg/registry/apis/secret/encryption" + encryption2 "github.com/grafana/grafana/pkg/registry/apis/secret/encryption" manager4 "github.com/grafana/grafana/pkg/registry/apis/secret/encryption/manager" + "github.com/grafana/grafana/pkg/registry/apis/secret/secretkeeper" + service11 "github.com/grafana/grafana/pkg/registry/apis/secret/service" + "github.com/grafana/grafana/pkg/registry/apis/secret/worker" "github.com/grafana/grafana/pkg/registry/apis/userstorage" "github.com/grafana/grafana/pkg/registry/apps" advisor2 "github.com/grafana/grafana/pkg/registry/apps/advisor" @@ -111,7 +114,7 @@ import ( "github.com/grafana/grafana/pkg/services/datasources" "github.com/grafana/grafana/pkg/services/datasources/guardian" service7 "github.com/grafana/grafana/pkg/services/datasources/service" - "github.com/grafana/grafana/pkg/services/encryption" + encryption3 "github.com/grafana/grafana/pkg/services/encryption" "github.com/grafana/grafana/pkg/services/encryption/provider" service2 "github.com/grafana/grafana/pkg/services/encryption/service" "github.com/grafana/grafana/pkg/services/extsvcauth" @@ -233,7 +236,7 @@ import ( "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/legacysql/dualwrite" database5 "github.com/grafana/grafana/pkg/storage/secret/database" - encryption2 "github.com/grafana/grafana/pkg/storage/secret/encryption" + "github.com/grafana/grafana/pkg/storage/secret/encryption" "github.com/grafana/grafana/pkg/storage/secret/metadata" migrator2 "github.com/grafana/grafana/pkg/storage/secret/migrator" "github.com/grafana/grafana/pkg/storage/unified" @@ -702,6 +705,38 @@ func Initialize(cfg *setting.Cfg, opts Options, apiOpts api.ServerOptions) (*Ser } importDashboardService := service9.ProvideService(routeRegisterImpl, quotaService, service12, pluginstoreService, libraryPanelService, dashboardService, accessControl, folderimplService, featureToggles) dashboardUpdater := service6.ProvideDashboardUpdater(inProcBus, pluginstoreService, service12, importDashboardService, service11, pluginService, dashboardService) + config := worker.ProvideWorkerConfig() + databaseDatabase := database5.ProvideDatabase(sqlStore, tracer) + outboxQueue := metadata.ProvideOutboxQueue(databaseDatabase, tracer, registerer) + secureValueMetadataStorage, err := metadata.ProvideSecureValueMetadataStorage(databaseDatabase, tracer, featureToggles, registerer) + if err != nil { + return nil, err + } + keeperMetadataStorage, err := metadata.ProvideKeeperMetadataStorage(databaseDatabase, tracer, featureToggles, registerer) + if err != nil { + return nil, err + } + encryptedValueStorage, err := encryption.ProvideEncryptedValueStorage(databaseDatabase, tracer, featureToggles) + if err != nil { + return nil, err + } + dataKeyStorage, err := encryption.ProvideDataKeyStorage(databaseDatabase, tracer, featureToggles, registerer) + if err != nil { + return nil, err + } + providerMap := encryption2.ProvideThirdPartyProviderMap() + encryptionManager, err := manager4.ProvideEncryptionManager(tracer, dataKeyStorage, cfg, usageStats, providerMap) + if err != nil { + return nil, err + } + ossKeeperService, err := secretkeeper.ProvideService(tracer, encryptedValueStorage, encryptionManager, registerer) + if err != nil { + return nil, err + } + workerWorker, err := worker.NewWorker(config, tracer, databaseDatabase, outboxQueue, secureValueMetadataStorage, keeperMetadataStorage, ossKeeperService, encryptionManager, featureToggles, registerer) + if err != nil { + return nil, err + } sanitizerProvider := sanitizer.ProvideService(renderingService) healthService, err := grpcserver.ProvideHealthService(cfg, grpcserverProvider) if err != nil { @@ -768,7 +803,7 @@ func Initialize(cfg *setting.Cfg, opts Options, apiOpts api.ServerOptions) (*Ser } ossUserProtectionImpl := authinfoimpl.ProvideOSSUserProtectionService() registration := authnimpl.ProvideRegistration(cfg, authnService, orgService, userAuthTokenService, acimplService, permissionRegistry, apikeyService, userService, authService, ossUserProtectionImpl, loginattemptimplService, quotaService, authinfoimplService, renderingService, featureToggles, oauthtokenService, socialService, remoteCache, ldapImpl, ossImpl, tracingService, tempuserService, notificationService) - backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, serviceImpl, serviceAccountsProxy, sanitizerProvider, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration) + backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, workerWorker, serviceImpl, serviceAccountsProxy, sanitizerProvider, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration) usageStatsProvidersRegistry := usagestatssvcs.ProvideUsageStatsProvidersRegistry(acimplService, userService) server, err := New(opts, cfg, httpServer, acimplService, provisioningServiceImpl, backgroundServiceRegistry, usageStatsProvidersRegistry, statscollectorService, registerer) if err != nil { @@ -1215,6 +1250,38 @@ func InitializeForTest(t sqlutil.ITestDB, testingT interface { } importDashboardService := service9.ProvideService(routeRegisterImpl, quotaService, service12, pluginstoreService, libraryPanelService, dashboardService, accessControl, folderimplService, featureToggles) dashboardUpdater := service6.ProvideDashboardUpdater(inProcBus, pluginstoreService, service12, importDashboardService, service11, pluginService, dashboardService) + config := worker.ProvideWorkerConfig() + databaseDatabase := database5.ProvideDatabase(sqlStore, tracer) + outboxQueue := metadata.ProvideOutboxQueue(databaseDatabase, tracer, registerer) + secureValueMetadataStorage, err := metadata.ProvideSecureValueMetadataStorage(databaseDatabase, tracer, featureToggles, registerer) + if err != nil { + return nil, err + } + keeperMetadataStorage, err := metadata.ProvideKeeperMetadataStorage(databaseDatabase, tracer, featureToggles, registerer) + if err != nil { + return nil, err + } + encryptedValueStorage, err := encryption.ProvideEncryptedValueStorage(databaseDatabase, tracer, featureToggles) + if err != nil { + return nil, err + } + dataKeyStorage, err := encryption.ProvideDataKeyStorage(databaseDatabase, tracer, featureToggles, registerer) + if err != nil { + return nil, err + } + providerMap := encryption2.ProvideThirdPartyProviderMap() + encryptionManager, err := manager4.ProvideEncryptionManager(tracer, dataKeyStorage, cfg, usageStats, providerMap) + if err != nil { + return nil, err + } + ossKeeperService, err := secretkeeper.ProvideService(tracer, encryptedValueStorage, encryptionManager, registerer) + if err != nil { + return nil, err + } + workerWorker, err := worker.NewWorker(config, tracer, databaseDatabase, outboxQueue, secureValueMetadataStorage, keeperMetadataStorage, ossKeeperService, encryptionManager, featureToggles, registerer) + if err != nil { + return nil, err + } sanitizerProvider := sanitizer.ProvideService(renderingService) healthService, err := grpcserver.ProvideHealthService(cfg, grpcserverProvider) if err != nil { @@ -1281,7 +1348,7 @@ func InitializeForTest(t sqlutil.ITestDB, testingT interface { } ossUserProtectionImpl := authinfoimpl.ProvideOSSUserProtectionService() registration := authnimpl.ProvideRegistration(cfg, authnService, orgService, userAuthTokenService, acimplService, permissionRegistry, apikeyService, userService, authService, ossUserProtectionImpl, loginattemptimplService, quotaService, authinfoimplService, renderingService, featureToggles, oauthtokentestService, socialService, remoteCache, ldapImpl, ossImpl, tracingService, tempuserService, notificationServiceMock) - backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, serviceImpl, serviceAccountsProxy, sanitizerProvider, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration) + backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, workerWorker, serviceImpl, serviceAccountsProxy, sanitizerProvider, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration) usageStatsProvidersRegistry := usagestatssvcs.ProvideUsageStatsProvidersRegistry(acimplService, userService) server, err := New(opts, cfg, httpServer, acimplService, provisioningServiceImpl, backgroundServiceRegistry, usageStatsProvidersRegistry, statscollectorService, registerer) if err != nil { @@ -1428,7 +1495,7 @@ var withOTelSet = wire.NewSet( otelTracer, grpcserver.ProvideService, interceptors.ProvideAuthenticator, ) -var wireBasicSet = wire.NewSet(annotationsimpl.ProvideService, wire.Bind(new(annotations.Repository), new(*annotationsimpl.RepositoryImpl)), New, api.ProvideHTTPServer, query.ProvideService, wire.Bind(new(query.Service), new(*query.ServiceImpl)), bus.ProvideBus, wire.Bind(new(bus.Bus), new(*bus.InProcBus)), rendering.ProvideService, wire.Bind(new(rendering.Service), new(*rendering.RenderingService)), routing.ProvideRegister, wire.Bind(new(routing.RouteRegister), new(*routing.RouteRegisterImpl)), hooks.ProvideService, kvstore.ProvideService, localcache.ProvideService, bundleregistry.ProvideService, wire.Bind(new(supportbundles.Service), new(*bundleregistry.Service)), updatemanager.ProvideGrafanaService, updatemanager.ProvidePluginsService, service.ProvideService, wire.Bind(new(usagestats.Service), new(*service.UsageStats)), validator2.ProvideService, legacy.ProvideLegacyMigrator, pluginsintegration.WireSet, dashboards.ProvideFileStoreManager, wire.Bind(new(dashboards.FileStore), new(*dashboards.FileStoreManager)), cloudwatch.ProvideService, cloudmonitoring.ProvideService, azuremonitor.ProvideService, postgres.ProvideService, mysql.ProvideService, mssql.ProvideService, store.ProvideEntityEventsService, dualwrite.ProvideService, httpclientprovider.New, wire.Bind(new(httpclient.Provider), new(*httpclient2.Provider)), serverlock.ProvideService, annotationsimpl.ProvideCleanupService, wire.Bind(new(annotations.Cleaner), new(*annotationsimpl.CleanupServiceImpl)), cleanup.ProvideService, shorturlimpl.ProvideService, wire.Bind(new(shorturls.Service), new(*shorturlimpl.ShortURLService)), queryhistory.ProvideService, wire.Bind(new(queryhistory.Service), new(*queryhistory.QueryHistoryService)), correlations.ProvideService, wire.Bind(new(correlations.Service), new(*correlations.CorrelationsService)), quotaimpl.ProvideService, remotecache.ProvideService, wire.Bind(new(remotecache.CacheStorage), new(*remotecache.RemoteCache)), authinfoimpl.ProvideService, wire.Bind(new(login.AuthInfoService), new(*authinfoimpl.Service)), authinfoimpl.ProvideStore, datasourceproxy.ProvideService, sort.ProvideService, search2.ProvideService, searchV2.ProvideService, searchV2.ProvideSearchHTTPService, store.ProvideService, store.ProvideSystemUsersService, live.ProvideService, pushhttp.ProvideService, contexthandler.ProvideService, service10.ProvideService, wire.Bind(new(service10.LDAP), new(*service10.LDAPImpl)), jwt.ProvideService, wire.Bind(new(jwt.JWTService), new(*jwt.AuthService)), store2.ProvideDBStore, image.ProvideDeleteExpiredService, ngalert.ProvideService, librarypanels.ProvideService, wire.Bind(new(librarypanels.Service), new(*librarypanels.LibraryPanelService)), libraryelements.ProvideService, wire.Bind(new(libraryelements.Service), new(*libraryelements.LibraryElementService)), notifications.ProvideService, notifications.ProvideSmtpService, github.ProvideFactory, tracing.ProvideService, tracing.ProvideTracingConfig, wire.Bind(new(tracing.Tracer), new(*tracing.TracingService)), withOTelSet, testdatasource.ProvideService, api4.ProvideService, opentsdb.ProvideService, socialimpl.ProvideService, influxdb.ProvideService, wire.Bind(new(social.Service), new(*socialimpl.SocialService)), tempo.ProvideService, loki.ProvideService, graphite.ProvideService, prometheus.ProvideService, elasticsearch.ProvideService, pyroscope.ProvideService, parca.ProvideService, zipkin.ProvideService, jaeger.ProvideService, service7.ProvideCacheService, wire.Bind(new(datasources.CacheService), new(*service7.CacheServiceImpl)), service2.ProvideEncryptionService, wire.Bind(new(encryption.Internal), new(*service2.Service)), manager.ProvideSecretsService, wire.Bind(new(secrets.Service), new(*manager.SecretsService)), database.ProvideSecretsStore, wire.Bind(new(secrets.Store), new(*database.SecretsStoreImpl)), grafanads.ProvideService, wire.Bind(new(dashboardsnapshots.Store), new(*database4.DashboardSnapshotStore)), database4.ProvideStore, wire.Bind(new(dashboardsnapshots.Service), new(*service8.ServiceImpl)), service8.ProvideService, service7.ProvideService, wire.Bind(new(datasources.DataSourceService), new(*service7.Service)), service7.ProvideLegacyDataSourceLookup, retriever.ProvideService, wire.Bind(new(serviceaccounts.ServiceAccountRetriever), new(*retriever.Service)), ossaccesscontrol.ProvideServiceAccountPermissions, wire.Bind(new(accesscontrol.ServiceAccountPermissionsService), new(*ossaccesscontrol.ServiceAccountPermissionsService)), manager2.ProvideServiceAccountsService, proxy.ProvideServiceAccountsProxy, wire.Bind(new(serviceaccounts.Service), new(*proxy.ServiceAccountsProxy)), expr.ProvideService, featuremgmt.ProvideManagerService, featuremgmt.ProvideToggles, featuremgmt.ProvideOpenFeatureService, featuremgmt.ProvideStaticEvaluator, service5.ProvideDashboardServiceImpl, wire.Bind(new(dashboards2.PermissionsRegistrationService), new(*service5.DashboardServiceImpl)), service5.ProvideDashboardService, service5.ProvideDashboardProvisioningService, service5.ProvideDashboardPluginService, database2.ProvideDashboardStore, folderimpl.ProvideService, wire.Bind(new(folder.Service), new(*folderimpl.Service)), folderimpl.ProvideStore, wire.Bind(new(folder.Store), new(*folderimpl.FolderStoreImpl)), folderimpl.ProvideDashboardFolderStore, wire.Bind(new(folder.FolderStore), new(*folderimpl.DashboardFolderStoreImpl)), service9.ProvideService, wire.Bind(new(dashboardimport.Service), new(*service9.ImportDashboardService)), service6.ProvideService, wire.Bind(new(plugindashboards.Service), new(*service6.Service)), service6.ProvideDashboardUpdater, sanitizer.ProvideService, kvstore2.ProvideService, avatar.ProvideAvatarCacheServer, statscollector.ProvideService, csrf.ProvideCSRFFilter, wire.Bind(new(csrf.Service), new(*csrf.CSRF)), ossaccesscontrol.ProvideTeamPermissions, wire.Bind(new(accesscontrol.TeamPermissionsService), new(*ossaccesscontrol.TeamPermissionsService)), ossaccesscontrol.ProvideFolderPermissions, wire.Bind(new(accesscontrol.FolderPermissionsService), new(*ossaccesscontrol.FolderPermissionsService)), ossaccesscontrol.ProvideDashboardPermissions, wire.Bind(new(accesscontrol.DashboardPermissionsService), new(*ossaccesscontrol.DashboardPermissionsService)), ossaccesscontrol.ProvideReceiverPermissionsService, wire.Bind(new(accesscontrol.ReceiverPermissionsService), new(*ossaccesscontrol.ReceiverPermissionsService)), starimpl.ProvideService, playlistimpl.ProvideService, apikeyimpl.ProvideService, dashverimpl.ProvideService, service3.ProvideService, wire.Bind(new(publicdashboards.Service), new(*service3.PublicDashboardServiceImpl)), database3.ProvideStore, wire.Bind(new(publicdashboards.Store), new(*database3.PublicDashboardStoreImpl)), metric.ProvideService, api2.ProvideApi, api3.ProvideApi, userimpl.ProvideService, orgimpl.ProvideService, orgimpl.ProvideDeletionService, statsimpl.ProvideService, grpccontext.ProvideContextHandler, grpcserver.ProvideHealthService, grpcserver.ProvideReflectionService, resolver.ProvideEntityReferenceResolver, teamimpl.ProvideService, teamapi.ProvideTeamAPI, tempuserimpl.ProvideService, loginattemptimpl.ProvideService, wire.Bind(new(loginattempt.Service), new(*loginattemptimpl.Service)), migrations2.ProvideDataSourceMigrationService, migrations2.ProvideSecretMigrationProvider, wire.Bind(new(migrations2.SecretMigrationProvider), new(*migrations2.SecretMigrationProviderImpl)), resourcepermissions.NewActionSetService, wire.Bind(new(accesscontrol.ActionResolver), new(resourcepermissions.ActionSetService)), wire.Bind(new(pluginaccesscontrol.ActionSetRegistry), new(resourcepermissions.ActionSetService)), permreg.ProvidePermissionRegistry, acimpl.ProvideAccessControl, dualwrite2.ProvideZanzanaReconciler, navtreeimpl.ProvideService, wire.Bind(new(accesscontrol.AccessControl), new(*acimpl.AccessControl)), wire.Bind(new(notifications.TempUserStore), new(tempuser.Service)), tagimpl.ProvideService, wire.Bind(new(tag.Service), new(*tagimpl.Service)), authnimpl.ProvideService, authnimpl.ProvideIdentitySynchronizer, authnimpl.ProvideAuthnService, authnimpl.ProvideAuthnServiceAuthenticateOnly, authnimpl.ProvideRegistration, supportbundlesimpl.ProvideService, extsvcaccounts.ProvideExtSvcAccountsService, wire.Bind(new(serviceaccounts.ExtSvcAccountsService), new(*extsvcaccounts.ExtSvcAccountsService)), registry2.ProvideExtSvcRegistry, wire.Bind(new(extsvcauth.ExternalServiceRegistry), new(*registry2.Registry)), anonstore.ProvideAnonDBStore, wire.Bind(new(anonstore.AnonStore), new(*anonstore.AnonDBStore)), loggermw.Provide, slogadapter.Provide, signingkeysimpl.ProvideEmbeddedSigningKeysService, wire.Bind(new(signingkeys.Service), new(*signingkeysimpl.Service)), ssosettingsimpl.ProvideService, wire.Bind(new(ssosettings.Service), new(*ssosettingsimpl.Service)), idimpl.ProvideService, wire.Bind(new(auth.IDService), new(*idimpl.Service)), cloudmigrationimpl.ProvideService, userimpl.ProvideVerifier, connectors.ProvideOrgRoleMapper, wire.Bind(new(user.Verifier), new(*userimpl.Verifier)), authz.WireSet, metadata.ProvideSecureValueMetadataStorage, metadata.ProvideKeeperMetadataStorage, metadata.ProvideDecryptStorage, encryption2.ProvideDataKeyStorage, encryption2.ProvideEncryptedValueStorage, metadata.ProvideOutboxQueue, migrator2.NewWithEngine, database5.ProvideDatabase, wire.Bind(new(contracts.Database), new(*database5.Database)), manager4.ProvideEncryptionManager, encryption3.ProvideThirdPartyProviderMap, decrypt.ProvideDecryptAuthorizer, decrypt.ProvideDecryptAllowList, resource.ProvideStorageMetrics, resource.ProvideIndexMetrics, apiserver.WireSet, apiregistry.WireSet, appregistry.WireSet) +var wireBasicSet = wire.NewSet(annotationsimpl.ProvideService, wire.Bind(new(annotations.Repository), new(*annotationsimpl.RepositoryImpl)), New, api.ProvideHTTPServer, query.ProvideService, wire.Bind(new(query.Service), new(*query.ServiceImpl)), bus.ProvideBus, wire.Bind(new(bus.Bus), new(*bus.InProcBus)), rendering.ProvideService, wire.Bind(new(rendering.Service), new(*rendering.RenderingService)), routing.ProvideRegister, wire.Bind(new(routing.RouteRegister), new(*routing.RouteRegisterImpl)), hooks.ProvideService, kvstore.ProvideService, localcache.ProvideService, bundleregistry.ProvideService, wire.Bind(new(supportbundles.Service), new(*bundleregistry.Service)), updatemanager.ProvideGrafanaService, updatemanager.ProvidePluginsService, service.ProvideService, wire.Bind(new(usagestats.Service), new(*service.UsageStats)), validator2.ProvideService, legacy.ProvideLegacyMigrator, pluginsintegration.WireSet, dashboards.ProvideFileStoreManager, wire.Bind(new(dashboards.FileStore), new(*dashboards.FileStoreManager)), cloudwatch.ProvideService, cloudmonitoring.ProvideService, azuremonitor.ProvideService, postgres.ProvideService, mysql.ProvideService, mssql.ProvideService, store.ProvideEntityEventsService, dualwrite.ProvideService, httpclientprovider.New, wire.Bind(new(httpclient.Provider), new(*httpclient2.Provider)), serverlock.ProvideService, annotationsimpl.ProvideCleanupService, wire.Bind(new(annotations.Cleaner), new(*annotationsimpl.CleanupServiceImpl)), cleanup.ProvideService, shorturlimpl.ProvideService, wire.Bind(new(shorturls.Service), new(*shorturlimpl.ShortURLService)), queryhistory.ProvideService, wire.Bind(new(queryhistory.Service), new(*queryhistory.QueryHistoryService)), correlations.ProvideService, wire.Bind(new(correlations.Service), new(*correlations.CorrelationsService)), quotaimpl.ProvideService, remotecache.ProvideService, wire.Bind(new(remotecache.CacheStorage), new(*remotecache.RemoteCache)), authinfoimpl.ProvideService, wire.Bind(new(login.AuthInfoService), new(*authinfoimpl.Service)), authinfoimpl.ProvideStore, datasourceproxy.ProvideService, sort.ProvideService, search2.ProvideService, searchV2.ProvideService, searchV2.ProvideSearchHTTPService, store.ProvideService, store.ProvideSystemUsersService, live.ProvideService, pushhttp.ProvideService, contexthandler.ProvideService, service10.ProvideService, wire.Bind(new(service10.LDAP), new(*service10.LDAPImpl)), jwt.ProvideService, wire.Bind(new(jwt.JWTService), new(*jwt.AuthService)), store2.ProvideDBStore, image.ProvideDeleteExpiredService, ngalert.ProvideService, librarypanels.ProvideService, wire.Bind(new(librarypanels.Service), new(*librarypanels.LibraryPanelService)), libraryelements.ProvideService, wire.Bind(new(libraryelements.Service), new(*libraryelements.LibraryElementService)), notifications.ProvideService, notifications.ProvideSmtpService, github.ProvideFactory, tracing.ProvideService, tracing.ProvideTracingConfig, wire.Bind(new(tracing.Tracer), new(*tracing.TracingService)), withOTelSet, testdatasource.ProvideService, api4.ProvideService, opentsdb.ProvideService, socialimpl.ProvideService, influxdb.ProvideService, wire.Bind(new(social.Service), new(*socialimpl.SocialService)), tempo.ProvideService, loki.ProvideService, graphite.ProvideService, prometheus.ProvideService, elasticsearch.ProvideService, pyroscope.ProvideService, parca.ProvideService, zipkin.ProvideService, jaeger.ProvideService, service7.ProvideCacheService, wire.Bind(new(datasources.CacheService), new(*service7.CacheServiceImpl)), service2.ProvideEncryptionService, wire.Bind(new(encryption3.Internal), new(*service2.Service)), manager.ProvideSecretsService, wire.Bind(new(secrets.Service), new(*manager.SecretsService)), database.ProvideSecretsStore, wire.Bind(new(secrets.Store), new(*database.SecretsStoreImpl)), grafanads.ProvideService, wire.Bind(new(dashboardsnapshots.Store), new(*database4.DashboardSnapshotStore)), database4.ProvideStore, wire.Bind(new(dashboardsnapshots.Service), new(*service8.ServiceImpl)), service8.ProvideService, service7.ProvideService, wire.Bind(new(datasources.DataSourceService), new(*service7.Service)), service7.ProvideLegacyDataSourceLookup, retriever.ProvideService, wire.Bind(new(serviceaccounts.ServiceAccountRetriever), new(*retriever.Service)), ossaccesscontrol.ProvideServiceAccountPermissions, wire.Bind(new(accesscontrol.ServiceAccountPermissionsService), new(*ossaccesscontrol.ServiceAccountPermissionsService)), manager2.ProvideServiceAccountsService, proxy.ProvideServiceAccountsProxy, wire.Bind(new(serviceaccounts.Service), new(*proxy.ServiceAccountsProxy)), expr.ProvideService, featuremgmt.ProvideManagerService, featuremgmt.ProvideToggles, featuremgmt.ProvideOpenFeatureService, featuremgmt.ProvideStaticEvaluator, service5.ProvideDashboardServiceImpl, wire.Bind(new(dashboards2.PermissionsRegistrationService), new(*service5.DashboardServiceImpl)), service5.ProvideDashboardService, service5.ProvideDashboardProvisioningService, service5.ProvideDashboardPluginService, database2.ProvideDashboardStore, folderimpl.ProvideService, wire.Bind(new(folder.Service), new(*folderimpl.Service)), folderimpl.ProvideStore, wire.Bind(new(folder.Store), new(*folderimpl.FolderStoreImpl)), folderimpl.ProvideDashboardFolderStore, wire.Bind(new(folder.FolderStore), new(*folderimpl.DashboardFolderStoreImpl)), service9.ProvideService, wire.Bind(new(dashboardimport.Service), new(*service9.ImportDashboardService)), service6.ProvideService, wire.Bind(new(plugindashboards.Service), new(*service6.Service)), service6.ProvideDashboardUpdater, sanitizer.ProvideService, kvstore2.ProvideService, avatar.ProvideAvatarCacheServer, statscollector.ProvideService, csrf.ProvideCSRFFilter, wire.Bind(new(csrf.Service), new(*csrf.CSRF)), ossaccesscontrol.ProvideTeamPermissions, wire.Bind(new(accesscontrol.TeamPermissionsService), new(*ossaccesscontrol.TeamPermissionsService)), ossaccesscontrol.ProvideFolderPermissions, wire.Bind(new(accesscontrol.FolderPermissionsService), new(*ossaccesscontrol.FolderPermissionsService)), ossaccesscontrol.ProvideDashboardPermissions, wire.Bind(new(accesscontrol.DashboardPermissionsService), new(*ossaccesscontrol.DashboardPermissionsService)), ossaccesscontrol.ProvideReceiverPermissionsService, wire.Bind(new(accesscontrol.ReceiverPermissionsService), new(*ossaccesscontrol.ReceiverPermissionsService)), starimpl.ProvideService, playlistimpl.ProvideService, apikeyimpl.ProvideService, dashverimpl.ProvideService, service3.ProvideService, wire.Bind(new(publicdashboards.Service), new(*service3.PublicDashboardServiceImpl)), database3.ProvideStore, wire.Bind(new(publicdashboards.Store), new(*database3.PublicDashboardStoreImpl)), metric.ProvideService, api2.ProvideApi, api3.ProvideApi, userimpl.ProvideService, orgimpl.ProvideService, orgimpl.ProvideDeletionService, statsimpl.ProvideService, grpccontext.ProvideContextHandler, grpcserver.ProvideHealthService, grpcserver.ProvideReflectionService, resolver.ProvideEntityReferenceResolver, teamimpl.ProvideService, teamapi.ProvideTeamAPI, tempuserimpl.ProvideService, loginattemptimpl.ProvideService, wire.Bind(new(loginattempt.Service), new(*loginattemptimpl.Service)), migrations2.ProvideDataSourceMigrationService, migrations2.ProvideSecretMigrationProvider, wire.Bind(new(migrations2.SecretMigrationProvider), new(*migrations2.SecretMigrationProviderImpl)), resourcepermissions.NewActionSetService, wire.Bind(new(accesscontrol.ActionResolver), new(resourcepermissions.ActionSetService)), wire.Bind(new(pluginaccesscontrol.ActionSetRegistry), new(resourcepermissions.ActionSetService)), permreg.ProvidePermissionRegistry, acimpl.ProvideAccessControl, dualwrite2.ProvideZanzanaReconciler, navtreeimpl.ProvideService, wire.Bind(new(accesscontrol.AccessControl), new(*acimpl.AccessControl)), wire.Bind(new(notifications.TempUserStore), new(tempuser.Service)), tagimpl.ProvideService, wire.Bind(new(tag.Service), new(*tagimpl.Service)), authnimpl.ProvideService, authnimpl.ProvideIdentitySynchronizer, authnimpl.ProvideAuthnService, authnimpl.ProvideAuthnServiceAuthenticateOnly, authnimpl.ProvideRegistration, supportbundlesimpl.ProvideService, extsvcaccounts.ProvideExtSvcAccountsService, wire.Bind(new(serviceaccounts.ExtSvcAccountsService), new(*extsvcaccounts.ExtSvcAccountsService)), registry2.ProvideExtSvcRegistry, wire.Bind(new(extsvcauth.ExternalServiceRegistry), new(*registry2.Registry)), anonstore.ProvideAnonDBStore, wire.Bind(new(anonstore.AnonStore), new(*anonstore.AnonDBStore)), loggermw.Provide, slogadapter.Provide, signingkeysimpl.ProvideEmbeddedSigningKeysService, wire.Bind(new(signingkeys.Service), new(*signingkeysimpl.Service)), ssosettingsimpl.ProvideService, wire.Bind(new(ssosettings.Service), new(*ssosettingsimpl.Service)), idimpl.ProvideService, wire.Bind(new(auth.IDService), new(*idimpl.Service)), cloudmigrationimpl.ProvideService, userimpl.ProvideVerifier, connectors.ProvideOrgRoleMapper, wire.Bind(new(user.Verifier), new(*userimpl.Verifier)), authz.WireSet, metadata.ProvideSecureValueMetadataStorage, metadata.ProvideKeeperMetadataStorage, metadata.ProvideDecryptStorage, decrypt.ProvideDecryptAuthorizer, decrypt.ProvideDecryptAllowList, encryption.ProvideDataKeyStorage, encryption.ProvideEncryptedValueStorage, metadata.ProvideOutboxQueue, service11.ProvideSecureValueService, migrator2.NewWithEngine, database5.ProvideDatabase, wire.Bind(new(contracts.Database), new(*database5.Database)), manager4.ProvideEncryptionManager, encryption2.ProvideThirdPartyProviderMap, worker.ProvideWorkerConfig, worker.NewWorker, resource.ProvideStorageMetrics, resource.ProvideIndexMetrics, apiserver.WireSet, apiregistry.WireSet, appregistry.WireSet) var wireSet = wire.NewSet( wireBasicSet, metrics.WireSet, sqlstore.ProvideService, metrics2.ProvideService, wire.Bind(new(notifications.Service), new(*notifications.NotificationService)), wire.Bind(new(notifications.WebhookSender), new(*notifications.NotificationService)), wire.Bind(new(notifications.EmailSender), new(*notifications.NotificationService)), wire.Bind(new(db.DB), new(*sqlstore.SQLStore)), prefimpl.ProvideService, oauthtoken.ProvideService, wire.Bind(new(oauthtoken.OAuthTokenService), new(*oauthtoken.Service)), wire.Bind(new(cleanup.AlertRuleService), new(*store2.DBstore)), diff --git a/pkg/storage/secret/metadata/outbox_store.go b/pkg/storage/secret/metadata/outbox_store.go index 1f183b7882b..efe569384f4 100644 --- a/pkg/storage/secret/metadata/outbox_store.go +++ b/pkg/storage/secret/metadata/outbox_store.go @@ -52,7 +52,6 @@ type outboxMessageDB struct { } func (s *outboxStore) Append(ctx context.Context, input contracts.AppendOutboxMessage) (messageID int64, err error) { - start := time.Now() ctx, span := s.tracer.Start(ctx, "outboxStore.Append", trace.WithAttributes( attribute.String("name", input.Name), attribute.String("namespace", input.Namespace), @@ -74,6 +73,7 @@ func (s *outboxStore) Append(ctx context.Context, input contracts.AppendOutboxMe assert.True(input.Type != "", "outboxStore.Append: outbox message type is required") + start := time.Now() messageID, err = s.insertMessage(ctx, input) if err != nil { return messageID, fmt.Errorf("inserting message into outbox table: %+w", err) @@ -156,7 +156,6 @@ func (s *outboxStore) insertMessage(ctx context.Context, input contracts.AppendO } func (s *outboxStore) ReceiveN(ctx context.Context, limit uint) ([]contracts.OutboxMessage, error) { - start := time.Now() messageIDs, err := s.fetchMessageIdsInQueue(ctx, limit) if err != nil { return nil, fmt.Errorf("fetching message ids from queue: %w", err) @@ -170,6 +169,7 @@ func (s *outboxStore) ReceiveN(ctx context.Context, limit uint) ([]contracts.Out MessageIDs: messageIDs, } + start := time.Now() query, err := sqltemplate.Execute(sqlSecureValueOutboxReceiveN, req) if err != nil { return nil, fmt.Errorf("execute template %q: %w", sqlSecureValueOutboxReceiveN.Name(), err)