Secrets: Refactor data_key_id out of the encoded secure value payload (#111852)
* everything compiles * tests pass * remove file included by accident * add entry to gitignore * some scaffolding for the migration executor * remove file * implement and test the migration * use xkube.Namespace in our interfaces * add todo * update wire deps * add some logs * fix wire dependency ordering * create tests to validate error conditions during migrations
This commit is contained in:
@@ -1,6 +1,10 @@
|
||||
package contracts
|
||||
|
||||
import "context"
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/xkube"
|
||||
)
|
||||
|
||||
// EncryptionManager is an envelope encryption service in charge of encrypting/decrypting secrets.
|
||||
type EncryptionManager interface {
|
||||
@@ -8,17 +12,23 @@ type EncryptionManager interface {
|
||||
// For those specific use cases where the encryption operation cannot be moved outside
|
||||
// the database transaction, look at database-specific methods present at the specific
|
||||
// implementation present at manager.EncryptionService.
|
||||
Encrypt(ctx context.Context, namespace string, payload []byte) ([]byte, error)
|
||||
Decrypt(ctx context.Context, namespace string, payload []byte) ([]byte, error)
|
||||
Encrypt(ctx context.Context, namespace xkube.Namespace, payload []byte) (EncryptedPayload, error)
|
||||
Decrypt(ctx context.Context, namespace xkube.Namespace, payload EncryptedPayload) ([]byte, error)
|
||||
}
|
||||
|
||||
type EncryptedPayload struct {
|
||||
DataKeyID string
|
||||
EncryptedData []byte
|
||||
}
|
||||
|
||||
type EncryptedValue struct {
|
||||
Namespace string
|
||||
Name string
|
||||
Version int64
|
||||
EncryptedData []byte
|
||||
Created int64
|
||||
Updated int64
|
||||
EncryptedPayload
|
||||
|
||||
Namespace string
|
||||
Name string
|
||||
Version int64
|
||||
Created int64
|
||||
Updated int64
|
||||
}
|
||||
|
||||
// ListOpts defines pagination options for listing encrypted values.
|
||||
@@ -28,10 +38,10 @@ type ListOpts struct {
|
||||
}
|
||||
|
||||
type EncryptedValueStorage interface {
|
||||
Create(ctx context.Context, namespace, name string, version int64, encryptedData []byte) (*EncryptedValue, error)
|
||||
Update(ctx context.Context, namespace, name string, version int64, encryptedData []byte) error
|
||||
Get(ctx context.Context, namespace, name string, version int64) (*EncryptedValue, error)
|
||||
Delete(ctx context.Context, namespace, name string, version int64) error
|
||||
Create(ctx context.Context, namespace xkube.Namespace, name string, version int64, encryptedData EncryptedPayload) (*EncryptedValue, error)
|
||||
Update(ctx context.Context, namespace xkube.Namespace, name string, version int64, encryptedData EncryptedPayload) error
|
||||
Get(ctx context.Context, namespace xkube.Namespace, name string, version int64) (*EncryptedValue, error)
|
||||
Delete(ctx context.Context, namespace xkube.Namespace, name string, version int64) error
|
||||
}
|
||||
|
||||
type GlobalEncryptedValueStorage interface {
|
||||
@@ -39,6 +49,10 @@ type GlobalEncryptedValueStorage interface {
|
||||
CountAll(ctx context.Context, untilTime *int64) (int64, error)
|
||||
}
|
||||
|
||||
type EncryptedValueMigrationExecutor interface {
|
||||
Execute(ctx context.Context) (int, error)
|
||||
}
|
||||
|
||||
type ConsolidationService interface {
|
||||
Consolidate(ctx context.Context) error
|
||||
}
|
||||
|
||||
@@ -96,10 +96,10 @@ func (s ExternalID) String() string {
|
||||
|
||||
// Keeper is the interface for secret keepers.
|
||||
type Keeper interface {
|
||||
Store(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace, name string, version int64, exposedValueOrRef string) (ExternalID, error)
|
||||
Update(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace, name string, version int64, exposedValueOrRef string) error
|
||||
Expose(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace, name string, version int64) (secretv1beta1.ExposedSecureValue, error)
|
||||
Delete(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace, name string, version int64) error
|
||||
Store(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace xkube.Namespace, name string, version int64, exposedValueOrRef string) (ExternalID, error)
|
||||
Update(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace xkube.Namespace, name string, version int64, exposedValueOrRef string) error
|
||||
Expose(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace xkube.Namespace, name string, version int64) (secretv1beta1.ExposedSecureValue, error)
|
||||
Delete(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace xkube.Namespace, name string, version int64) error
|
||||
}
|
||||
|
||||
// Service is the interface for secret keeper services.
|
||||
|
||||
@@ -1,10 +1,8 @@
|
||||
package manager
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/base64"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strconv"
|
||||
@@ -20,13 +18,10 @@ import (
|
||||
"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/cipher"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/xkube"
|
||||
"github.com/grafana/grafana/pkg/util"
|
||||
)
|
||||
|
||||
const (
|
||||
keyIdDelimiter = '#'
|
||||
)
|
||||
|
||||
type EncryptionManager struct {
|
||||
tracer trace.Tracer
|
||||
store contracts.DataKeyStorage
|
||||
@@ -99,12 +94,9 @@ func (s *EncryptionManager) registerUsageMetrics() {
|
||||
})
|
||||
}
|
||||
|
||||
// TODO: Why do we need to use a global variable for this?
|
||||
var b64 = base64.RawStdEncoding
|
||||
|
||||
func (s *EncryptionManager) Encrypt(ctx context.Context, namespace string, payload []byte) ([]byte, error) {
|
||||
func (s *EncryptionManager) Encrypt(ctx context.Context, namespace xkube.Namespace, payload []byte) (contracts.EncryptedPayload, error) {
|
||||
ctx, span := s.tracer.Start(ctx, "EnvelopeEncryptionManager.Encrypt", trace.WithAttributes(
|
||||
attribute.String("namespace", namespace),
|
||||
attribute.String("namespace", namespace.String()),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
@@ -128,34 +120,30 @@ func (s *EncryptionManager) Encrypt(ctx context.Context, namespace string, paylo
|
||||
id, dataKey, err = s.currentDataKey(ctx, namespace, label)
|
||||
if err != nil {
|
||||
s.log.Error("Failed to get current data key", "error", err, "label", label)
|
||||
return nil, err
|
||||
return contracts.EncryptedPayload{}, err
|
||||
}
|
||||
|
||||
var encrypted []byte
|
||||
encrypted, err = s.cipher.Encrypt(ctx, payload, string(dataKey))
|
||||
if err != nil {
|
||||
s.log.Error("Failed to encrypt secret", "error", err)
|
||||
return nil, err
|
||||
return contracts.EncryptedPayload{}, err
|
||||
}
|
||||
|
||||
prefix := make([]byte, b64.EncodedLen(len(id))+2)
|
||||
b64.Encode(prefix[1:], []byte(id))
|
||||
prefix[0] = keyIdDelimiter
|
||||
prefix[len(prefix)-1] = keyIdDelimiter
|
||||
encryptedPayload := contracts.EncryptedPayload{
|
||||
DataKeyID: id,
|
||||
EncryptedData: encrypted,
|
||||
}
|
||||
|
||||
blob := make([]byte, len(prefix)+len(encrypted))
|
||||
copy(blob, prefix)
|
||||
copy(blob[len(prefix):], encrypted)
|
||||
|
||||
return blob, nil
|
||||
return encryptedPayload, nil
|
||||
}
|
||||
|
||||
// currentDataKey looks up for current data key in cache or database by name, and decrypts it.
|
||||
// If there's no current data key in cache nor in database it generates a new random data key,
|
||||
// and stores it into both the in-memory cache and database (encrypted by the encryption provider).
|
||||
func (s *EncryptionManager) currentDataKey(ctx context.Context, namespace string, label string) (string, []byte, error) {
|
||||
func (s *EncryptionManager) currentDataKey(ctx context.Context, namespace xkube.Namespace, label string) (string, []byte, error) {
|
||||
ctx, span := s.tracer.Start(ctx, "EnvelopeEncryptionManager.CurrentDataKey", trace.WithAttributes(
|
||||
attribute.String("namespace", namespace),
|
||||
attribute.String("namespace", namespace.String()),
|
||||
attribute.String("label", label),
|
||||
))
|
||||
defer span.End()
|
||||
@@ -166,14 +154,14 @@ func (s *EncryptionManager) currentDataKey(ctx context.Context, namespace string
|
||||
defer s.mtx.Unlock()
|
||||
|
||||
// We try to fetch the data key, either from cache or database
|
||||
id, dataKey, err := s.dataKeyByLabel(ctx, namespace, label)
|
||||
id, dataKey, err := s.dataKeyByLabel(ctx, namespace.String(), label)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
|
||||
// If no existing data key was found, create a new one
|
||||
if dataKey == nil {
|
||||
id, dataKey, err = s.newDataKey(ctx, namespace, label)
|
||||
id, dataKey, err = s.newDataKey(ctx, namespace.String(), label)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
@@ -264,9 +252,9 @@ func newRandomDataKey() ([]byte, error) {
|
||||
return rawDataKey, nil
|
||||
}
|
||||
|
||||
func (s *EncryptionManager) Decrypt(ctx context.Context, namespace string, payload []byte) ([]byte, error) {
|
||||
func (s *EncryptionManager) Decrypt(ctx context.Context, namespace xkube.Namespace, payload contracts.EncryptedPayload) ([]byte, error) {
|
||||
ctx, span := s.tracer.Start(ctx, "EnvelopeEncryptionManager.Decrypt", trace.WithAttributes(
|
||||
attribute.String("namespace", namespace),
|
||||
attribute.String("namespace", namespace.String()),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
@@ -285,50 +273,28 @@ func (s *EncryptionManager) Decrypt(ctx context.Context, namespace string, paylo
|
||||
}
|
||||
}()
|
||||
|
||||
if len(payload) == 0 {
|
||||
if len(payload.EncryptedData) == 0 {
|
||||
err = fmt.Errorf("unable to decrypt empty payload")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
payload = payload[1:]
|
||||
endOfKey := bytes.Index(payload, []byte{keyIdDelimiter})
|
||||
if endOfKey == -1 {
|
||||
err = fmt.Errorf("could not find valid key id in encrypted payload")
|
||||
return nil, err
|
||||
}
|
||||
b64Key := payload[:endOfKey]
|
||||
payload = payload[endOfKey+1:]
|
||||
keyId := make([]byte, b64.DecodedLen(len(b64Key)))
|
||||
_, err = b64.Decode(keyId, b64Key)
|
||||
if err != nil {
|
||||
if payload.DataKeyID == "" {
|
||||
err = fmt.Errorf("unable to decrypt empty data key id")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
dataKey, err := s.dataKeyById(ctx, namespace, string(keyId))
|
||||
dataKey, err := s.dataKeyById(ctx, namespace.String(), payload.DataKeyID)
|
||||
if err != nil {
|
||||
s.log.FromContext(ctx).Error("Failed to lookup data key by id", "id", string(keyId), "error", err)
|
||||
s.log.FromContext(ctx).Error("Failed to lookup data key by id", "id", payload.DataKeyID, "error", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var decrypted []byte
|
||||
decrypted, err = s.cipher.Decrypt(ctx, payload, string(dataKey))
|
||||
decrypted, err = s.cipher.Decrypt(ctx, payload.EncryptedData, string(dataKey))
|
||||
|
||||
return decrypted, err
|
||||
}
|
||||
|
||||
func (s *EncryptionManager) GetDecryptedValue(ctx context.Context, namespace string, sjd map[string][]byte, key, fallback string) string {
|
||||
if value, ok := sjd[key]; ok {
|
||||
decryptedData, err := s.Decrypt(ctx, namespace, value)
|
||||
if err != nil {
|
||||
return fallback
|
||||
}
|
||||
|
||||
return string(decryptedData)
|
||||
}
|
||||
|
||||
return fallback
|
||||
}
|
||||
|
||||
// dataKeyById looks up for data key in the database and returns it decrypted.
|
||||
func (s *EncryptionManager) dataKeyById(ctx context.Context, namespace, id string) ([]byte, error) {
|
||||
ctx, span := s.tracer.Start(ctx, "EnvelopeEncryptionManager.GetDataKey", trace.WithAttributes(
|
||||
|
||||
@@ -17,6 +17,7 @@ import (
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/encryption"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/encryption/cipher/service"
|
||||
osskmsproviders "github.com/grafana/grafana/pkg/registry/apis/secret/encryption/kmsproviders"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/xkube"
|
||||
"github.com/grafana/grafana/pkg/services/sqlstore"
|
||||
"github.com/grafana/grafana/pkg/setting"
|
||||
"github.com/grafana/grafana/pkg/storage/secret/database"
|
||||
@@ -34,7 +35,7 @@ func TestMain(m *testing.M) {
|
||||
func TestEncryptionService_EnvelopeEncryption(t *testing.T) {
|
||||
svc := setupTestService(t)
|
||||
ctx := context.Background()
|
||||
namespace := "test-namespace"
|
||||
namespace := xkube.Namespace("test-namespace")
|
||||
|
||||
t.Run("encrypting should create DEK", func(t *testing.T) {
|
||||
plaintext := []byte("very secret string")
|
||||
@@ -46,7 +47,7 @@ func TestEncryptionService_EnvelopeEncryption(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, plaintext, decrypted)
|
||||
|
||||
keys, err := svc.store.ListDataKeys(ctx, namespace)
|
||||
keys, err := svc.store.ListDataKeys(ctx, namespace.String())
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, len(keys), 1)
|
||||
})
|
||||
@@ -61,7 +62,7 @@ func TestEncryptionService_EnvelopeEncryption(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, plaintext, decrypted)
|
||||
|
||||
keys, err := svc.store.ListDataKeys(ctx, namespace)
|
||||
keys, err := svc.store.ListDataKeys(ctx, namespace.String())
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, len(keys), 1)
|
||||
})
|
||||
@@ -212,7 +213,7 @@ func TestEncryptionService_UseCurrentProvider(t *testing.T) {
|
||||
}
|
||||
encryptionManager.providerConfig.CurrentProvider = encryption.ProviderID("fakeProvider.v1")
|
||||
|
||||
namespace := "test-namespace"
|
||||
namespace := xkube.Namespace("test-namespace")
|
||||
encrypted, _ := encryptionManager.Encrypt(context.Background(), namespace, []byte{})
|
||||
assert.True(t, fake.encryptCalled)
|
||||
assert.False(t, fake.decryptCalled)
|
||||
@@ -241,7 +242,7 @@ func TestEncryptionService_UseCurrentProvider(t *testing.T) {
|
||||
|
||||
func TestEncryptionService_SecretKeyVersionUpgrade(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
namespace := "test-namespace"
|
||||
namespace := xkube.Namespace("test-namespace")
|
||||
|
||||
// Generate random keys for testing
|
||||
oldKey := util.GenerateShortUID() + util.GenerateShortUID() // 32 chars
|
||||
@@ -416,16 +417,30 @@ func (p *fakeProvider) Decrypt(_ context.Context, _ []byte) ([]byte, error) {
|
||||
|
||||
func TestEncryptionService_Decrypt(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
namespace := "test-namespace"
|
||||
namespace := xkube.Namespace("test-namespace")
|
||||
|
||||
t.Run("empty payload should fail", func(t *testing.T) {
|
||||
svc := setupTestService(t)
|
||||
_, err := svc.Decrypt(context.Background(), namespace, []byte(""))
|
||||
_, err := svc.Decrypt(context.Background(), namespace, contracts.EncryptedPayload{
|
||||
DataKeyID: "test-data-key-id",
|
||||
EncryptedData: []byte(""),
|
||||
})
|
||||
require.Error(t, err)
|
||||
|
||||
assert.Equal(t, "unable to decrypt empty payload", err.Error())
|
||||
})
|
||||
|
||||
t.Run("empty data key id should fail", func(t *testing.T) {
|
||||
svc := setupTestService(t)
|
||||
_, err := svc.Decrypt(context.Background(), namespace, contracts.EncryptedPayload{
|
||||
DataKeyID: "",
|
||||
EncryptedData: []byte("some payload"),
|
||||
})
|
||||
require.Error(t, err)
|
||||
|
||||
assert.Equal(t, "unable to decrypt empty data key id", err.Error())
|
||||
})
|
||||
|
||||
t.Run("ee encrypted payload with ee enabled should work", func(t *testing.T) {
|
||||
svc := setupTestService(t)
|
||||
ciphertext, err := svc.Encrypt(ctx, namespace, []byte("grafana"))
|
||||
@@ -442,7 +457,7 @@ func TestIntegration_SecretsService(t *testing.T) {
|
||||
|
||||
ctx := context.Background()
|
||||
someData := []byte(`some-data`)
|
||||
namespace := "test-namespace"
|
||||
namespace := xkube.Namespace("test-namespace")
|
||||
|
||||
tcs := map[string]func(*testing.T, db.DB, contracts.EncryptionManager){
|
||||
"regular": func(t *testing.T, _ db.DB, svc contracts.EncryptionManager) {
|
||||
@@ -562,7 +577,7 @@ func TestIntegration_SecretsService(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
|
||||
ctx := context.Background()
|
||||
namespace := "test-namespace"
|
||||
namespace := xkube.Namespace("test-namespace")
|
||||
|
||||
// Here's what actually matters and varies on each test: look at the test case name.
|
||||
//
|
||||
|
||||
@@ -104,7 +104,7 @@ func (w *Worker) Cleanup(ctx context.Context, sv *secretv1beta1.SecureValue) err
|
||||
}
|
||||
|
||||
// Keeper deletion is idempotent
|
||||
if err := keeper.Delete(ctx, keeperCfg, sv.Namespace, sv.Name, sv.Status.Version); err != nil {
|
||||
if err := keeper.Delete(ctx, keeperCfg, xkube.Namespace(sv.Namespace), sv.Name, sv.Status.Version); err != nil {
|
||||
return fmt.Errorf("deleting secure value from keeper: %w", err)
|
||||
}
|
||||
|
||||
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
secretv1beta1 "github.com/grafana/grafana/apps/secret/pkg/apis/secret/v1beta1"
|
||||
"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/xkube"
|
||||
"github.com/grafana/grafana/pkg/storage/secret/encryption"
|
||||
"github.com/mitchellh/copystructure"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -58,7 +59,7 @@ func TestBasic(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
|
||||
// Get the secret value once to make sure it's reachable
|
||||
exposedValue, err := keeper.Expose(t.Context(), keeperCfg, sv.Namespace, sv.Name, sv.Status.Version)
|
||||
exposedValue, err := keeper.Expose(t.Context(), keeperCfg, xkube.Namespace(sv.Namespace), sv.Name, sv.Status.Version)
|
||||
require.NoError(t, err)
|
||||
require.NotEmpty(t, exposedValue.DangerouslyExposeAndConsumeValue())
|
||||
|
||||
@@ -78,7 +79,7 @@ func TestBasic(t *testing.T) {
|
||||
require.Empty(t, svs)
|
||||
|
||||
// Try to get the secreet value again to make sure it's been deleted from the keeper
|
||||
exposedValue, err = keeper.Expose(t.Context(), keeperCfg, sv.Namespace, sv.Name, sv.Status.Version)
|
||||
exposedValue, err = keeper.Expose(t.Context(), keeperCfg, xkube.Namespace(sv.Namespace), sv.Name, sv.Status.Version)
|
||||
require.ErrorIs(t, err, encryption.ErrEncryptedValueNotFound)
|
||||
require.Empty(t, exposedValue)
|
||||
})
|
||||
|
||||
@@ -1,9 +1,12 @@
|
||||
package secretkeeper
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
|
||||
secretv1beta1 "github.com/grafana/grafana/apps/secret/pkg/apis/secret/v1beta1"
|
||||
"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/secretkeeper/sqlkeeper"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
@@ -20,11 +23,17 @@ func ProvideService(
|
||||
tracer trace.Tracer,
|
||||
store contracts.EncryptedValueStorage,
|
||||
encryptionManager contracts.EncryptionManager,
|
||||
migrationExecutor contracts.EncryptedValueMigrationExecutor,
|
||||
reg prometheus.Registerer,
|
||||
_ *secret.DependencyRegisterer, // noop import so wire runs DB migrations before instantiating this service -- can be nil when manually instantiating
|
||||
) (*OSSKeeperService, error) {
|
||||
systemKeeper, err := sqlkeeper.NewSQLKeeper(tracer, encryptionManager, store, migrationExecutor, reg)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create system keeper: %w", err)
|
||||
}
|
||||
|
||||
return &OSSKeeperService{
|
||||
// TODO: rename to system keeper or something like that
|
||||
systemKeeper: sqlkeeper.NewSQLKeeper(tracer, encryptionManager, store, reg),
|
||||
systemKeeper: systemKeeper,
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
osskmsproviders "github.com/grafana/grafana/pkg/registry/apis/secret/encryption/kmsproviders"
|
||||
"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/testutils"
|
||||
"github.com/grafana/grafana/pkg/services/sqlstore"
|
||||
"github.com/grafana/grafana/pkg/setting"
|
||||
"github.com/grafana/grafana/pkg/storage/secret/database"
|
||||
@@ -65,7 +66,8 @@ func setupTestService(t *testing.T, cfg *setting.Cfg) (*OSSKeeperService, error)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Initialize the keeper service
|
||||
keeperService, err := ProvideService(tracer, encValueStore, encryptionManager, nil)
|
||||
keeperService, err := ProvideService(tracer, encValueStore, encryptionManager, &testutils.NoopMigrationExecutor{}, nil, nil)
|
||||
require.NoError(t, err)
|
||||
|
||||
return keeperService, err
|
||||
}
|
||||
|
||||
@@ -5,9 +5,11 @@ import (
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-app-sdk/logging"
|
||||
secretv1beta1 "github.com/grafana/grafana/apps/secret/pkg/apis/secret/v1beta1"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/contracts"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/secretkeeper/metrics"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/xkube"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
@@ -26,20 +28,32 @@ func NewSQLKeeper(
|
||||
tracer trace.Tracer,
|
||||
encryptionManager contracts.EncryptionManager,
|
||||
store contracts.EncryptedValueStorage,
|
||||
migrationExecutor contracts.EncryptedValueMigrationExecutor,
|
||||
reg prometheus.Registerer,
|
||||
) *SQLKeeper {
|
||||
) (*SQLKeeper, error) {
|
||||
// Run the encrypted value store migration before anything else, otherwise operations may fail
|
||||
// TODO: This does not need to be here forever, but we may currently have on-prem deployments using GSM, so it needs to be here for now.
|
||||
// Periodically assess whether it is safe to remove - most likely for G13 should be fine.
|
||||
log := logging.FromContext(context.Background())
|
||||
log.Debug("sqlkeeper: executing encrypted value store migration")
|
||||
rowsAffected, err := migrationExecutor.Execute(context.Background())
|
||||
log.Debug("sqlkeeper: encrypted value store migration completed", "rows_affected", rowsAffected)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error encountered during encrypted value store migration: %w", err)
|
||||
}
|
||||
|
||||
return &SQLKeeper{
|
||||
tracer: tracer,
|
||||
encryptionManager: encryptionManager,
|
||||
store: store,
|
||||
metrics: metrics.NewKeeperMetrics(reg),
|
||||
}
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *SQLKeeper) Store(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace, name string, version int64, exposedValueOrRef string) (contracts.ExternalID, error) {
|
||||
func (s *SQLKeeper) Store(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace xkube.Namespace, name string, version int64, exposedValueOrRef string) (contracts.ExternalID, error) {
|
||||
ctx, span := s.tracer.Start(ctx, "SQLKeeper.Store",
|
||||
trace.WithAttributes(
|
||||
attribute.String("namespace", namespace),
|
||||
attribute.String("namespace", namespace.String()),
|
||||
attribute.String("name", name),
|
||||
attribute.Int64("version", version)),
|
||||
)
|
||||
@@ -63,9 +77,9 @@ func (s *SQLKeeper) Store(ctx context.Context, cfg secretv1beta1.KeeperConfig, n
|
||||
return contracts.ExternalID(""), nil
|
||||
}
|
||||
|
||||
func (s *SQLKeeper) Expose(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace, name string, version int64) (secretv1beta1.ExposedSecureValue, error) {
|
||||
func (s *SQLKeeper) Expose(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace xkube.Namespace, name string, version int64) (secretv1beta1.ExposedSecureValue, error) {
|
||||
ctx, span := s.tracer.Start(ctx, "SQLKeeper.Expose", trace.WithAttributes(
|
||||
attribute.String("namespace", namespace),
|
||||
attribute.String("namespace", namespace.String()),
|
||||
attribute.String("name", name),
|
||||
attribute.Int64("version", version),
|
||||
))
|
||||
@@ -77,7 +91,7 @@ func (s *SQLKeeper) Expose(ctx context.Context, cfg secretv1beta1.KeeperConfig,
|
||||
return "", fmt.Errorf("unable to get encrypted value: %w", err)
|
||||
}
|
||||
|
||||
exposedBytes, err := s.encryptionManager.Decrypt(ctx, namespace, encryptedValue.EncryptedData)
|
||||
exposedBytes, err := s.encryptionManager.Decrypt(ctx, namespace, encryptedValue.EncryptedPayload)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("unable to decrypt value: %w", err)
|
||||
}
|
||||
@@ -88,9 +102,9 @@ func (s *SQLKeeper) Expose(ctx context.Context, cfg secretv1beta1.KeeperConfig,
|
||||
return exposedValue, nil
|
||||
}
|
||||
|
||||
func (s *SQLKeeper) Delete(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace, name string, version int64) error {
|
||||
func (s *SQLKeeper) Delete(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace xkube.Namespace, name string, version int64) error {
|
||||
ctx, span := s.tracer.Start(ctx, "SQLKeeper.Delete", trace.WithAttributes(
|
||||
attribute.String("namespace", namespace),
|
||||
attribute.String("namespace", namespace.String()),
|
||||
attribute.String("name", name),
|
||||
attribute.Int64("version", version),
|
||||
))
|
||||
@@ -107,9 +121,9 @@ func (s *SQLKeeper) Delete(ctx context.Context, cfg secretv1beta1.KeeperConfig,
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *SQLKeeper) Update(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace, name string, version int64, exposedValueOrRef string) error {
|
||||
func (s *SQLKeeper) Update(ctx context.Context, cfg secretv1beta1.KeeperConfig, namespace xkube.Namespace, name string, version int64, exposedValueOrRef string) error {
|
||||
ctx, span := s.tracer.Start(ctx, "SQLKeeper.Update", trace.WithAttributes(
|
||||
attribute.String("namespace", namespace),
|
||||
attribute.String("namespace", namespace.String()),
|
||||
attribute.String("name", name),
|
||||
attribute.Int64("version", version),
|
||||
))
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
|
||||
secretv1beta1 "github.com/grafana/grafana/apps/secret/pkg/apis/secret/v1beta1"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/testutils"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/xkube"
|
||||
"github.com/grafana/grafana/pkg/tests/testsuite"
|
||||
)
|
||||
|
||||
@@ -16,10 +17,10 @@ func TestMain(m *testing.M) {
|
||||
}
|
||||
|
||||
func Test_SQLKeeperSetup(t *testing.T) {
|
||||
namespace1 := "namespace1"
|
||||
namespace1 := xkube.Namespace("namespace1")
|
||||
name1 := "name1"
|
||||
version1 := int64(1)
|
||||
namespace2 := "namespace2"
|
||||
namespace2 := xkube.Namespace("namespace2")
|
||||
name2 := "name2"
|
||||
plaintext1 := "very secret string in namespace 1"
|
||||
plaintext2 := "very secret string in namespace 2"
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
|
||||
"github.com/grafana/grafana-app-sdk/logging"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/contracts"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/xkube"
|
||||
otelcodes "go.opentelemetry.io/otel/codes"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
)
|
||||
@@ -60,21 +61,21 @@ func (s *ConsolidationService) Consolidate(ctx context.Context) (err error) {
|
||||
|
||||
for _, ev := range encryptedValues {
|
||||
// Decrypt the value using its old data key.
|
||||
decryptedValue, err := s.encryptionManager.Decrypt(ctx, ev.Namespace, ev.EncryptedData)
|
||||
decryptedValue, err := s.encryptionManager.Decrypt(ctx, xkube.Namespace(ev.Namespace), ev.EncryptedPayload)
|
||||
if err != nil {
|
||||
logging.FromContext(ctx).Error("Failed to decrypt value", "namespace", ev.Namespace, "name", ev.Name, "error", err)
|
||||
continue
|
||||
}
|
||||
|
||||
// Re-encrypt the value using a new data key.
|
||||
reEncryptedValue, err := s.encryptionManager.Encrypt(ctx, ev.Namespace, decryptedValue)
|
||||
reEncryptedValue, err := s.encryptionManager.Encrypt(ctx, xkube.Namespace(ev.Namespace), decryptedValue)
|
||||
if err != nil {
|
||||
logging.FromContext(ctx).Error("Failed to re-encrypt value", "namespace", ev.Namespace, "name", ev.Name, "error", err)
|
||||
continue
|
||||
}
|
||||
|
||||
// Update the encrypted value in the store.
|
||||
err = s.encryptedValueStore.Update(ctx, ev.Namespace, ev.Name, ev.Version, reEncryptedValue)
|
||||
err = s.encryptedValueStore.Update(ctx, xkube.Namespace(ev.Namespace), ev.Name, ev.Version, reEncryptedValue)
|
||||
if err != nil {
|
||||
logging.FromContext(ctx).Error("Failed to update encrypted value", "namespace", ev.Namespace, "name", ev.Name, "error", err)
|
||||
continue
|
||||
|
||||
@@ -97,7 +97,7 @@ func TestConsolidation(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
originalDecryptedValues = append(originalDecryptedValues, decryptedValue.DangerouslyExposeAndConsumeValue())
|
||||
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, tc.namespace, tc.name, 1)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, xkube.Namespace(tc.namespace), tc.name, 1)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, encryptedValue)
|
||||
originalEncryptedData = append(originalEncryptedData, encryptedValue.EncryptedData)
|
||||
@@ -115,7 +115,7 @@ func TestConsolidation(t *testing.T) {
|
||||
require.Equal(t, originalDecryptedValues[i], decryptedValue.DangerouslyExposeAndConsumeValue())
|
||||
|
||||
// Verify that the encrypted data has changed (indicating re-encryption)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, tc.namespace, tc.name, 1)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, xkube.Namespace(tc.namespace), tc.name, 1)
|
||||
require.NoError(t, err)
|
||||
require.NotEqual(t, originalEncryptedData[i], encryptedValue.EncryptedData)
|
||||
}
|
||||
@@ -174,7 +174,7 @@ func TestConsolidation(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
initialDecryptedValues = append(initialDecryptedValues, decryptedValue.DangerouslyExposeAndConsumeValue())
|
||||
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, tc.namespace, tc.name, 1)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, xkube.Namespace(tc.namespace), tc.name, 1)
|
||||
require.NoError(t, err)
|
||||
initialEncryptedData = append(initialEncryptedData, encryptedValue.EncryptedData)
|
||||
}
|
||||
@@ -223,7 +223,7 @@ func TestConsolidation(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
newSecretDecryptedValues = append(newSecretDecryptedValues, decryptedValue.DangerouslyExposeAndConsumeValue())
|
||||
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, tc.namespace, tc.name, 1)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, xkube.Namespace(tc.namespace), tc.name, 1)
|
||||
require.NoError(t, err)
|
||||
newSecretEncryptedData = append(newSecretEncryptedData, encryptedValue.EncryptedData)
|
||||
}
|
||||
@@ -252,7 +252,7 @@ func TestConsolidation(t *testing.T) {
|
||||
require.Equal(t, initialDecryptedValues[i], decryptedValue.DangerouslyExposeAndConsumeValue())
|
||||
|
||||
// Verify that the encrypted data has changed (indicating re-encryption)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, tc.namespace, tc.name, 1)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, xkube.Namespace(tc.namespace), tc.name, 1)
|
||||
require.NoError(t, err)
|
||||
require.NotEqual(t, initialEncryptedData[i], encryptedValue.EncryptedData)
|
||||
}
|
||||
@@ -275,7 +275,7 @@ func TestConsolidation(t *testing.T) {
|
||||
|
||||
// Verify that the encrypted data has changed from what it was when first created
|
||||
// (indicating it was re-encrypted during consolidation)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, tc.namespace, tc.name, 1)
|
||||
encryptedValue, err := sut.EncryptedValueStorage.Get(ctx, xkube.Namespace(tc.namespace), tc.name, 1)
|
||||
require.NoError(t, err)
|
||||
require.NotEqual(t, newSecretEncryptedData[i], encryptedValue.EncryptedData)
|
||||
}
|
||||
|
||||
@@ -146,7 +146,7 @@ func (s *SecureValueService) Update(ctx context.Context, newSecureValue *secretv
|
||||
}
|
||||
logging.FromContext(ctx).Debug("retrieved keeper", "namespace", newSecureValue.Namespace, "keeperName", newSecureValue.Spec.Keeper, "type", keeperCfg.Type())
|
||||
|
||||
secret, err := keeper.Expose(ctx, keeperCfg, newSecureValue.Namespace, newSecureValue.Name, currentVersion.Status.Version)
|
||||
secret, err := keeper.Expose(ctx, keeperCfg, xkube.Namespace(newSecureValue.Namespace), newSecureValue.Name, currentVersion.Status.Version)
|
||||
if err != nil {
|
||||
return nil, false, fmt.Errorf("reading secret value from keeper: %w", err)
|
||||
}
|
||||
@@ -191,7 +191,7 @@ func (s *SecureValueService) createNewVersion(ctx context.Context, sv *secretv1b
|
||||
// TODO: can we stop using external id?
|
||||
// TODO: store uses only the namespace and returns and id. It could be a kv instead.
|
||||
// TODO: check that the encrypted store works with multiple versions
|
||||
externalID, err := keeper.Store(ctx, keeperCfg, createdSv.Namespace, createdSv.Name, createdSv.Status.Version, sv.Spec.Value.DangerouslyExposeAndConsumeValue())
|
||||
externalID, err := keeper.Store(ctx, keeperCfg, xkube.Namespace(createdSv.Namespace), createdSv.Name, createdSv.Status.Version, sv.Spec.Value.DangerouslyExposeAndConsumeValue())
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("storing secure value in keeper: %w", err)
|
||||
}
|
||||
|
||||
@@ -126,7 +126,14 @@ func Setup(t *testing.T, opts ...func(*SetupConfig)) Sut {
|
||||
globalEncryptedValueStorage, err := encryptionstorage.ProvideGlobalEncryptedValueStorage(database, tracer)
|
||||
require.NoError(t, err)
|
||||
|
||||
sqlKeeper := sqlkeeper.NewSQLKeeper(tracer, encryptionManager, encryptedValueStorage, nil)
|
||||
// Initialize a noop migration executor for the sql keeper so it doesn't interfere with initialization
|
||||
noopMigrationExecutor := &NoopMigrationExecutor{}
|
||||
sqlKeeper, err := sqlkeeper.NewSQLKeeper(tracer, encryptionManager, encryptedValueStorage, noopMigrationExecutor, nil)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Initialize a real migration executor for test
|
||||
realMigrationExecutor, err := encryptionstorage.ProvideEncryptedValueMigrationExecutor(database, tracer, encryptedValueStorage, globalEncryptedValueStorage)
|
||||
require.NoError(t, err)
|
||||
|
||||
var keeperService contracts.KeeperService = newKeeperServiceWrapper(sqlKeeper)
|
||||
|
||||
@@ -158,39 +165,41 @@ func Setup(t *testing.T, opts ...func(*SetupConfig)) Sut {
|
||||
keeperService)
|
||||
|
||||
return Sut{
|
||||
SecureValueService: secureValueService,
|
||||
SecureValueMetadataStorage: secureValueMetadataStorage,
|
||||
DecryptStorage: decryptStorage,
|
||||
DecryptService: decryptService,
|
||||
EncryptedValueStorage: encryptedValueStorage,
|
||||
GlobalEncryptedValueStorage: globalEncryptedValueStorage,
|
||||
SQLKeeper: sqlKeeper,
|
||||
Database: database,
|
||||
AccessClient: accessClient,
|
||||
ConsolidationService: consolidationService,
|
||||
EncryptionManager: encryptionManager,
|
||||
GlobalDataKeyStore: globalDataKeyStore,
|
||||
GarbageCollectionWorker: garbageCollectionWorker,
|
||||
Clock: clock,
|
||||
KeeperService: keeperService,
|
||||
KeeperMetadataStorage: keeperMetadataStorage,
|
||||
SecureValueService: secureValueService,
|
||||
SecureValueMetadataStorage: secureValueMetadataStorage,
|
||||
DecryptStorage: decryptStorage,
|
||||
DecryptService: decryptService,
|
||||
EncryptedValueStorage: encryptedValueStorage,
|
||||
GlobalEncryptedValueStorage: globalEncryptedValueStorage,
|
||||
EncryptedValueMigrationExecutor: realMigrationExecutor,
|
||||
SQLKeeper: sqlKeeper,
|
||||
Database: database,
|
||||
AccessClient: accessClient,
|
||||
ConsolidationService: consolidationService,
|
||||
EncryptionManager: encryptionManager,
|
||||
GlobalDataKeyStore: globalDataKeyStore,
|
||||
GarbageCollectionWorker: garbageCollectionWorker,
|
||||
Clock: clock,
|
||||
KeeperService: keeperService,
|
||||
KeeperMetadataStorage: keeperMetadataStorage,
|
||||
}
|
||||
}
|
||||
|
||||
type Sut struct {
|
||||
SecureValueService contracts.SecureValueService
|
||||
SecureValueMetadataStorage contracts.SecureValueMetadataStorage
|
||||
DecryptStorage contracts.DecryptStorage
|
||||
DecryptService decryptcontracts.DecryptService
|
||||
EncryptedValueStorage contracts.EncryptedValueStorage
|
||||
GlobalEncryptedValueStorage contracts.GlobalEncryptedValueStorage
|
||||
SQLKeeper *sqlkeeper.SQLKeeper
|
||||
Database *database.Database
|
||||
AccessClient types.AccessClient
|
||||
ConsolidationService contracts.ConsolidationService
|
||||
EncryptionManager contracts.EncryptionManager
|
||||
GlobalDataKeyStore contracts.GlobalDataKeyStorage
|
||||
GarbageCollectionWorker *garbagecollectionworker.Worker
|
||||
SecureValueService contracts.SecureValueService
|
||||
SecureValueMetadataStorage contracts.SecureValueMetadataStorage
|
||||
DecryptStorage contracts.DecryptStorage
|
||||
DecryptService decryptcontracts.DecryptService
|
||||
EncryptedValueStorage contracts.EncryptedValueStorage
|
||||
GlobalEncryptedValueStorage contracts.GlobalEncryptedValueStorage
|
||||
EncryptedValueMigrationExecutor contracts.EncryptedValueMigrationExecutor
|
||||
SQLKeeper *sqlkeeper.SQLKeeper
|
||||
Database *database.Database
|
||||
AccessClient types.AccessClient
|
||||
ConsolidationService contracts.ConsolidationService
|
||||
EncryptionManager contracts.EncryptionManager
|
||||
GlobalDataKeyStore contracts.GlobalDataKeyStorage
|
||||
GarbageCollectionWorker *garbagecollectionworker.Worker
|
||||
// The fake clock passed to implementations to make testing easier
|
||||
Clock *FakeClock
|
||||
KeeperService contracts.KeeperService
|
||||
@@ -366,3 +375,10 @@ func (c *FakeClock) Now() time.Time {
|
||||
func (c *FakeClock) AdvanceBy(duration time.Duration) {
|
||||
c.Current = c.Current.Add(duration)
|
||||
}
|
||||
|
||||
type NoopMigrationExecutor struct {
|
||||
}
|
||||
|
||||
func (e *NoopMigrationExecutor) Execute(ctx context.Context) (int, error) {
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user