Compare commits

..
Author SHA1 Message Date
Georges Chaudy 3c83c64d72 Implement batch operation support in KV storage
- Introduced a new Batch method to execute multiple operations atomically within a single transaction.
- Defined BatchOp types for various operations: put, create, update, and delete, with specific semantics for each.
- Updated the KV interface and implementation to support batch operations, ensuring rollback on failure.
- Enhanced testing suite to cover various batch scenarios, including success, failure, and edge cases, ensuring robust functionality.
2026-01-09 15:24:34 +01:00
2 changed files with 263 additions and 478 deletions
+74 -177
View File
@@ -35,108 +35,32 @@ type ListOptions struct {
Limit int64 // maximum number of results to return. 0 means no limit.
}
// CompareTarget specifies what to compare in a transaction
type CompareTarget int
// BatchOpMode controls the semantics of each operation in a batch
type BatchOpMode int
const (
CompareExists CompareTarget = iota // Check if key exists (Value: bool)
CompareValue // Compare actual value (Value: []byte)
// BatchOpPut performs an upsert: create or update (never fails on key state)
BatchOpPut BatchOpMode = iota
// BatchOpCreate creates a new key, fails if the key already exists
BatchOpCreate
// BatchOpUpdate updates an existing key, fails if the key doesn't exist
BatchOpUpdate
// BatchOpDelete removes a key, idempotent (never fails on key state)
BatchOpDelete
)
// CompareResult specifies the comparison operator
type CompareResult int
const (
CompareEqual CompareResult = iota
CompareNotEqual
CompareGreater
CompareLess
)
// Compare represents a single comparison in a transaction.
// Use the constructor functions CompareKeyExists, CompareKeyNotExists, and CompareKeyValue to create comparisons.
type Compare struct {
Key string
Target CompareTarget
Result CompareResult // Only used for CompareValue
Exists bool // Used when Target == CompareExists
Value []byte // Used when Target == CompareValue
}
// CompareKeyExists creates a comparison that succeeds if the key exists.
func CompareKeyExists(key string) Compare {
return Compare{Key: key, Target: CompareExists, Exists: true}
}
// CompareKeyNotExists creates a comparison that succeeds if the key does not exist.
func CompareKeyNotExists(key string) Compare {
return Compare{Key: key, Target: CompareExists, Exists: false}
}
// CompareKeyValue creates a comparison that compares the value of a key.
// The comparison succeeds if the stored value matches the expected value according to the result operator.
func CompareKeyValue(key string, result CompareResult, value []byte) Compare {
return Compare{Key: key, Target: CompareValue, Result: result, Value: value}
}
// TxnOpType specifies the type of operation in a transaction
type TxnOpType int
const (
TxnOpPut TxnOpType = iota
TxnOpDelete
)
// TxnOp represents an operation in a transaction.
// Use the constructor functions TxnPut and TxnDelete to create operations.
type TxnOp struct {
Type TxnOpType
// BatchOp represents a single operation in an atomic batch
type BatchOp struct {
Mode BatchOpMode
Key string
Value []byte // For Put operations
Value []byte // For Put/Create/Update operations, nil for Delete
}
// TxnPut creates a Put operation that stores a value at the given key.
func TxnPut(key string, value []byte) TxnOp {
return TxnOp{Type: TxnOpPut, Key: key, Value: value}
}
// Maximum limit for batch operations
const MaxBatchOps = 20
// TxnDelete creates a Delete operation that removes the given key.
func TxnDelete(key string) TxnOp {
return TxnOp{Type: TxnOpDelete, Key: key}
}
// TxnResponse contains the result of a transaction
type TxnResponse struct {
// Succeeded indicates whether the comparisons passed (true) or failed (false)
Succeeded bool
}
// Maximum limits for transaction operations
const (
MaxTxnCompares = 8
MaxTxnOps = 8
)
// ValidateTxnRequest validates the transaction request parameters.
func ValidateTxnRequest(section string, cmps []Compare, successOps []TxnOp, failureOps []TxnOp) error {
if section == "" {
return fmt.Errorf("section is required")
}
if len(cmps) > MaxTxnCompares {
return fmt.Errorf("too many comparisons: %d > %d", len(cmps), MaxTxnCompares)
}
if len(successOps) > MaxTxnOps {
return fmt.Errorf("too many success operations: %d > %d", len(successOps), MaxTxnOps)
}
if len(failureOps) > MaxTxnOps {
return fmt.Errorf("too many failure operations: %d > %d", len(failureOps), MaxTxnOps)
}
return nil
}
// ErrKeyAlreadyExists is returned when BatchOpCreate is used on an existing key
var ErrKeyAlreadyExists = errors.New("key already exists")
type KV interface {
// Keys returns all the keys in the store
@@ -164,10 +88,16 @@ type KV interface {
// This is used to ensure the server and client are not too far apart in time.
UnixTimestamp(ctx context.Context) (int64, error)
// Txn executes a transaction with compare-and-swap semantics.
// If all comparisons succeed, successOps are executed; otherwise failureOps are executed.
// Limited to MaxTxnCompares comparisons and MaxTxnOps operations each for success/failure.
Txn(ctx context.Context, section string, cmps []Compare, successOps []TxnOp, failureOps []TxnOp) (*TxnResponse, error)
// Batch executes all operations atomically within a single transaction.
// If any operation fails, all operations are rolled back.
// Operations are executed in order; the batch stops on first failure.
//
// Operation semantics:
// - BatchOpPut: Upsert (create or update), never fails on key state
// - BatchOpCreate: Fail with ErrKeyAlreadyExists if key exists
// - BatchOpUpdate: Fail with ErrNotFound if key doesn't exist
// - BatchOpDelete: Idempotent, never fails on key state
Batch(ctx context.Context, section string, ops []BatchOp) error
}
var _ KV = &badgerKV{}
@@ -469,101 +399,68 @@ func (k *badgerKV) BatchDelete(ctx context.Context, section string, keys []strin
return txn.Commit()
}
func (k *badgerKV) Txn(ctx context.Context, section string, cmps []Compare, successOps []TxnOp, failureOps []TxnOp) (*TxnResponse, error) {
func (k *badgerKV) Batch(ctx context.Context, section string, ops []BatchOp) error {
if k.db.IsClosed() {
return nil, fmt.Errorf("database is closed")
return fmt.Errorf("database is closed")
}
if err := ValidateTxnRequest(section, cmps, successOps, failureOps); err != nil {
return nil, err
if section == "" {
return fmt.Errorf("section is required")
}
if len(ops) > MaxBatchOps {
return fmt.Errorf("too many operations: %d > %d", len(ops), MaxBatchOps)
}
txn := k.db.NewTransaction(true)
defer txn.Discard()
// Evaluate all comparisons
succeeded := true
for _, cmp := range cmps {
keyWithSection := section + "/" + cmp.Key
item, err := txn.Get([]byte(keyWithSection))
keyExists := err == nil
if err != nil && !errors.Is(err, badger.ErrKeyNotFound) {
return nil, err
}
switch cmp.Target {
case CompareExists:
if keyExists != cmp.Exists {
succeeded = false
}
case CompareValue:
if !keyExists {
// Key doesn't exist, comparison fails unless comparing for not-equal
if cmp.Result != CompareNotEqual {
succeeded = false
}
} else {
itemValue, err := item.ValueCopy(nil)
if err != nil {
return nil, err
}
if !compareBytes(itemValue, cmp.Value, cmp.Result) {
succeeded = false
}
}
default:
return nil, fmt.Errorf("unknown compare target: %d", cmp.Target)
}
if !succeeded {
break
}
}
// Execute the appropriate operations
ops := successOps
if !succeeded {
ops = failureOps
}
for _, op := range ops {
keyWithSection := section + "/" + op.Key
switch op.Type {
case TxnOpPut:
switch op.Mode {
case BatchOpCreate:
// Check that key doesn't exist, then set
_, err := txn.Get([]byte(keyWithSection))
if err == nil {
return ErrKeyAlreadyExists
}
if !errors.Is(err, badger.ErrKeyNotFound) {
return err
}
if err := txn.Set([]byte(keyWithSection), op.Value); err != nil {
return nil, err
return err
}
case TxnOpDelete:
case BatchOpUpdate:
// Check that key exists, then set
_, err := txn.Get([]byte(keyWithSection))
if errors.Is(err, badger.ErrKeyNotFound) {
return ErrNotFound
}
if err != nil {
return err
}
if err := txn.Set([]byte(keyWithSection), op.Value); err != nil {
return err
}
case BatchOpPut:
// Upsert: create or update
if err := txn.Set([]byte(keyWithSection), op.Value); err != nil {
return err
}
case BatchOpDelete:
// Idempotent delete - don't error if not found
if err := txn.Delete([]byte(keyWithSection)); err != nil {
return nil, err
return err
}
default:
return nil, fmt.Errorf("unknown operation type: %d", op.Type)
return fmt.Errorf("unknown operation mode: %d", op.Mode)
}
}
if err := txn.Commit(); err != nil {
return nil, err
}
return &TxnResponse{Succeeded: succeeded}, nil
}
// compareBytes compares two byte slices based on the comparison result type
func compareBytes(a, b []byte, result CompareResult) bool {
cmp := bytes.Compare(a, b)
switch result {
case CompareEqual:
return cmp == 0
case CompareNotEqual:
return cmp != 0
case CompareGreater:
return cmp > 0
case CompareLess:
return cmp < 0
default:
return false
}
return txn.Commit()
}
+189 -301
View File
@@ -28,7 +28,7 @@ const (
TestKVUnixTimestamp = "unix timestamp"
TestKVBatchGet = "batch get operations"
TestKVBatchDelete = "batch delete operations"
TestKVTxn = "transaction operations"
TestKVBatch = "batch operations"
)
// NewKVFunc is a function that creates a new KV instance for testing
@@ -70,7 +70,7 @@ func RunKVTest(t *testing.T, newKV NewKVFunc, opts *KVTestOptions) {
{TestKVUnixTimestamp, runTestKVUnixTimestamp},
{TestKVBatchGet, runTestKVBatchGet},
{TestKVBatchDelete, runTestKVBatchDelete},
{TestKVTxn, runTestKVTxn},
{TestKVBatch, runTestKVBatch},
}
for _, tc := range cases {
@@ -804,150 +804,70 @@ func saveKVHelper(t *testing.T, kv resource.KV, ctx context.Context, section, ke
require.NoError(t, err)
}
func runTestKVTxn(t *testing.T, kv resource.KV, nsPrefix string) {
func runTestKVBatch(t *testing.T, kv resource.KV, nsPrefix string) {
ctx := testutil.NewTestContext(t, time.Now().Add(30*time.Second))
section := nsPrefix + "-txn"
section := nsPrefix + "-batch"
t.Run("txn with empty section", func(t *testing.T) {
_, err := kv.Txn(ctx, "", nil, nil, nil)
t.Run("batch with empty section", func(t *testing.T) {
err := kv.Batch(ctx, "", nil)
assert.Error(t, err)
assert.Contains(t, err.Error(), "section is required")
})
t.Run("txn with no comparisons executes success ops", func(t *testing.T) {
// With no comparisons, all comparisons "pass" so success ops should run
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "no-cmp-key", Value: []byte("success-value")},
t.Run("batch with empty ops succeeds", func(t *testing.T) {
err := kv.Batch(ctx, section, nil)
require.NoError(t, err)
})
t.Run("batch put creates new key", func(t *testing.T) {
ops := []resource.BatchOp{
{Mode: resource.BatchOpPut, Key: "put-key", Value: []byte("put-value")},
}
resp, err := kv.Txn(ctx, section, nil, successOps, nil)
err := kv.Batch(ctx, section, ops)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
// Verify the key was created
reader, err := kv.Get(ctx, section, "no-cmp-key")
reader, err := kv.Get(ctx, section, "put-key")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "success-value", string(value))
assert.Equal(t, "put-value", string(value))
err = reader.Close()
require.NoError(t, err)
})
t.Run("txn compare exists false for non-existent key", func(t *testing.T) {
// Key doesn't exist, so exists should be false
cmps := []resource.Compare{
resource.CompareKeyNotExists("non-existent-key"),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "created-key", Value: []byte("created-value")},
}
failureOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "non-existent-failure-marker", Value: []byte("failure-executed")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, failureOps)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
// Verify the key was created
reader, err := kv.Get(ctx, section, "created-key")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "created-value", string(value))
err = reader.Close()
require.NoError(t, err)
// Verify failure op was not executed
_, err = kv.Get(ctx, section, "non-existent-failure-marker")
assert.Error(t, err)
})
t.Run("txn compare exists true for existing key", func(t *testing.T) {
t.Run("batch put updates existing key", func(t *testing.T) {
// First create a key
saveKVHelper(t, kv, ctx, section, "existing-key", strings.NewReader("existing-value"))
saveKVHelper(t, kv, ctx, section, "put-update-key", strings.NewReader("original-value"))
// Key exists, so exists should be true
cmps := []resource.Compare{
resource.CompareKeyExists("existing-key"),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "existing-key", Value: []byte("updated-value")},
}
failureOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "existing-failure-marker", Value: []byte("failure-executed")},
ops := []resource.BatchOp{
{Mode: resource.BatchOpPut, Key: "put-update-key", Value: []byte("updated-value")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, failureOps)
err := kv.Batch(ctx, section, ops)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
// Verify the key was updated
reader, err := kv.Get(ctx, section, "existing-key")
reader, err := kv.Get(ctx, section, "put-update-key")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "updated-value", string(value))
err = reader.Close()
require.NoError(t, err)
// Verify failure op was not executed
_, err = kv.Get(ctx, section, "existing-failure-marker")
assert.Error(t, err)
})
t.Run("txn compare exists fails when key exists but expected not to", func(t *testing.T) {
// First create a key
saveKVHelper(t, kv, ctx, section, "exists-fail-key", strings.NewReader("some-value"))
// Key exists but we expect it not to
cmps := []resource.Compare{
resource.CompareKeyNotExists("exists-fail-key"),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "exists-fail-result", Value: []byte("should-not-exist")},
}
failureOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "exists-fail-marker", Value: []byte("failure-executed")},
t.Run("batch create succeeds for new key", func(t *testing.T) {
ops := []resource.BatchOp{
{Mode: resource.BatchOpCreate, Key: "create-new-key", Value: []byte("new-value")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, failureOps)
err := kv.Batch(ctx, section, ops)
require.NoError(t, err)
assert.False(t, resp.Succeeded)
// Verify success op was not executed
_, err = kv.Get(ctx, section, "exists-fail-result")
assert.Error(t, err)
assert.Equal(t, resource.ErrNotFound, err)
// Verify failure ops were executed
reader, err := kv.Get(ctx, section, "exists-fail-marker")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "failure-executed", string(value))
err = reader.Close()
require.NoError(t, err)
})
t.Run("txn compare value equals", func(t *testing.T) {
// First create a key with known value
saveKVHelper(t, kv, ctx, section, "value-cmp-key", strings.NewReader("expected-value"))
cmps := []resource.Compare{
resource.CompareKeyValue("value-cmp-key", resource.CompareEqual, []byte("expected-value")),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "value-cmp-key", Value: []byte("new-value")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, nil)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
// Verify the key was updated
reader, err := kv.Get(ctx, section, "value-cmp-key")
// Verify the key was created
reader, err := kv.Get(ctx, section, "create-new-key")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
@@ -956,218 +876,186 @@ func runTestKVTxn(t *testing.T, kv resource.KV, nsPrefix string) {
require.NoError(t, err)
})
t.Run("txn compare value fails executes failure ops", func(t *testing.T) {
// First create a key with known value
saveKVHelper(t, kv, ctx, section, "value-fail-key", strings.NewReader("actual-value"))
t.Run("batch create fails for existing key", func(t *testing.T) {
// First create a key
saveKVHelper(t, kv, ctx, section, "create-exists-key", strings.NewReader("existing-value"))
cmps := []resource.Compare{
resource.CompareKeyValue("value-fail-key", resource.CompareEqual, []byte("wrong-value")),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "value-fail-key", Value: []byte("should-not-be-set")},
}
failureOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "failure-marker", Value: []byte("failure-executed")},
ops := []resource.BatchOp{
{Mode: resource.BatchOpCreate, Key: "create-exists-key", Value: []byte("new-value")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, failureOps)
require.NoError(t, err)
assert.False(t, resp.Succeeded)
err := kv.Batch(ctx, section, ops)
assert.ErrorIs(t, err, resource.ErrKeyAlreadyExists)
// Verify the original key was not changed
reader, err := kv.Get(ctx, section, "value-fail-key")
// Verify the original value is unchanged
reader, err := kv.Get(ctx, section, "create-exists-key")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "actual-value", string(value))
assert.Equal(t, "existing-value", string(value))
err = reader.Close()
require.NoError(t, err)
})
t.Run("batch update succeeds for existing key", func(t *testing.T) {
// First create a key
saveKVHelper(t, kv, ctx, section, "update-exists-key", strings.NewReader("original-value"))
ops := []resource.BatchOp{
{Mode: resource.BatchOpUpdate, Key: "update-exists-key", Value: []byte("updated-value")},
}
err := kv.Batch(ctx, section, ops)
require.NoError(t, err)
// Verify the key was updated
reader, err := kv.Get(ctx, section, "update-exists-key")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "updated-value", string(value))
err = reader.Close()
require.NoError(t, err)
})
t.Run("batch update fails for non-existent key", func(t *testing.T) {
ops := []resource.BatchOp{
{Mode: resource.BatchOpUpdate, Key: "update-nonexistent-key", Value: []byte("new-value")},
}
err := kv.Batch(ctx, section, ops)
assert.ErrorIs(t, err, resource.ErrNotFound)
// Verify the key was not created
_, err = kv.Get(ctx, section, "update-nonexistent-key")
assert.ErrorIs(t, err, resource.ErrNotFound)
})
t.Run("batch delete removes existing key", func(t *testing.T) {
// First create a key
saveKVHelper(t, kv, ctx, section, "delete-exists-key", strings.NewReader("to-be-deleted"))
ops := []resource.BatchOp{
{Mode: resource.BatchOpDelete, Key: "delete-exists-key"},
}
err := kv.Batch(ctx, section, ops)
require.NoError(t, err)
// Verify the key was deleted
_, err = kv.Get(ctx, section, "delete-exists-key")
assert.ErrorIs(t, err, resource.ErrNotFound)
})
t.Run("batch delete is idempotent for non-existent key", func(t *testing.T) {
ops := []resource.BatchOp{
{Mode: resource.BatchOpDelete, Key: "delete-nonexistent-key"},
}
err := kv.Batch(ctx, section, ops)
require.NoError(t, err) // Should succeed even though key doesn't exist
})
t.Run("batch multiple operations atomic success", func(t *testing.T) {
ops := []resource.BatchOp{
{Mode: resource.BatchOpPut, Key: "multi-key1", Value: []byte("value1")},
{Mode: resource.BatchOpPut, Key: "multi-key2", Value: []byte("value2")},
{Mode: resource.BatchOpPut, Key: "multi-key3", Value: []byte("value3")},
}
err := kv.Batch(ctx, section, ops)
require.NoError(t, err)
// Verify all keys were created
for i := 1; i <= 3; i++ {
key := fmt.Sprintf("multi-key%d", i)
reader, err := kv.Get(ctx, section, key)
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, fmt.Sprintf("value%d", i), string(value))
err = reader.Close()
require.NoError(t, err)
}
})
t.Run("batch multiple operations atomic rollback on failure", func(t *testing.T) {
// First create a key that will cause the batch to fail
saveKVHelper(t, kv, ctx, section, "rollback-exists", strings.NewReader("existing"))
ops := []resource.BatchOp{
{Mode: resource.BatchOpPut, Key: "rollback-new1", Value: []byte("value1")},
{Mode: resource.BatchOpCreate, Key: "rollback-exists", Value: []byte("should-fail")}, // This will fail
{Mode: resource.BatchOpPut, Key: "rollback-new2", Value: []byte("value2")},
}
err := kv.Batch(ctx, section, ops)
assert.ErrorIs(t, err, resource.ErrKeyAlreadyExists)
// Verify rollback: the first operation should NOT have persisted
_, err = kv.Get(ctx, section, "rollback-new1")
assert.ErrorIs(t, err, resource.ErrNotFound)
// Verify the third operation was not executed
_, err = kv.Get(ctx, section, "rollback-new2")
assert.ErrorIs(t, err, resource.ErrNotFound)
})
t.Run("batch mixed operations", func(t *testing.T) {
// Setup: create a key to update and one to delete
saveKVHelper(t, kv, ctx, section, "mixed-update", strings.NewReader("original"))
saveKVHelper(t, kv, ctx, section, "mixed-delete", strings.NewReader("to-delete"))
ops := []resource.BatchOp{
{Mode: resource.BatchOpCreate, Key: "mixed-create", Value: []byte("created")},
{Mode: resource.BatchOpUpdate, Key: "mixed-update", Value: []byte("updated")},
{Mode: resource.BatchOpDelete, Key: "mixed-delete"},
{Mode: resource.BatchOpPut, Key: "mixed-put", Value: []byte("put")},
}
err := kv.Batch(ctx, section, ops)
require.NoError(t, err)
// Verify create
reader, err := kv.Get(ctx, section, "mixed-create")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "created", string(value))
err = reader.Close()
require.NoError(t, err)
// Verify failure ops were executed
reader, err = kv.Get(ctx, section, "failure-marker")
// Verify update
reader, err = kv.Get(ctx, section, "mixed-update")
require.NoError(t, err)
value, err = io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "failure-executed", string(value))
assert.Equal(t, "updated", string(value))
err = reader.Close()
require.NoError(t, err)
// Verify delete
_, err = kv.Get(ctx, section, "mixed-delete")
assert.ErrorIs(t, err, resource.ErrNotFound)
// Verify put
reader, err = kv.Get(ctx, section, "mixed-put")
require.NoError(t, err)
value, err = io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "put", string(value))
err = reader.Close()
require.NoError(t, err)
})
t.Run("txn delete operation", func(t *testing.T) {
// First create a key
saveKVHelper(t, kv, ctx, section, "delete-txn-key", strings.NewReader("to-be-deleted"))
cmps := []resource.Compare{
resource.CompareKeyExists("delete-txn-key"),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpDelete, Key: "delete-txn-key"},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, nil)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
// Verify the key was deleted
_, err = kv.Get(ctx, section, "delete-txn-key")
assert.Error(t, err)
assert.Equal(t, resource.ErrNotFound, err)
})
t.Run("txn multiple comparisons all must pass", func(t *testing.T) {
// Create two keys
saveKVHelper(t, kv, ctx, section, "multi-cmp-key1", strings.NewReader("value1"))
saveKVHelper(t, kv, ctx, section, "multi-cmp-key2", strings.NewReader("value2"))
cmps := []resource.Compare{
resource.CompareKeyValue("multi-cmp-key1", resource.CompareEqual, []byte("value1")),
resource.CompareKeyValue("multi-cmp-key2", resource.CompareEqual, []byte("value2")),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "multi-success", Value: []byte("both-passed")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, nil)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
// Verify success op was executed
reader, err := kv.Get(ctx, section, "multi-success")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "both-passed", string(value))
err = reader.Close()
require.NoError(t, err)
})
t.Run("txn multiple comparisons one fails", func(t *testing.T) {
// Create two keys
saveKVHelper(t, kv, ctx, section, "multi-fail-key1", strings.NewReader("value1"))
saveKVHelper(t, kv, ctx, section, "multi-fail-key2", strings.NewReader("value2"))
cmps := []resource.Compare{
resource.CompareKeyValue("multi-fail-key1", resource.CompareEqual, []byte("value1")),
resource.CompareKeyValue("multi-fail-key2", resource.CompareEqual, []byte("wrong-value")),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "multi-fail-success", Value: []byte("should-not-exist")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, nil)
require.NoError(t, err)
assert.False(t, resp.Succeeded)
// Verify success op was not executed
_, err = kv.Get(ctx, section, "multi-fail-success")
assert.Error(t, err)
assert.Equal(t, resource.ErrNotFound, err)
})
t.Run("txn too many comparisons", func(t *testing.T) {
cmps := make([]resource.Compare, resource.MaxTxnCompares+1)
for i := range cmps {
cmps[i] = resource.CompareKeyNotExists(fmt.Sprintf("key-%d", i))
}
_, err := kv.Txn(ctx, section, cmps, nil, nil)
assert.Error(t, err)
assert.Contains(t, err.Error(), "too many comparisons")
})
t.Run("txn too many success operations", func(t *testing.T) {
ops := make([]resource.TxnOp, resource.MaxTxnOps+1)
t.Run("batch too many operations", func(t *testing.T) {
ops := make([]resource.BatchOp, resource.MaxBatchOps+1)
for i := range ops {
ops[i] = resource.TxnOp{Type: resource.TxnOpPut, Key: fmt.Sprintf("key-%d", i), Value: []byte("value")}
ops[i] = resource.BatchOp{Mode: resource.BatchOpPut, Key: fmt.Sprintf("key-%d", i), Value: []byte("value")}
}
_, err := kv.Txn(ctx, section, nil, ops, nil)
err := kv.Batch(ctx, section, ops)
assert.Error(t, err)
assert.Contains(t, err.Error(), "too many success operations")
})
t.Run("txn too many failure operations", func(t *testing.T) {
ops := make([]resource.TxnOp, resource.MaxTxnOps+1)
for i := range ops {
ops[i] = resource.TxnOp{Type: resource.TxnOpPut, Key: fmt.Sprintf("key-%d", i), Value: []byte("value")}
}
_, err := kv.Txn(ctx, section, nil, nil, ops)
assert.Error(t, err)
assert.Contains(t, err.Error(), "too many failure operations")
})
t.Run("txn compare greater than", func(t *testing.T) {
saveKVHelper(t, kv, ctx, section, "greater-key", strings.NewReader("bbb"))
cmps := []resource.Compare{
resource.CompareKeyValue("greater-key", resource.CompareGreater, []byte("aaa")),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "greater-result", Value: []byte("passed")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, nil)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
})
t.Run("txn compare less than", func(t *testing.T) {
saveKVHelper(t, kv, ctx, section, "less-key", strings.NewReader("aaa"))
cmps := []resource.Compare{
resource.CompareKeyValue("less-key", resource.CompareLess, []byte("bbb")),
}
successOps := []resource.TxnOp{
{Type: resource.TxnOpPut, Key: "less-result", Value: []byte("passed")},
}
resp, err := kv.Txn(ctx, section, cmps, successOps, nil)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
})
t.Run("txn compare not equal for non-existent value key", func(t *testing.T) {
// Key doesn't exist, comparing value should fail for equal but pass for not-equal
cmps := []resource.Compare{
resource.CompareKeyValue("non-existent-value-key", resource.CompareNotEqual, []byte("any-value")),
}
successOps := []resource.TxnOp{
resource.TxnPut("not-equal-result", []byte("passed")),
}
resp, err := kv.Txn(ctx, section, cmps, successOps, nil)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
})
t.Run("txn using constructor functions", func(t *testing.T) {
// Test that constructor functions work correctly
saveKVHelper(t, kv, ctx, section, "constructor-key", strings.NewReader("constructor-value"))
cmps := []resource.Compare{
resource.CompareKeyExists("constructor-key"),
}
successOps := []resource.TxnOp{
resource.TxnPut("constructor-new", []byte("new-value")),
resource.TxnDelete("constructor-key"),
}
resp, err := kv.Txn(ctx, section, cmps, successOps, nil)
require.NoError(t, err)
assert.True(t, resp.Succeeded)
// Verify the operations actually happened
_, err = kv.Get(ctx, section, "constructor-key")
assert.Error(t, err)
assert.Equal(t, resource.ErrNotFound, err)
reader, err := kv.Get(ctx, section, "constructor-new")
require.NoError(t, err)
value, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "new-value", string(value))
err = reader.Close()
require.NoError(t, err)
assert.Contains(t, err.Error(), "too many operations")
})
}