diff --git a/pkg/storage/unified/resource/kv.go b/pkg/storage/unified/resource/kv.go index b9218411a51..bcea806689d 100644 --- a/pkg/storage/unified/resource/kv.go +++ b/pkg/storage/unified/resource/kv.go @@ -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() } diff --git a/pkg/storage/unified/testing/kv.go b/pkg/storage/unified/testing/kv.go index 59b1338d041..b65b769e8b1 100644 --- a/pkg/storage/unified/testing/kv.go +++ b/pkg/storage/unified/testing/kv.go @@ -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") }) }