SecretsManager: Introduce worker and secret async service (#107614)

SecretsManager: Introduce worker and secret aysnc service

Co-authored-by: PoorlyDefinedBehaviour <brunotj2015@hotmail.com>
Co-authored-by: Matheus Macabu <macabu@users.noreply.github.com>
Co-authored-by: Michael Mandrus <michael.mandrus@grafana.com>
This commit is contained in:
Dana Axinte
2025-07-04 13:13:48 +01:00
committed by GitHub
co-authored by PoorlyDefinedBehaviour Matheus Macabu Michael Mandrus
parent 76a21fb2e2
commit 46c38fdbb7
11 changed files with 1264 additions and 10 deletions
@@ -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
}
@@ -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
}
@@ -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
}
@@ -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)
})
}
@@ -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()
}
+274
View File
@@ -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
}
@@ -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)
})
}
@@ -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,
)
}
+7 -2
View File
@@ -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,
+73 -6
View File
File diff suppressed because one or more lines are too long
+2 -2
View File
@@ -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)