Skip to content
File

Blob: chord/local_kv_test.go

go613 lines
1package chord
2 
3import (
4 "bytes"
5 "container/ring"
6 "context"
7 "crypto/rand"
8 "fmt"
9 mathRand "math/rand"
10 "sort"
11 "testing"
12 "time"
13 
14 "go.miragespace.co/specter/spec/chord"
15 "go.miragespace.co/specter/util/testcond"
16 
17 "github.com/stretchr/testify/require"
18)
19 
20func makeKV(as *require.Assertions, num int, length int) (keys [][]byte, values [][]byte) {
21 keys = make([][]byte, num)
22 values = make([][]byte, num)
23 
24 var err error
25 for i := range keys {
26 keys[i] = make([]byte, length)
27 values[i] = make([]byte, length)
28 l := copy(keys[i], fmt.Appendf(nil, "key %d: ", i))
29 _, err = rand.Read(keys[i][l:])
30 as.NoError(err)
31 _, err = rand.Read(values[i])
32 as.NoError(err)
33 }
34 return
35}
36 
37func TestKVOperation(t *testing.T) {
38 as := require.New(t)
39 
40 nodes, done := makeRing(t, as, 5)
41 defer done()
42 
43 key := make([]byte, 16)
44 
45 for _, local := range nodes {
46 value := make([]byte, 16)
47 
48 _, err := rand.Read(key)
49 as.NoError(err)
50 _, err = rand.Read(value)
51 as.NoError(err)
52 
53 // Put
54 err = local.Put(context.Background(), key, value)
55 as.NoError(err)
56 
57 fsck(as, nodes)
58 
59 // Get
60 for _, remote := range nodes {
61 r, err := remote.Get(context.Background(), key)
62 as.NoError(err)
63 as.EqualValues(value, r)
64 }
65 
66 // Overwrite
67 _, err = rand.Read(value)
68 as.NoError(err)
69 err = local.Put(context.Background(), key, value)
70 as.NoError(err)
71 for _, remote := range nodes {
72 r, err := remote.Get(context.Background(), key)
73 as.NoError(err)
74 as.EqualValues(value, r)
75 }
76 
77 // Delete
78 err = local.Delete(context.Background(), key)
79 as.NoError(err)
80 for _, remote := range nodes {
81 r, err := remote.Get(context.Background(), key)
82 as.NoError(err)
83 as.Nil(r)
84 }
85 
86 // PrefixAppend
87 _, err = rand.Read(value)
88 as.NoError(err)
89 err = local.PrefixAppend(context.Background(), key, value)
90 as.NoError(err)
91 err = local.PrefixAppend(context.Background(), key, value)
92 as.ErrorIs(err, chord.ErrKVPrefixConflict)
93 for _, remote := range nodes {
94 // PrefixList
95 ret, err := remote.PrefixList(context.Background(), key)
96 as.NoError(err)
97 as.Len(ret, 1)
98 as.EqualValues(value, ret[0])
99 }
100 
101 // PrefixRemove
102 err = local.PrefixRemove(context.Background(), key, value)
103 as.NoError(err)
104 for _, remote := range nodes {
105 // PrefixList
106 ret, err := remote.PrefixList(context.Background(), key)
107 as.NoError(err)
108 as.Len(ret, 0)
109 }
110 err = local.PrefixAppend(context.Background(), key, value)
111 as.NoError(err)
112 }
113}
114 
115func kvFsck(kv chord.KVProvider, low, high uint64) bool {
116 valid := true
117 
118 keys, _ := kv.RangeKeys(context.Background(), 0, 0)
119 for _, key := range keys {
120 if !chord.Between(low, chord.Hash(key), high, true) {
121 valid = false
122 }
123 }
124 return valid
125}
126 
127func fsck(as *require.Assertions, nodes []*LocalNode) {
128 for _, node := range nodes {
129 pre := node.getPredecessor()
130 as.True(kvFsck(node.kv, pre.ID(), node.ID()), "node %d contains out of range keys", node.ID())
131 }
132}
133 
134func TestKeyTransferOut(t *testing.T) {
135 as := require.New(t)
136 
137 numNodes := 3
138 nodes, done := makeRing(t, as, numNodes)
139 defer done()
140 
141 keys, values := makeKV(as, 30, 8)
142 
143 for i := range keys {
144 as.Nil(nodes[0].Put(context.Background(), keys[i], values[i]))
145 }
146 
147 randomNode := nodes[mathRand.Intn(numNodes)]
148 
149 successor := randomNode.getSuccessor()
150 predecessor := randomNode.getPredecessor()
151 t.Logf("predecessor: %d, leaving: %d, successor: %d", predecessor.ID(), randomNode.ID(), successor.ID())
152 
153 leavingKeys, err := randomNode.kv.RangeKeys(context.Background(), 0, 0)
154 as.NoError(err)
155 
156 randomNode.Leave()
157 <-time.After(waitInterval)
158 
159 c := make([]*LocalNode, 0)
160 for _, node := range nodes {
161 if node == randomNode {
162 continue
163 }
164 c = append(c, node)
165 }
166 fsck(as, c)
167 
168 succVals, _ := successor.(*LocalNode).kv.Export(context.Background(), leavingKeys)
169 as.Len(succVals, len(leavingKeys))
170 
171 indices := make([]int, 0)
172 for _, k := range leavingKeys {
173 for i := range keys {
174 if string(keys[i]) == string(k) {
175 indices = append(indices, i)
176 }
177 }
178 }
179 as.Len(indices, len(leavingKeys))
180 
181 for i, v := range succVals {
182 as.EqualValues(values[indices[i]], v.GetSimpleValue())
183 }
184 
185 preVals, _ := predecessor.(*LocalNode).kv.Export(context.TODO(), leavingKeys)
186 for _, v := range preVals {
187 as.Nil(v.GetSimpleValue())
188 }
189}
190 
191func TestKeyTransferIn(t *testing.T) {
192 as := require.New(t)
193 
194 seedCfg := devConfig(t, as)
195 seedCfg.Identity.Id = chord.MaxIdentitifer / 2 // halfway
196 seed := NewLocalNode(seedCfg)
197 as.NoError(seed.Create())
198 defer seed.Leave()
199 waitRing(as, seed)
200 
201 keys, values := makeKV(as, 400, 8)
202 
203 for i := range keys {
204 err := seed.Put(context.Background(), keys[i], values[i])
205 as.NoError(err)
206 }
207 
208 n1Cfg := devConfig(t, as)
209 n1Cfg.Identity.Id = chord.MaxIdentitifer / 4 // 1 quarter
210 n1 := NewLocalNode(n1Cfg)
211 as.NoError(n1.Join(seed))
212 defer n1.Leave()
213 waitRing(as, n1)
214 
215 <-time.After(waitInterval * 2)
216 
217 keys, _ = n1.kv.RangeKeys(context.Background(), 0, 0)
218 as.Greater(len(keys), 0)
219 vals, _ := n1.kv.Export(context.Background(), keys)
220 for _, val := range vals {
221 as.Greater(len(val.GetSimpleValue()), 0)
222 }
223 
224 fsck(as, []*LocalNode{n1, seed})
225 
226 n2Cfg := devConfig(t, as)
227 n2Cfg.Identity.Id = (chord.MaxIdentitifer / 4 * 3) // 3 quarter
228 n2 := NewLocalNode(n2Cfg)
229 as.NoError(n2.Join(seed))
230 defer n2.Leave()
231 waitRing(as, n2)
232 
233 keys, _ = n2.kv.RangeKeys(context.Background(), 0, 0)
234 as.Greater(len(keys), 0)
235 vals, _ = n2.kv.Export(context.Background(), keys)
236 for _, val := range vals {
237 as.Greater(len(val.GetSimpleValue()), 0)
238 }
239 
240 fsck(as, []*LocalNode{n2, n1, seed})
241}
242 
243func TestListKeys(t *testing.T) {
244 as := require.New(t)
245 
246 numNodes := 3
247 nodes, done := makeRing(t, as, numNodes)
248 defer done()
249 
250 keys, values := makeKV(as, 30, 8)
251 
252 for i := range keys {
253 as.Nil(nodes[0].Put(context.Background(), keys[i], values[i]))
254 }
255 
256 composite, err := nodes[0].ListKeys(context.Background(), []byte(""))
257 as.NoError(err)
258 as.Len(composite, len(keys))
259 found := 0
260 for _, k1 := range composite {
261 for _, k2 := range keys {
262 if bytes.Equal(k1.GetKey(), k2) {
263 found++
264 }
265 }
266 }
267 as.Equal(len(keys), found)
268}
269 
270type concurrentTest struct {
271 numNodes int
272 numKeys int
273}
274 
275var concurrentParams = []concurrentTest{
276 // 64
277 {
278 numNodes: 64,
279 numKeys: 100,
280 },
281 {
282 numNodes: 64,
283 numKeys: 200,
284 },
285 {
286 numNodes: 64,
287 numKeys: 300,
288 },
289 {
290 numNodes: 64,
291 numKeys: 600,
292 },
293 // 128
294 {
295 numNodes: 128,
296 numKeys: 100,
297 },
298 {
299 numNodes: 128,
300 numKeys: 200,
301 },
302 {
303 numNodes: 128,
304 numKeys: 300,
305 },
306 {
307 numNodes: 128,
308 numKeys: 600,
309 },
310}
311 
312func TestConcurrentJoinKV(t *testing.T) {
313 if testing.Short() {
314 t.Skip("skipping many nodes concurrent join kv in short mode")
315 }
316 
317 for _, tc := range concurrentParams {
318 t.Run(fmt.Sprintf("test with %d nodes and %d keys", tc.numNodes, tc.numKeys), func(t *testing.T) {
319 t.Parallel()
320 concurrentJoinKVOps(t, tc.numNodes, tc.numKeys)
321 })
322 }
323}
324 
325func awaitStablizedGlobally(t *testing.T, as *require.Assertions, timeout time.Duration, nodes []*LocalNode) {
326 // create a sorted list of test nodes
327 refNodes := append([]*LocalNode{}, nodes...)
328 sort.SliceStable(refNodes, func(i, j int) bool {
329 return refNodes[i].ID() < refNodes[j].ID()
330 })
331 // use container/ring to build the stablized list of successors for each node
332 fullRing := ring.New(len(nodes))
333 succsMap := make(map[uint64]*ring.Ring)
334 for _, n := range refNodes {
335 fullRing.Value = n.ID()
336 fullRing = fullRing.Next()
337 succsMap[n.ID()] = fullRing
338 }
339 // then we check for condition, where every node has the correct local states
340 // with respect to global order of the ring
341 as.NoError(testcond.WaitForCondition(func() bool {
342 for _, node := range nodes {
343 if node.getPredecessor() == nil {
344 return false
345 }
346 if node.getSuccessor() == nil {
347 return false
348 }
349 }
350 // counter clockwise
351 for i := 0; i < len(refNodes)-1; i++ {
352 if refNodes[i].ID() != refNodes[i+1].getPredecessor().ID() {
353 return false
354 }
355 }
356 if refNodes[len(refNodes)-1].ID() != refNodes[0].getPredecessor().ID() {
357 return false
358 }
359 // clockwise
360 for i := 0; i < len(nodes)-1; i++ {
361 if refNodes[i+1].ID() != refNodes[i].getSuccessor().ID() {
362 return false
363 }
364 }
365 if refNodes[0].ID() != refNodes[len(refNodes)-1].getSuccessor().ID() {
366 return false
367 }
368 // ensure that successor list is correct (eventual consistency),
369 // otherwise lookup may go to the wrong node and failing the test incorrectly
370 for _, n := range nodes {
371 expectSuccs := succsMap[n.ID()]
372 actualSuccs := n.getSuccessors()
373 for _, s := range actualSuccs {
374 if s.ID() != expectSuccs.Value {
375 return false
376 }
377 expectSuccs = expectSuccs.Next()
378 }
379 }
380 t.Logf("[stablized] valid ring:\n")
381 fullRing.Do(func(a any) {
382 t.Logf(" %v\n", a)
383 })
384 return true
385 }, waitInterval, timeout))
386}
387 
388func concurrentJoinKVOps(t *testing.T, numNodes, numKeys int) {
389 as := require.New(t)
390 
391 // can't use makeRing here as we need to manually control joining
392 nodes := make([]*LocalNode, numNodes)
393 for i := range numNodes {
394 node := NewLocalNode(devConfig(t, as))
395 nodes[i] = node
396 }
397 
398 keys, values := makeKV(as, numKeys, 64)
399 syncA := make(chan struct{})
400 
401 nodes[0].Create()
402 
403 stale := 0
404 go func() {
405 defer close(syncA)
406 
407 for i := range keys {
408 RETRY:
409 err := nodes[0].Put(context.Background(), keys[i], values[i])
410 if err != nil {
411 if chord.ErrorIsRetryable(err) {
412 stale++
413 t.Logf("[put] outdated ownership at key %d", i)
414 time.Sleep(defaultInterval)
415 goto RETRY
416 }
417 as.NoError(err)
418 return
419 }
420 t.Logf("message %d inserted\n", i)
421 time.Sleep(defaultInterval) // used to pace insertions
422 }
423 }()
424 
425 for i := 1; i < numNodes; i++ {
426 as.NoError(nodes[i].Join(nodes[0]))
427 }
428 
429 <-syncA
430 
431 // wait until the ring is fully stablized before we check for missing values
432 awaitStablizedGlobally(t, as, time.Second*10, nodes)
433 
434 nodes[0].logger.Debug("Starting test validation")
435 
436 found := 0
437 missingIndices := make([]int, 0)
438 mismatchedIndices := make([]int, 0)
439 for i := range keys {
440 val, err := nodes[0].Get(context.Background(), keys[i])
441 as.NoError(err)
442 
443 if bytes.Equal(values[i], val) {
444 found++
445 } else if len(val) == 0 {
446 missingIndices = append(missingIndices, i)
447 } else {
448 mismatchedIndices = append(mismatchedIndices, i)
449 }
450 }
451 
452 t.Logf("stale ownership counts: %d", stale)
453 t.Logf("missing indices: %+v\n", missingIndices)
454 t.Logf("mismatched indices: %+v\n", mismatchedIndices)
455 
456 if len(missingIndices) > 0 {
457 for _, i := range missingIndices {
458 k := keys[i]
459 got, err := nodes[0].FindSuccessor(chord.Hash(k))
460 as.NoError(err)
461 as.NotNil(got)
462 for j := 1; j < numNodes; j++ {
463 v, _ := nodes[j].kv.Get(context.Background(), k)
464 if v != nil {
465 t.Logf("missing key index %d routed to node %d but found in node %d", i, got.ID(), nodes[j].ID())
466 }
467 }
468 }
469 }
470 
471 as.Equal(numKeys, found, "expect %d keys to be found, but only %d keys found with %d missing and %d mismatched", numKeys, found, len(missingIndices), len(mismatchedIndices))
472 
473 for i := numNodes - 1; i >= 0; i-- {
474 nodes[i].Leave()
475 }
476 
477 lostIndices := make([]int, 0)
478 k, _ := nodes[0].kv.RangeKeys(context.Background(), 0, 0)
479 if len(k) != numKeys {
480 for i := range keys {
481 val, err := nodes[0].kv.Get(context.Background(), keys[i])
482 as.NoError(err)
483 
484 if !bytes.Equal(values[i], val) {
485 lostIndices = append(lostIndices, i)
486 }
487 }
488 }
489 t.Logf("lost indices: %+v\n", lostIndices)
490 as.Equal(0, len(lostIndices), "expect no keys to be lost when one node remains, but %d keys were lost", len(lostIndices))
491}
492 
493func TestConcurrentLeaveKV(t *testing.T) {
494 if testing.Short() {
495 t.Skip("skipping many nodes concurrent leave kv in short mode")
496 }
497 
498 for _, tc := range concurrentParams {
499 t.Run(fmt.Sprintf("test with %d nodes and %d keys", tc.numNodes, tc.numKeys), func(t *testing.T) {
500 t.Parallel()
501 concurrentLeaveKVOps(t, tc.numNodes, tc.numKeys)
502 })
503 }
504}
505 
506func concurrentLeaveKVOps(t *testing.T, numNodes, numKeys int) {
507 as := require.New(t)
508 
509 nodes, done := makeRing(t, as, numNodes)
510 defer done()
511 
512 // wait until the ring is fully stablized before we insert values
513 awaitStablizedGlobally(t, as, time.Second*10, nodes)
514 
515 keys, values := makeKV(as, numKeys, 64)
516 syncA := make(chan struct{})
517 
518 stale := 0
519 go func() {
520 defer close(syncA)
521 
522 for i := range keys {
523 RETRY:
524 err := nodes[0].Put(context.Background(), keys[i], values[i])
525 if err != nil {
526 if chord.ErrorIsRetryable(err) {
527 stale++
528 t.Logf("[put] outdated ownership at key %d", i)
529 time.Sleep(defaultInterval)
530 goto RETRY
531 }
532 as.NoError(err)
533 return
534 }
535 t.Logf("message %d inserted\n", i)
536 time.Sleep(defaultInterval) // used to pace insertions
537 }
538 }()
539 
540 // kill every node except the first node
541 for i := 1; i < numNodes; i++ {
542 nodes[i].Leave()
543 }
544 
545 <-syncA
546 
547 nodes[0].logger.Debug("Starting test validation")
548 
549 found := 0
550 missingIndices := make([]int, 0)
551 mismatchedIndices := make([]int, 0)
552 for i := range keys {
553 // all keys should be in the first node
554 val, err := nodes[0].kv.Get(context.Background(), keys[i])
555 as.NoError(err)
556 if bytes.Equal(values[i], val) {
557 found++
558 } else if len(val) == 0 {
559 missingIndices = append(missingIndices, i)
560 } else {
561 mismatchedIndices = append(mismatchedIndices, i)
562 }
563 }
564 
565 t.Logf("stale ownership counts: %d", stale)
566 t.Logf("missing indices: %+v\n", missingIndices)
567 t.Logf("mismatched indices: %+v\n", mismatchedIndices)
568 
569 if len(missingIndices) > 0 {
570 for _, i := range missingIndices {
571 k := keys[i]
572 got, err := nodes[0].FindSuccessor(chord.Hash(k))
573 as.NoError(err)
574 as.NotNil(got)
575 for j := 1; j < numNodes; j++ {
576 v, _ := nodes[j].kv.Get(context.Background(), k)
577 if v != nil {
578 t.Logf("missing key index %d routed to node %d but found in node %d", i, got.ID(), nodes[j].ID())
579 }
580 }
581 }
582 }
583 
584 if len(mismatchedIndices) > 0 {
585 for _, i := range mismatchedIndices {
586 k := keys[i]
587 for j := 1; j < numNodes; j++ {
588 v, _ := nodes[j].kv.Get(context.Background(), k)
589 if v != nil {
590 t.Logf("mismatched key index %d found in node %d", i, nodes[j].ID())
591 }
592 }
593 }
594 }
595 
596 as.Equal(numKeys, found, "expect %d keys to be found, but only %d keys found with %d missing and %d mismatched", numKeys, found, len(missingIndices), len(mismatchedIndices))
597 
598 lostIndices := make([]int, 0)
599 k, _ := nodes[0].kv.RangeKeys(context.Background(), 0, 0)
600 if len(k) != numKeys {
601 for i := range keys {
602 val, err := nodes[0].kv.Get(context.Background(), keys[i])
603 as.NoError(err)
604 
605 if !bytes.Equal(values[i], val) {
606 lostIndices = append(lostIndices, i)
607 }
608 }
609 }
610 t.Logf("lost indices: %+v\n", lostIndices)
611 as.Equal(0, len(lostIndices), "expect no keys to be lost when one node remains, but %d keys were lost", len(lostIndices))
612}