SecretsManager: Add outbox store (#106613)
SecretsManager: add outbox store Co-authored-by: PoorlyDefinedBehaviour <brunotj2015@hotmail.com> Co-authored-by: Matheus Macabu <macabu@users.noreply.github.com>
This commit is contained in:
co-authored by
PoorlyDefinedBehaviour
Matheus Macabu
parent
5a34447ad9
commit
de28231f2f
@@ -0,0 +1,63 @@
|
||||
package contracts
|
||||
|
||||
import (
|
||||
"context"
|
||||
)
|
||||
|
||||
type contextRequestIdKey struct{}
|
||||
|
||||
type OutboxMessageType string
|
||||
|
||||
func GetRequestId(ctx context.Context) string {
|
||||
v := ctx.Value(contextRequestIdKey{})
|
||||
requestId, ok := v.(string)
|
||||
if !ok {
|
||||
return ""
|
||||
}
|
||||
|
||||
return requestId
|
||||
}
|
||||
|
||||
func ContextWithRequestID(ctx context.Context, requestId string) context.Context {
|
||||
return context.WithValue(ctx, contextRequestIdKey{}, requestId)
|
||||
}
|
||||
|
||||
const (
|
||||
CreateSecretOutboxMessage OutboxMessageType = "create"
|
||||
UpdateSecretOutboxMessage OutboxMessageType = "update"
|
||||
DeleteSecretOutboxMessage OutboxMessageType = "delete"
|
||||
)
|
||||
|
||||
type AppendOutboxMessage struct {
|
||||
RequestID string
|
||||
Type OutboxMessageType
|
||||
Name string
|
||||
Namespace string
|
||||
EncryptedSecret string
|
||||
KeeperName *string
|
||||
ExternalID *string
|
||||
}
|
||||
|
||||
type OutboxMessage struct {
|
||||
RequestID string
|
||||
Type OutboxMessageType
|
||||
MessageID string
|
||||
Name string
|
||||
Namespace string
|
||||
EncryptedSecret string
|
||||
KeeperName *string
|
||||
ExternalID *string
|
||||
// How many times this message has been received
|
||||
ReceiveCount int
|
||||
}
|
||||
|
||||
type OutboxQueue interface {
|
||||
// Appends a message to the outbox queue
|
||||
Append(ctx context.Context, message AppendOutboxMessage) (string, error)
|
||||
// Receives at most n messages from the outbox queue
|
||||
ReceiveN(ctx context.Context, n uint) ([]OutboxMessage, error)
|
||||
// Deletes a message from the outbox queue
|
||||
Delete(ctx context.Context, messageID string) error
|
||||
// Increments the number of times each message has been received by 1. Must be atomic.
|
||||
IncrementReceiveCount(ctx context.Context, messageIDs []string) error
|
||||
}
|
||||
@@ -421,6 +421,7 @@ var wireBasicSet = wire.NewSet(
|
||||
// Secrets Manager
|
||||
secretmetadata.ProvideSecureValueMetadataStorage,
|
||||
secretmetadata.ProvideKeeperMetadataStorage,
|
||||
secretmetadata.ProvideOutboxQueue,
|
||||
secretencryption.ProvideEncryptedValueStorage,
|
||||
secretmigrator.NewWithEngine,
|
||||
secretdatabase.ProvideDatabase,
|
||||
|
||||
@@ -0,0 +1,35 @@
|
||||
INSERT INTO {{ .Ident "secret_secure_value_outbox" }} (
|
||||
{{ .Ident "request_id" }},
|
||||
{{ .Ident "uid" }},
|
||||
{{ .Ident "message_type" }},
|
||||
{{ .Ident "name" }},
|
||||
{{ .Ident "namespace" }},
|
||||
{{ if .Row.EncryptedSecret.Valid }}
|
||||
{{ .Ident "encrypted_secret" }},
|
||||
{{ end }}
|
||||
{{ if .Row.KeeperName.Valid }}
|
||||
{{ .Ident "keeper_name" }},
|
||||
{{ end }}
|
||||
{{ if .Row.ExternalID.Valid }}
|
||||
{{ .Ident "external_id" }},
|
||||
{{ end }}
|
||||
{{ .Ident "receive_count" }},
|
||||
{{ .Ident "created" }}
|
||||
) VALUES (
|
||||
{{ .Arg .Row.RequestID }},
|
||||
{{ .Arg .Row.MessageID }},
|
||||
{{ .Arg .Row.MessageType }},
|
||||
{{ .Arg .Row.Name }},
|
||||
{{ .Arg .Row.Namespace }},
|
||||
{{ if .Row.EncryptedSecret.Valid }}
|
||||
{{ .Arg .Row.EncryptedSecret.String }},
|
||||
{{ end }}
|
||||
{{ if .Row.KeeperName.Valid }}
|
||||
{{ .Arg .Row.KeeperName.String }},
|
||||
{{ end }}
|
||||
{{ if .Row.ExternalID.Valid }}
|
||||
{{ .Arg .Row.ExternalID.String }},
|
||||
{{ end }}
|
||||
{{ .Arg .Row.ReceiveCount }},
|
||||
{{ .Arg .Row.Created }}
|
||||
);
|
||||
@@ -0,0 +1,5 @@
|
||||
DELETE FROM
|
||||
{{ .Ident "secret_secure_value_outbox" }}
|
||||
WHERE
|
||||
{{ .Ident "uid" }} = {{ .Arg .MessageID }}
|
||||
;
|
||||
@@ -0,0 +1,19 @@
|
||||
SELECT
|
||||
{{ .Ident "request_id" }},
|
||||
{{ .Ident "uid" }},
|
||||
{{ .Ident "message_type" }},
|
||||
{{ .Ident "name" }},
|
||||
{{ .Ident "namespace" }},
|
||||
{{ .Ident "encrypted_secret" }},
|
||||
{{ .Ident "keeper_name" }},
|
||||
{{ .Ident "external_id" }},
|
||||
{{ .Ident "receive_count" }},
|
||||
{{ .Ident "created" }}
|
||||
FROM
|
||||
{{ .Ident "secret_secure_value_outbox" }}
|
||||
ORDER BY
|
||||
{{ .Ident "created" }} ASC
|
||||
LIMIT
|
||||
{{ .Arg .ReceiveLimit }}
|
||||
{{ .SelectFor "UPDATE SKIP LOCKED" }}
|
||||
;
|
||||
@@ -0,0 +1,7 @@
|
||||
UPDATE
|
||||
{{ .Ident "secret_secure_value_outbox" }}
|
||||
SET
|
||||
{{ .Ident "receive_count" }} = {{ .Ident "receive_count" }} + 1
|
||||
WHERE
|
||||
{{ .Ident "uid" }} IN ({{ .ArgList .MessageIDs }})
|
||||
;
|
||||
@@ -0,0 +1,250 @@
|
||||
package metadata
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
unifiedsql "github.com/grafana/grafana/pkg/storage/unified/sql"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/assert"
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/contracts"
|
||||
"github.com/grafana/grafana/pkg/storage/unified/sql/sqltemplate"
|
||||
)
|
||||
|
||||
type outboxStore struct {
|
||||
db contracts.Database
|
||||
dialect sqltemplate.Dialect
|
||||
}
|
||||
|
||||
func ProvideOutboxQueue(db contracts.Database) contracts.OutboxQueue {
|
||||
return &outboxStore{
|
||||
db: db,
|
||||
dialect: sqltemplate.DialectForDriver(db.DriverName()),
|
||||
}
|
||||
}
|
||||
|
||||
type outboxMessageDB struct {
|
||||
RequestID string
|
||||
MessageID string
|
||||
MessageType contracts.OutboxMessageType
|
||||
Name string
|
||||
Namespace string
|
||||
EncryptedSecret sql.NullString
|
||||
KeeperName sql.NullString
|
||||
ExternalID sql.NullString
|
||||
ReceiveCount int
|
||||
Created int64
|
||||
}
|
||||
|
||||
func (s *outboxStore) Append(ctx context.Context, input contracts.AppendOutboxMessage) (string, error) {
|
||||
assert.True(input.Type != "", "outboxStore.Append: outbox message type is required")
|
||||
|
||||
messageID, err := s.insertMessage(ctx, input)
|
||||
if err != nil {
|
||||
return messageID, fmt.Errorf("inserting message into outbox table: %+w", err)
|
||||
}
|
||||
|
||||
return messageID, nil
|
||||
}
|
||||
|
||||
func (s *outboxStore) insertMessage(ctx context.Context, input contracts.AppendOutboxMessage) (string, error) {
|
||||
keeperName := sql.NullString{}
|
||||
if input.KeeperName != nil {
|
||||
keeperName = sql.NullString{
|
||||
Valid: true,
|
||||
String: *input.KeeperName,
|
||||
}
|
||||
}
|
||||
|
||||
externalID := sql.NullString{}
|
||||
if input.ExternalID != nil {
|
||||
externalID = sql.NullString{
|
||||
Valid: true,
|
||||
String: *input.ExternalID,
|
||||
}
|
||||
}
|
||||
|
||||
encryptedSecret := sql.NullString{}
|
||||
if input.Type == contracts.CreateSecretOutboxMessage || input.Type == contracts.UpdateSecretOutboxMessage {
|
||||
encryptedSecret = sql.NullString{
|
||||
Valid: true,
|
||||
String: input.EncryptedSecret,
|
||||
}
|
||||
}
|
||||
|
||||
messageID := uuid.New().String()
|
||||
|
||||
req := appendSecureValueOutbox{
|
||||
SQLTemplate: sqltemplate.New(s.dialect),
|
||||
Row: &outboxMessageDB{
|
||||
RequestID: input.RequestID,
|
||||
MessageID: messageID,
|
||||
MessageType: input.Type,
|
||||
Name: input.Name,
|
||||
Namespace: input.Namespace,
|
||||
EncryptedSecret: encryptedSecret,
|
||||
KeeperName: keeperName,
|
||||
ExternalID: externalID,
|
||||
ReceiveCount: 0,
|
||||
Created: time.Now().UTC().UnixMilli(),
|
||||
},
|
||||
}
|
||||
|
||||
query, err := sqltemplate.Execute(sqlSecureValueOutboxAppend, req)
|
||||
if err != nil {
|
||||
return messageID, fmt.Errorf("execute template %q: %w", sqlSecureValueOutboxAppend.Name(), err)
|
||||
}
|
||||
|
||||
result, err := s.db.ExecContext(ctx, query, req.GetArgs()...)
|
||||
if err != nil {
|
||||
if unifiedsql.IsRowAlreadyExistsError(err) {
|
||||
return messageID, contracts.ErrSecureValueOperationInProgress
|
||||
}
|
||||
return messageID, fmt.Errorf("inserting message into secure value outbox table: %w", err)
|
||||
}
|
||||
|
||||
rowsAffected, err := result.RowsAffected()
|
||||
if err != nil {
|
||||
return messageID, fmt.Errorf("get rows affected: %w", err)
|
||||
}
|
||||
|
||||
if rowsAffected != 1 {
|
||||
return messageID, fmt.Errorf("expected to affect 1 row, but affected %d", rowsAffected)
|
||||
}
|
||||
|
||||
return messageID, nil
|
||||
}
|
||||
|
||||
func (s *outboxStore) ReceiveN(ctx context.Context, n uint) ([]contracts.OutboxMessage, error) {
|
||||
req := receiveNSecureValueOutbox{
|
||||
SQLTemplate: sqltemplate.New(s.dialect),
|
||||
ReceiveLimit: n,
|
||||
}
|
||||
|
||||
query, err := sqltemplate.Execute(sqlSecureValueOutboxReceiveN, req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("execute template %q: %w", sqlSecureValueOutboxReceiveN.Name(), err)
|
||||
}
|
||||
|
||||
rows, err := s.db.QueryContext(ctx, query, req.GetArgs()...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("fetching rows from secure value outbox table: %w", err)
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
|
||||
messages := make([]contracts.OutboxMessage, 0)
|
||||
|
||||
for rows.Next() {
|
||||
var row outboxMessageDB
|
||||
if err := rows.Scan(
|
||||
&row.RequestID,
|
||||
&row.MessageID,
|
||||
&row.MessageType,
|
||||
&row.Name,
|
||||
&row.Namespace,
|
||||
&row.EncryptedSecret,
|
||||
&row.KeeperName,
|
||||
&row.ExternalID,
|
||||
&row.ReceiveCount,
|
||||
&row.Created,
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("scanning row from secure value outbox table: %w", err)
|
||||
}
|
||||
|
||||
var keeperName *string
|
||||
if row.KeeperName.Valid {
|
||||
keeperName = &row.KeeperName.String
|
||||
}
|
||||
|
||||
var externalID *string
|
||||
if row.ExternalID.Valid {
|
||||
externalID = &row.ExternalID.String
|
||||
}
|
||||
|
||||
msg := contracts.OutboxMessage{
|
||||
RequestID: row.RequestID,
|
||||
Type: row.MessageType,
|
||||
MessageID: row.MessageID,
|
||||
Name: row.Name,
|
||||
Namespace: row.Namespace,
|
||||
KeeperName: keeperName,
|
||||
ExternalID: externalID,
|
||||
ReceiveCount: row.ReceiveCount,
|
||||
}
|
||||
|
||||
if row.MessageType != contracts.DeleteSecretOutboxMessage && row.EncryptedSecret.Valid {
|
||||
msg.EncryptedSecret = row.EncryptedSecret.String
|
||||
}
|
||||
|
||||
messages = append(messages, msg)
|
||||
}
|
||||
|
||||
if err := rows.Err(); err != nil {
|
||||
return messages, fmt.Errorf("reading rows: %w", err)
|
||||
}
|
||||
|
||||
return messages, nil
|
||||
}
|
||||
|
||||
func (s *outboxStore) Delete(ctx context.Context, messageID string) error {
|
||||
assert.True(messageID != "", "outboxStore.Delete: messageID is required")
|
||||
|
||||
if err := s.deleteMessage(ctx, messageID); err != nil {
|
||||
return fmt.Errorf("deleting message from outbox table %+w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *outboxStore) deleteMessage(ctx context.Context, messageID string) error {
|
||||
req := deleteSecureValueOutbox{
|
||||
SQLTemplate: sqltemplate.New(s.dialect),
|
||||
MessageID: messageID,
|
||||
}
|
||||
|
||||
query, err := sqltemplate.Execute(sqlSecureValueOutboxDelete, req)
|
||||
if err != nil {
|
||||
return fmt.Errorf("execute template %q: %w", sqlSecureValueOutboxDelete.Name(), err)
|
||||
}
|
||||
|
||||
result, err := s.db.ExecContext(ctx, query, req.GetArgs()...)
|
||||
if err != nil {
|
||||
return fmt.Errorf("deleting message id=%v from secure value outbox table: %w", messageID, err)
|
||||
}
|
||||
|
||||
rowsAffected, err := result.RowsAffected()
|
||||
if err != nil {
|
||||
return fmt.Errorf("get rows affected: %w", err)
|
||||
}
|
||||
|
||||
if rowsAffected > 1 {
|
||||
return fmt.Errorf("bug: deleted more than one row from the outbox table, should delete only one at a time: deleted=%v", rowsAffected)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *outboxStore) IncrementReceiveCount(ctx context.Context, messageIDs []string) error {
|
||||
if len(messageIDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
req := incrementReceiveCountOutbox{
|
||||
SQLTemplate: sqltemplate.New(s.dialect),
|
||||
MessageIDs: messageIDs,
|
||||
}
|
||||
query, err := sqltemplate.Execute(sqlSecureValueOutboxUpdateReceiveCount, req)
|
||||
if err != nil {
|
||||
return fmt.Errorf("execute template %q: %w", sqlSecureValueOutboxUpdateReceiveCount.Name(), err)
|
||||
}
|
||||
|
||||
_, err = s.db.ExecContext(ctx, query, req.GetArgs()...)
|
||||
if err != nil {
|
||||
return fmt.Errorf("updating outbox messages receive count: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,263 @@
|
||||
package metadata
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"slices"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana/pkg/registry/apis/secret/contracts"
|
||||
"github.com/grafana/grafana/pkg/services/sqlstore"
|
||||
"github.com/grafana/grafana/pkg/storage/secret/database"
|
||||
"github.com/grafana/grafana/pkg/storage/secret/migrator"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
type outboxStoreModel struct {
|
||||
rows []contracts.OutboxMessage
|
||||
}
|
||||
|
||||
func newOutboxStoreModel() *outboxStoreModel {
|
||||
return &outboxStoreModel{}
|
||||
}
|
||||
|
||||
func (model *outboxStoreModel) Append(messageID string, message contracts.AppendOutboxMessage) {
|
||||
model.rows = append(model.rows, contracts.OutboxMessage{
|
||||
Type: message.Type,
|
||||
MessageID: messageID,
|
||||
Name: message.Name,
|
||||
Namespace: message.Namespace,
|
||||
EncryptedSecret: message.EncryptedSecret,
|
||||
KeeperName: message.KeeperName,
|
||||
ExternalID: message.ExternalID,
|
||||
})
|
||||
}
|
||||
|
||||
func (model *outboxStoreModel) ReceiveN(n uint) []contracts.OutboxMessage {
|
||||
maxMessages := min(len(model.rows), int(n))
|
||||
if maxMessages == 0 {
|
||||
return []contracts.OutboxMessage{}
|
||||
}
|
||||
return model.rows[:maxMessages]
|
||||
}
|
||||
|
||||
func (model *outboxStoreModel) Delete(messageID string) {
|
||||
oldLen := len(model.rows)
|
||||
model.rows = slices.DeleteFunc(model.rows, func(m contracts.OutboxMessage) bool {
|
||||
return m.MessageID == messageID
|
||||
})
|
||||
if len(model.rows) != oldLen-1 {
|
||||
panic("Delete: deleted more than one message")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOutboxStoreModel(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
model := newOutboxStoreModel()
|
||||
|
||||
require.Empty(t, model.ReceiveN(10))
|
||||
|
||||
appendOutboxMessage := contracts.AppendOutboxMessage{
|
||||
Type: contracts.CreateSecretOutboxMessage,
|
||||
Name: "s-1",
|
||||
Namespace: "n-1",
|
||||
EncryptedSecret: "value",
|
||||
ExternalID: nil,
|
||||
}
|
||||
|
||||
outboxMessage1 := contracts.OutboxMessage{
|
||||
MessageID: "message_id_1",
|
||||
Type: contracts.CreateSecretOutboxMessage,
|
||||
Name: "s-1",
|
||||
Namespace: "n-1",
|
||||
EncryptedSecret: "value",
|
||||
ExternalID: nil,
|
||||
}
|
||||
|
||||
outboxMessage2 := contracts.OutboxMessage{
|
||||
MessageID: "message_id_2",
|
||||
Type: contracts.CreateSecretOutboxMessage,
|
||||
Name: "s-1",
|
||||
Namespace: "n-1",
|
||||
EncryptedSecret: "value",
|
||||
ExternalID: nil,
|
||||
}
|
||||
|
||||
model.Append("message_id_1", appendOutboxMessage)
|
||||
|
||||
require.Equal(t, []contracts.OutboxMessage{outboxMessage1}, model.ReceiveN(10))
|
||||
|
||||
model.Append("message_id_2", appendOutboxMessage)
|
||||
|
||||
require.Equal(t, []contracts.OutboxMessage{outboxMessage1, outboxMessage2}, model.ReceiveN(10))
|
||||
|
||||
model.Delete(outboxMessage1.MessageID)
|
||||
|
||||
require.Equal(t, []contracts.OutboxMessage{outboxMessage2}, model.ReceiveN(5))
|
||||
|
||||
model.Delete(outboxMessage2.MessageID)
|
||||
|
||||
require.Empty(t, model.ReceiveN(1))
|
||||
}
|
||||
|
||||
func TestOutboxStoreSecureValueOperationInProgress(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
t.Run("Append returns error when the queue already contains an operation for the secure value", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
testDB := sqlstore.NewTestStore(t, sqlstore.WithMigrator(migrator.New()))
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
outbox := ProvideOutboxQueue(database.ProvideDatabase(testDB))
|
||||
|
||||
_, err := outbox.Append(ctx, contracts.AppendOutboxMessage{
|
||||
RequestID: "1",
|
||||
Type: contracts.CreateSecretOutboxMessage,
|
||||
Name: "name1",
|
||||
Namespace: "ns1",
|
||||
EncryptedSecret: "v1",
|
||||
KeeperName: nil,
|
||||
ExternalID: nil,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
_, err = outbox.Append(ctx, contracts.AppendOutboxMessage{
|
||||
RequestID: "1",
|
||||
Type: contracts.UpdateSecretOutboxMessage,
|
||||
Name: "name1",
|
||||
Namespace: "ns1",
|
||||
EncryptedSecret: "v1",
|
||||
KeeperName: nil,
|
||||
ExternalID: nil,
|
||||
})
|
||||
|
||||
require.ErrorIs(t, err, contracts.ErrSecureValueOperationInProgress)
|
||||
})
|
||||
}
|
||||
|
||||
func TestOutboxStore(t *testing.T) {
|
||||
testDB := sqlstore.NewTestStore(t, sqlstore.WithMigrator(migrator.New()))
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
outbox := ProvideOutboxQueue(database.ProvideDatabase(testDB))
|
||||
|
||||
m1 := contracts.AppendOutboxMessage{
|
||||
Type: contracts.CreateSecretOutboxMessage,
|
||||
Name: "s-1",
|
||||
Namespace: "n-1",
|
||||
EncryptedSecret: "value",
|
||||
ExternalID: nil,
|
||||
}
|
||||
m2 := contracts.AppendOutboxMessage{
|
||||
Type: contracts.CreateSecretOutboxMessage,
|
||||
Name: "s-1",
|
||||
Namespace: "n-2",
|
||||
EncryptedSecret: "value",
|
||||
ExternalID: nil,
|
||||
}
|
||||
|
||||
messages, err := outbox.ReceiveN(ctx, 10)
|
||||
require.NoError(t, err)
|
||||
require.Empty(t, messages)
|
||||
|
||||
messageID1, err := outbox.Append(ctx, m1)
|
||||
require.NoError(t, err)
|
||||
|
||||
for range 2 {
|
||||
messages, err = outbox.ReceiveN(ctx, 10)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 1, len(messages))
|
||||
require.Equal(t, messageID1, messages[0].MessageID)
|
||||
}
|
||||
|
||||
messageID2, err := outbox.Append(ctx, m2)
|
||||
require.NoError(t, err)
|
||||
|
||||
messages, err = outbox.ReceiveN(ctx, 3)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 2, len(messages))
|
||||
require.Equal(t, messageID1, messages[0].MessageID)
|
||||
require.Equal(t, messageID2, messages[1].MessageID)
|
||||
|
||||
require.NoError(t, outbox.Delete(ctx, messageID1))
|
||||
messages, err = outbox.ReceiveN(ctx, 10)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 1, len(messages))
|
||||
require.Equal(t, messageID2, messages[0].MessageID)
|
||||
|
||||
require.NoError(t, outbox.Delete(ctx, messageID2))
|
||||
|
||||
messages, err = outbox.ReceiveN(ctx, 100)
|
||||
require.NoError(t, err)
|
||||
require.Empty(t, messages)
|
||||
}
|
||||
|
||||
func TestOutboxStoreProperty(t *testing.T) {
|
||||
seed := time.Now().UnixMicro()
|
||||
rng := rand.New(rand.NewSource(seed))
|
||||
|
||||
defer func() {
|
||||
if err := recover(); err != nil || t.Failed() {
|
||||
fmt.Printf("TestOutboxStoreProperty: err=%+v\n\nSEED=%+v", err, seed)
|
||||
t.FailNow()
|
||||
}
|
||||
}()
|
||||
|
||||
// The number of iterations was decided arbitrarily based on the time the test takes to run
|
||||
for range 10 {
|
||||
testDB := sqlstore.NewTestStore(t, sqlstore.WithMigrator(migrator.New()))
|
||||
|
||||
outbox := ProvideOutboxQueue(database.ProvideDatabase(testDB))
|
||||
|
||||
model := newOutboxStoreModel()
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
for i := range 100 {
|
||||
n := rng.Intn(3)
|
||||
switch n {
|
||||
case 0:
|
||||
message := contracts.AppendOutboxMessage{
|
||||
Type: contracts.CreateSecretOutboxMessage,
|
||||
Name: fmt.Sprintf("s-%d", i),
|
||||
Namespace: fmt.Sprintf("n-%d", i),
|
||||
EncryptedSecret: "value",
|
||||
ExternalID: nil,
|
||||
}
|
||||
messageID, err := outbox.Append(ctx, message)
|
||||
require.NoError(t, err)
|
||||
|
||||
model.Append(messageID, message)
|
||||
|
||||
case 1:
|
||||
n := uint(rng.Intn(10))
|
||||
messages, err := outbox.ReceiveN(ctx, n)
|
||||
require.NoError(t, err)
|
||||
|
||||
modelMessages := model.ReceiveN(n)
|
||||
|
||||
require.Equal(t, len(modelMessages), len(messages))
|
||||
require.Equal(t, modelMessages, messages)
|
||||
|
||||
case 2:
|
||||
if len(model.rows) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
message := model.rows[rng.Intn(len(model.rows))]
|
||||
|
||||
model.Delete(message.MessageID)
|
||||
require.NoError(t, outbox.Delete(ctx, message.MessageID))
|
||||
|
||||
default:
|
||||
panic(fmt.Sprintf("unhandled action: %+v", n))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -22,6 +22,11 @@ var (
|
||||
sqlKeeperDelete = mustTemplate("keeper_delete.sql")
|
||||
|
||||
sqlKeeperListByName = mustTemplate("keeper_listByName.sql")
|
||||
|
||||
sqlSecureValueOutboxAppend = mustTemplate("secure_value_outbox_append.sql")
|
||||
sqlSecureValueOutboxReceiveN = mustTemplate("secure_value_outbox_receiveN.sql")
|
||||
sqlSecureValueOutboxDelete = mustTemplate("secure_value_outbox_delete.sql")
|
||||
sqlSecureValueOutboxUpdateReceiveCount = mustTemplate("secure_value_outbox_update_receive_count.sql")
|
||||
)
|
||||
|
||||
func mustTemplate(filename string) *template.Template {
|
||||
@@ -104,3 +109,34 @@ type listByNameKeeper struct {
|
||||
func (r listByNameKeeper) Validate() error {
|
||||
return nil // TODO
|
||||
}
|
||||
|
||||
/*************************************/
|
||||
/**-- Secure Value Outbox Queries --**/
|
||||
/*************************************/
|
||||
type appendSecureValueOutbox struct {
|
||||
sqltemplate.SQLTemplate
|
||||
Row *outboxMessageDB
|
||||
}
|
||||
|
||||
func (appendSecureValueOutbox) Validate() error { return nil }
|
||||
|
||||
type receiveNSecureValueOutbox struct {
|
||||
sqltemplate.SQLTemplate
|
||||
ReceiveLimit uint
|
||||
}
|
||||
|
||||
func (receiveNSecureValueOutbox) Validate() error { return nil }
|
||||
|
||||
type deleteSecureValueOutbox struct {
|
||||
sqltemplate.SQLTemplate
|
||||
MessageID string
|
||||
}
|
||||
|
||||
func (deleteSecureValueOutbox) Validate() error { return nil }
|
||||
|
||||
type incrementReceiveCountOutbox struct {
|
||||
sqltemplate.SQLTemplate
|
||||
MessageIDs []string
|
||||
}
|
||||
|
||||
func (incrementReceiveCountOutbox) Validate() error { return nil }
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package metadata
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"testing"
|
||||
"text/template"
|
||||
|
||||
@@ -106,3 +107,104 @@ func TestKeeperQueries(t *testing.T) {
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
func TestSecureValueOutboxQueries(t *testing.T) {
|
||||
mocks.CheckQuerySnapshots(t, mocks.TemplateTestSetup{
|
||||
RootDir: "testdata",
|
||||
Templates: map[*template.Template][]mocks.TemplateTestCase{
|
||||
sqlSecureValueOutboxUpdateReceiveCount: {
|
||||
{
|
||||
|
||||
Name: "update-receive-count",
|
||||
Data: &incrementReceiveCountOutbox{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
MessageIDs: []string{"id1", "id2", "id3"},
|
||||
},
|
||||
},
|
||||
},
|
||||
sqlSecureValueOutboxAppend: {
|
||||
{
|
||||
Name: "no-encrypted-secret",
|
||||
Data: &appendSecureValueOutbox{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
Row: &outboxMessageDB{
|
||||
MessageID: "my-uuid",
|
||||
MessageType: "some-type",
|
||||
Name: "name",
|
||||
Namespace: "namespace",
|
||||
ExternalID: sql.NullString{Valid: true, String: "external-id"},
|
||||
KeeperName: sql.NullString{Valid: true, String: "keeper"},
|
||||
Created: 1234,
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "no-external-id",
|
||||
Data: &appendSecureValueOutbox{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
Row: &outboxMessageDB{
|
||||
MessageID: "my-uuid",
|
||||
MessageType: "some-type",
|
||||
Name: "name",
|
||||
Namespace: "namespace",
|
||||
EncryptedSecret: sql.NullString{Valid: true, String: "encrypted"},
|
||||
KeeperName: sql.NullString{Valid: true, String: "keeper"},
|
||||
Created: 1234,
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "no-keeper-name",
|
||||
Data: &appendSecureValueOutbox{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
Row: &outboxMessageDB{
|
||||
MessageID: "my-uuid",
|
||||
MessageType: "some-type",
|
||||
Name: "name",
|
||||
Namespace: "namespace",
|
||||
EncryptedSecret: sql.NullString{Valid: true, String: "encrypted"},
|
||||
ExternalID: sql.NullString{Valid: true, String: "external-id"},
|
||||
Created: 1234,
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "all-fields-present",
|
||||
Data: &appendSecureValueOutbox{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
Row: &outboxMessageDB{
|
||||
MessageID: "my-uuid",
|
||||
MessageType: "some-type",
|
||||
Name: "name",
|
||||
Namespace: "namespace",
|
||||
EncryptedSecret: sql.NullString{Valid: true, String: "encrypted"},
|
||||
ExternalID: sql.NullString{Valid: true, String: ""}, // can be empty string
|
||||
KeeperName: sql.NullString{Valid: true, String: "keeper"},
|
||||
Created: 1234,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
|
||||
sqlSecureValueOutboxReceiveN: {
|
||||
{
|
||||
Name: "basic",
|
||||
Data: &receiveNSecureValueOutbox{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
ReceiveLimit: 10,
|
||||
},
|
||||
},
|
||||
},
|
||||
|
||||
sqlSecureValueOutboxDelete: {
|
||||
{
|
||||
Name: "basic",
|
||||
Data: &deleteSecureValueOutbox{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
MessageID: "my-uuid",
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
Vendored
Executable
+23
@@ -0,0 +1,23 @@
|
||||
INSERT INTO `secret_secure_value_outbox` (
|
||||
`request_id`,
|
||||
`uid`,
|
||||
`message_type`,
|
||||
`name`,
|
||||
`namespace`,
|
||||
`encrypted_secret`,
|
||||
`keeper_name`,
|
||||
`external_id`,
|
||||
`receive_count`,
|
||||
`created`
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'keeper',
|
||||
'',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO `secret_secure_value_outbox` (
|
||||
`request_id`,
|
||||
`uid`,
|
||||
`message_type`,
|
||||
`name`,
|
||||
`namespace`,
|
||||
`keeper_name`,
|
||||
`external_id`,
|
||||
`receive_count`,
|
||||
`created`
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'keeper',
|
||||
'external-id',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO `secret_secure_value_outbox` (
|
||||
`request_id`,
|
||||
`uid`,
|
||||
`message_type`,
|
||||
`name`,
|
||||
`namespace`,
|
||||
`encrypted_secret`,
|
||||
`keeper_name`,
|
||||
`receive_count`,
|
||||
`created`
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'keeper',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO `secret_secure_value_outbox` (
|
||||
`request_id`,
|
||||
`uid`,
|
||||
`message_type`,
|
||||
`name`,
|
||||
`namespace`,
|
||||
`encrypted_secret`,
|
||||
`external_id`,
|
||||
`receive_count`,
|
||||
`created`
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'external-id',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+5
@@ -0,0 +1,5 @@
|
||||
DELETE FROM
|
||||
`secret_secure_value_outbox`
|
||||
WHERE
|
||||
`uid` = 'my-uuid'
|
||||
;
|
||||
Vendored
Executable
+19
@@ -0,0 +1,19 @@
|
||||
SELECT
|
||||
`request_id`,
|
||||
`uid`,
|
||||
`message_type`,
|
||||
`name`,
|
||||
`namespace`,
|
||||
`encrypted_secret`,
|
||||
`keeper_name`,
|
||||
`external_id`,
|
||||
`receive_count`,
|
||||
`created`
|
||||
FROM
|
||||
`secret_secure_value_outbox`
|
||||
ORDER BY
|
||||
`created` ASC
|
||||
LIMIT
|
||||
10
|
||||
FOR UPDATE SKIP LOCKED
|
||||
;
|
||||
Vendored
Executable
+7
@@ -0,0 +1,7 @@
|
||||
UPDATE
|
||||
`secret_secure_value_outbox`
|
||||
SET
|
||||
`receive_count` = `receive_count` + 1
|
||||
WHERE
|
||||
`uid` IN ('id1', 'id2', 'id3')
|
||||
;
|
||||
Vendored
Executable
+23
@@ -0,0 +1,23 @@
|
||||
INSERT INTO "secret_secure_value_outbox" (
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"encrypted_secret",
|
||||
"keeper_name",
|
||||
"external_id",
|
||||
"receive_count",
|
||||
"created"
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'keeper',
|
||||
'',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO "secret_secure_value_outbox" (
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"keeper_name",
|
||||
"external_id",
|
||||
"receive_count",
|
||||
"created"
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'keeper',
|
||||
'external-id',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO "secret_secure_value_outbox" (
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"encrypted_secret",
|
||||
"keeper_name",
|
||||
"receive_count",
|
||||
"created"
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'keeper',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO "secret_secure_value_outbox" (
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"encrypted_secret",
|
||||
"external_id",
|
||||
"receive_count",
|
||||
"created"
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'external-id',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+5
@@ -0,0 +1,5 @@
|
||||
DELETE FROM
|
||||
"secret_secure_value_outbox"
|
||||
WHERE
|
||||
"uid" = 'my-uuid'
|
||||
;
|
||||
Vendored
Executable
+19
@@ -0,0 +1,19 @@
|
||||
SELECT
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"encrypted_secret",
|
||||
"keeper_name",
|
||||
"external_id",
|
||||
"receive_count",
|
||||
"created"
|
||||
FROM
|
||||
"secret_secure_value_outbox"
|
||||
ORDER BY
|
||||
"created" ASC
|
||||
LIMIT
|
||||
10
|
||||
FOR UPDATE SKIP LOCKED
|
||||
;
|
||||
Vendored
Executable
+7
@@ -0,0 +1,7 @@
|
||||
UPDATE
|
||||
"secret_secure_value_outbox"
|
||||
SET
|
||||
"receive_count" = "receive_count" + 1
|
||||
WHERE
|
||||
"uid" IN ('id1', 'id2', 'id3')
|
||||
;
|
||||
Vendored
Executable
+23
@@ -0,0 +1,23 @@
|
||||
INSERT INTO "secret_secure_value_outbox" (
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"encrypted_secret",
|
||||
"keeper_name",
|
||||
"external_id",
|
||||
"receive_count",
|
||||
"created"
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'keeper',
|
||||
'',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO "secret_secure_value_outbox" (
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"keeper_name",
|
||||
"external_id",
|
||||
"receive_count",
|
||||
"created"
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'keeper',
|
||||
'external-id',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO "secret_secure_value_outbox" (
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"encrypted_secret",
|
||||
"keeper_name",
|
||||
"receive_count",
|
||||
"created"
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'keeper',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+21
@@ -0,0 +1,21 @@
|
||||
INSERT INTO "secret_secure_value_outbox" (
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"encrypted_secret",
|
||||
"external_id",
|
||||
"receive_count",
|
||||
"created"
|
||||
) VALUES (
|
||||
'',
|
||||
'my-uuid',
|
||||
'some-type',
|
||||
'name',
|
||||
'namespace',
|
||||
'encrypted',
|
||||
'external-id',
|
||||
0,
|
||||
1234
|
||||
);
|
||||
Vendored
Executable
+5
@@ -0,0 +1,5 @@
|
||||
DELETE FROM
|
||||
"secret_secure_value_outbox"
|
||||
WHERE
|
||||
"uid" = 'my-uuid'
|
||||
;
|
||||
Vendored
Executable
+18
@@ -0,0 +1,18 @@
|
||||
SELECT
|
||||
"request_id",
|
||||
"uid",
|
||||
"message_type",
|
||||
"name",
|
||||
"namespace",
|
||||
"encrypted_secret",
|
||||
"keeper_name",
|
||||
"external_id",
|
||||
"receive_count",
|
||||
"created"
|
||||
FROM
|
||||
"secret_secure_value_outbox"
|
||||
ORDER BY
|
||||
"created" ASC
|
||||
LIMIT
|
||||
10
|
||||
;
|
||||
Vendored
Executable
+7
@@ -0,0 +1,7 @@
|
||||
UPDATE
|
||||
"secret_secure_value_outbox"
|
||||
SET
|
||||
"receive_count" = "receive_count" + 1
|
||||
WHERE
|
||||
"uid" IN ('id1', 'id2', 'id3')
|
||||
;
|
||||
@@ -12,8 +12,9 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
TableNameKeeper = "secret_keeper"
|
||||
TableNameEncryptedValue = "secret_encrypted_value"
|
||||
TableNameKeeper = "secret_keeper"
|
||||
TableNameSecureValueOutbox = "secret_secure_value_outbox"
|
||||
TableNameEncryptedValue = "secret_encrypted_value"
|
||||
)
|
||||
|
||||
type SecretDB struct {
|
||||
@@ -80,6 +81,29 @@ func (*SecretDB) AddMigration(mg *migrator.Migrator) {
|
||||
Indices: []*migrator.Index{}, // TODO: add indexes based on the queries we make.
|
||||
})
|
||||
|
||||
tables = append(tables, migrator.Table{
|
||||
Name: TableNameSecureValueOutbox,
|
||||
Columns: []*migrator.Column{
|
||||
{Name: "request_id", Type: migrator.DB_NVarchar, Length: 253, Nullable: false},
|
||||
{Name: "uid", Type: migrator.DB_NVarchar, Length: 36, IsPrimaryKey: true}, // Fixed size of a UUID.
|
||||
{Name: "message_type", Type: migrator.DB_NVarchar, Length: 16, Nullable: false},
|
||||
{Name: "name", Type: migrator.DB_NVarchar, Length: 253, Nullable: false}, // Limit enforced by K8s.
|
||||
{Name: "namespace", Type: migrator.DB_NVarchar, Length: 253, Nullable: false}, // Limit enforced by K8s.
|
||||
{Name: "encrypted_secret", Type: migrator.DB_Blob, Nullable: true},
|
||||
{Name: "keeper_name", Type: migrator.DB_NVarchar, Length: 253, Nullable: true}, // Keeper name, if not set, use default keeper.
|
||||
{Name: "external_id", Type: migrator.DB_NVarchar, Length: 36, Nullable: true}, // Fixed size of a UUID.
|
||||
{Name: "receive_count", Type: migrator.DB_SmallInt, Nullable: false},
|
||||
{Name: "created", Type: migrator.DB_BigInt, Nullable: false},
|
||||
},
|
||||
Indices: []*migrator.Index{
|
||||
// There's only one operation per secret in the queue at all times,
|
||||
// meaning the namespace + name combination should be unique
|
||||
{Cols: []string{"namespace", "name"}, Type: migrator.UniqueIndex},
|
||||
// Used for sorting
|
||||
{Cols: []string{"created"}, Type: migrator.IndexType},
|
||||
},
|
||||
})
|
||||
|
||||
// Initialize all tables
|
||||
for t := range tables {
|
||||
mg.AddMigration("drop table "+tables[t].Name, migrator.NewDropTableMigration(tables[t].Name))
|
||||
|
||||
Reference in New Issue
Block a user