forked from antflydb/antfly
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtransform_test.go
More file actions
362 lines (314 loc) · 13 KB
/
Copy pathtransform_test.go
File metadata and controls
362 lines (314 loc) · 13 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
// Copyright 2025 Antfly, Inc.
//
// Licensed under the Elastic License 2.0 (ELv2); you may not use this file
// except in compliance with the Elastic License 2.0. You may obtain a copy of
// the Elastic License 2.0 at
//
// https://www.antfly.io/licensing/ELv2-license
//
// Unless required by applicable law or agreed to in writing, software distributed
// under the Elastic License 2.0 is distributed on an "AS IS" BASIS, WITHOUT
// WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
// Elastic License 2.0 for the specific language governing permissions and
// limitations.
package e2e
import (
"sync"
"testing"
"time"
antfly "github.com/antflydb/antfly/pkg/client"
"github.com/stretchr/testify/require"
)
// Test configuration constants
const (
transformTestTableName = "transform_test_table"
transformTestNumShards = 4
)
// TestE2E_Transform_MaxKeepsLatestValue tests that using $max operator
// ensures the latest (highest) version value is always kept, regardless
// of operation order.
func TestE2E_Transform_MaxKeepsLatestValue(t *testing.T) {
skipInShortMode(t)
ctx := testContext(t, 3*time.Minute)
cluster := setupClusterWithTable(t, ctx, transformTestTableName, transformTestNumShards)
// Insert initial document with version 5
t.Log("Inserting initial document with version 5...")
initialDoc := map[string]any{
"name": "test-item",
"version": 5,
"data": "initial data",
}
_, err := cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Inserts: map[string]any{"item-1": initialDoc},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Failed to insert initial document")
err = cluster.WaitForKeyAvailable(ctx, transformTestTableName, "item-1", 10*time.Second)
require.NoError(t, err, "Key not available")
doc, err := cluster.Client.LookupKey(ctx, transformTestTableName, "item-1")
require.NoError(t, err, "Failed to lookup initial document")
require.Equal(t, float64(5), doc["version"], "Initial version should be 5")
// Use $max to try updating with lower version (should be ignored)
t.Log("Attempting to update with lower version 3 using $max (should be ignored)...")
_, err = cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Transforms: []antfly.Transform{
{
Key: "item-1",
Operations: []antfly.TransformOp{
{Op: antfly.TransformOpTypeMax, Path: "version", Value: 3},
{Op: antfly.TransformOpTypeSet, Path: "data", Value: "updated with v3"},
},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Transform failed")
doc, err = cluster.Client.LookupKey(ctx, transformTestTableName, "item-1")
require.NoError(t, err, "Failed to lookup document after low version update")
require.Equal(t, float64(5), doc["version"], "Version should still be 5 after $max with 3")
// Data should have been updated even though version wasn't (operations are independent)
require.Equal(t, "updated with v3", doc["data"], "Data should be updated")
// Use $max to update with higher version (should be applied)
t.Log("Updating with higher version 10 using $max (should be applied)...")
_, err = cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Transforms: []antfly.Transform{
{
Key: "item-1",
Operations: []antfly.TransformOp{
{Op: antfly.TransformOpTypeMax, Path: "version", Value: 10},
{Op: antfly.TransformOpTypeSet, Path: "data", Value: "updated with v10"},
},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Transform failed")
doc, err = cluster.Client.LookupKey(ctx, transformTestTableName, "item-1")
require.NoError(t, err, "Failed to lookup document after high version update")
require.Equal(t, float64(10), doc["version"], "Version should be 10 after $max with 10")
require.Equal(t, "updated with v10", doc["data"], "Data should be updated")
}
// TestE2E_Transform_ConcurrentMaxUpdates tests that concurrent updates
// using $max all converge to the highest version value.
func TestE2E_Transform_ConcurrentMaxUpdates(t *testing.T) {
skipInShortMode(t)
ctx := testContext(t, 3*time.Minute)
cluster := setupClusterWithTable(t, ctx, transformTestTableName, transformTestNumShards)
// Insert initial document with version 0
t.Log("Inserting initial document with version 0...")
_, err := cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Inserts: map[string]any{
"concurrent-item": map[string]any{
"name": "concurrent-test",
"version": 0,
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Failed to insert initial document")
err = cluster.WaitForKeyAvailable(ctx, transformTestTableName, "concurrent-item", 10*time.Second)
require.NoError(t, err, "Key not available")
// Launch concurrent updates with different versions
t.Log("Launching concurrent $max updates with versions 1-20...")
numUpdates := 20
var wg sync.WaitGroup
for version := 1; version <= numUpdates; version++ {
wg.Go(func() {
_, err := cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Transforms: []antfly.Transform{
{
Key: "concurrent-item",
Operations: []antfly.TransformOp{
{Op: antfly.TransformOpTypeMax, Path: "version", Value: version},
},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
if err != nil {
t.Logf("Update with version %d failed: %v", version, err)
}
})
}
wg.Wait()
// Verify the final version is the maximum (20)
t.Log("Verifying final version is 20...")
doc, err := cluster.Client.LookupKey(ctx, transformTestTableName, "concurrent-item")
require.NoError(t, err, "Failed to lookup document after concurrent updates")
require.Equal(t, float64(numUpdates), doc["version"],
"Version should be %d after concurrent $max updates", numUpdates)
}
// TestE2E_Transform_UpsertWithMax tests that $max works correctly with upsert
// to atomically create a document with a version if it doesn't exist,
// or update the version if the incoming value is higher.
func TestE2E_Transform_UpsertWithMax(t *testing.T) {
skipInShortMode(t)
ctx := testContext(t, 3*time.Minute)
cluster := setupClusterWithTable(t, ctx, transformTestTableName, transformTestNumShards)
// Use upsert to create a new document with version
t.Log("Creating new document using upsert with $max...")
_, err := cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Transforms: []antfly.Transform{
{
Key: "upsert-item",
Upsert: true,
Operations: []antfly.TransformOp{
{Op: antfly.TransformOpTypeMax, Path: "version", Value: 5},
{Op: antfly.TransformOpTypeSet, Path: "name", Value: "upserted-item"},
},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Upsert transform failed")
err = cluster.WaitForKeyAvailable(ctx, transformTestTableName, "upsert-item", 10*time.Second)
require.NoError(t, err, "Key not available after upsert")
doc, err := cluster.Client.LookupKey(ctx, transformTestTableName, "upsert-item")
require.NoError(t, err, "Failed to lookup upserted document")
require.Equal(t, float64(5), doc["version"], "Version should be 5")
require.Equal(t, "upserted-item", doc["name"], "Name should be set")
// Upsert again with lower version (should not change version)
t.Log("Upserting with lower version 3 (should not change version)...")
_, err = cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Transforms: []antfly.Transform{
{
Key: "upsert-item",
Upsert: true,
Operations: []antfly.TransformOp{
{Op: antfly.TransformOpTypeMax, Path: "version", Value: 3},
{Op: antfly.TransformOpTypeSet, Path: "status", Value: "updated"},
},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Second upsert failed")
doc, err = cluster.Client.LookupKey(ctx, transformTestTableName, "upsert-item")
require.NoError(t, err, "Failed to lookup document after second upsert")
require.Equal(t, float64(5), doc["version"], "Version should still be 5")
require.Equal(t, "updated", doc["status"], "Status should be updated")
// Upsert with higher version (should update version)
t.Log("Upserting with higher version 10 (should update version)...")
_, err = cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Transforms: []antfly.Transform{
{
Key: "upsert-item",
Upsert: true,
Operations: []antfly.TransformOp{
{Op: antfly.TransformOpTypeMax, Path: "version", Value: 10},
},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Third upsert failed")
doc, err = cluster.Client.LookupKey(ctx, transformTestTableName, "upsert-item")
require.NoError(t, err, "Failed to lookup document after third upsert")
require.Equal(t, float64(10), doc["version"], "Version should be 10")
}
// TestE2E_Transform_IncAtomicCounter tests that $inc operator provides
// atomic counter increments without race conditions.
func TestE2E_Transform_IncAtomicCounter(t *testing.T) {
skipInShortMode(t)
ctx := testContext(t, 3*time.Minute)
cluster := setupClusterWithTable(t, ctx, transformTestTableName, transformTestNumShards)
// Create document with counter at 0
t.Log("Creating document with counter at 0...")
_, err := cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Inserts: map[string]any{
"counter-item": map[string]any{
"name": "atomic-counter",
"counter": 0,
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Failed to create counter document")
err = cluster.WaitForKeyAvailable(ctx, transformTestTableName, "counter-item", 10*time.Second)
require.NoError(t, err, "Key not available")
// Launch concurrent increments
t.Log("Launching 50 concurrent $inc operations...")
numIncrements := 50
var wg sync.WaitGroup
for range numIncrements {
wg.Go(func() {
_, err := cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Transforms: []antfly.Transform{
{
Key: "counter-item",
Operations: []antfly.TransformOp{
{Op: antfly.TransformOpTypeInc, Path: "counter", Value: 1},
},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
if err != nil {
t.Logf("Increment failed: %v", err)
}
})
}
wg.Wait()
// Verify final counter value equals number of increments
t.Log("Verifying final counter value...")
doc, err := cluster.Client.LookupKey(ctx, transformTestTableName, "counter-item")
require.NoError(t, err, "Failed to lookup counter document")
require.Equal(t, float64(numIncrements), doc["counter"],
"Counter should be %d after %d concurrent increments", numIncrements, numIncrements)
}
// TestE2E_Transform_MultipleOperators tests combining multiple transform
// operators in a single batch to perform complex atomic updates.
func TestE2E_Transform_MultipleOperators(t *testing.T) {
skipInShortMode(t)
ctx := testContext(t, 3*time.Minute)
cluster := setupClusterWithTable(t, ctx, transformTestTableName, transformTestNumShards)
// Create initial document
t.Log("Creating initial document...")
_, err := cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Inserts: map[string]any{
"multi-op-item": map[string]any{
"name": "multi-operator-test",
"version": 1,
"views": 0,
"tags": []string{"initial"},
"metadata": map[string]any{"created": true},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Failed to create document")
err = cluster.WaitForKeyAvailable(ctx, transformTestTableName, "multi-op-item", 10*time.Second)
require.NoError(t, err, "Key not available")
// Apply multiple operators atomically
t.Log("Applying multiple operators atomically...")
_, err = cluster.Client.Batch(ctx, transformTestTableName, antfly.BatchRequest{
Transforms: []antfly.Transform{
{
Key: "multi-op-item",
Operations: []antfly.TransformOp{
{Op: antfly.TransformOpTypeMax, Path: "version", Value: 5},
{Op: antfly.TransformOpTypeInc, Path: "views", Value: 1},
{Op: antfly.TransformOpTypeAddToSet, Path: "tags", Value: "updated"},
{Op: antfly.TransformOpTypeSet, Path: "metadata.lastUpdated", Value: "2025-01-26"},
},
},
},
SyncLevel: antfly.SyncLevelWrite,
})
require.NoError(t, err, "Multi-operator transform failed")
// Verify all operators applied correctly
t.Log("Verifying all operations applied...")
doc, err := cluster.Client.LookupKey(ctx, transformTestTableName, "multi-op-item")
require.NoError(t, err, "Failed to lookup document")
require.Equal(t, float64(5), doc["version"], "Version should be 5 ($max)")
require.Equal(t, float64(1), doc["views"], "Views should be 1 ($inc)")
tags, ok := doc["tags"].([]any)
require.True(t, ok, "Tags should be an array")
require.Len(t, tags, 2, "Should have 2 tags")
require.Contains(t, tags, "initial", "Should contain 'initial' tag")
require.Contains(t, tags, "updated", "Should contain 'updated' tag ($addToSet)")
metadata, ok := doc["metadata"].(map[string]any)
require.True(t, ok, "Metadata should be a map")
require.Equal(t, true, metadata["created"], "Should preserve existing metadata")
require.Equal(t, "2025-01-26", metadata["lastUpdated"], "Should have lastUpdated ($set)")
}