Skip to content
File

Blob: kv/transfer_test.go

go154 lines
1package kv
2 
3import (
4 "context"
5 "fmt"
6 "testing"
7 "time"
8 
9 "go.miragespace.co/specter/kv/aof"
10 "go.miragespace.co/specter/kv/memory"
11 "go.miragespace.co/specter/kv/sqlite3"
12 "go.miragespace.co/specter/spec/chord"
13 "go.miragespace.co/specter/spec/protocol"
14 
15 "github.com/stretchr/testify/require"
16 "go.uber.org/zap"
17 "google.golang.org/protobuf/proto"
18)
19 
20func transferProvider(t *testing.T, name, dir string) (chord.KVProvider, func()) {
21 t.Helper()
22 switch name {
23 case "memory":
24 return memory.WithHashFn(chord.Hash), func() {}
25 case "aof":
26 store, err := aof.New(aof.Config{
27 Logger: zap.NewNop(), HasnFn: chord.Hash, DataDir: dir, FlushInterval: time.Second,
28 })
29 require.NoError(t, err)
30 go store.Start()
31 t.Cleanup(store.Stop)
32 return store, store.Stop
33 case "sqlite":
34 store, err := sqlite3.New(sqlite3.Config{
35 Logger: zap.NewNop(), HashFn: chord.Hash, DataDir: dir,
36 })
37 require.NoError(t, err)
38 t.Cleanup(store.Close)
39 return store, store.Close
40 default:
41 t.Fatalf("unknown provider %q", name)
42 return nil, nil
43 }
44}
45 
46// Membership transfers only live keys returned by RangeKeys. Test the actual
47// ImportRequest envelope, since direct provider calls preserve Go's nil/empty
48// distinction even when the wire format loses it.
49func TestSerializedTransferAcrossProviders(t *testing.T) {
50 type record struct {
51 key []byte
52 value []byte
53 prefix bool
54 lease uint64
55 }
56 providers := []string{"memory", "aof", "sqlite"}
57 for _, source := range providers {
58 for _, destination := range providers {
59 t.Run(source+"_to_"+destination, func(t *testing.T) {
60 ctx := context.Background()
61 from, _ := transferProvider(t, source, t.TempDir())
62 destinationDir := t.TempDir()
63 to, closeDestination := transferProvider(t, destination, destinationDir)
64 var records []record
65 var liveKeys [][]byte
66 var composites []*protocol.KeyComposite
67 for _, simple := range []struct {
68 name string
69 value []byte
70 }{{"absent", nil}, {"empty", []byte{}}, {"populated", []byte("value")}} {
71 for flags := range 4 {
72 r := record{
73 key: []byte(fmt.Sprintf("%s-%d", simple.name, flags)),
74 value: simple.value, prefix: flags&1 != 0,
75 }
76 if r.value != nil {
77 require.NoError(t, from.Put(ctx, r.key, r.value))
78 composites = append(composites, &protocol.KeyComposite{Key: r.key, Type: protocol.KeyComposite_SIMPLE})
79 }
80 if r.prefix {
81 require.NoError(t, from.PrefixAppend(ctx, r.key, []byte("child")))
82 composites = append(composites, &protocol.KeyComposite{Key: r.key, Type: protocol.KeyComposite_PREFIX})
83 }
84 if flags&2 != 0 {
85 var err error
86 r.lease, err = from.Acquire(ctx, r.key, time.Minute)
87 require.NoError(t, err)
88 composites = append(composites, &protocol.KeyComposite{Key: r.key, Type: protocol.KeyComposite_LEASE})
89 }
90 records = append(records, r)
91 if r.value != nil || r.prefix || r.lease != 0 {
92 liveKeys = append(liveKeys, r.key)
93 }
94 }
95 }
96 
97 keys, err := from.RangeKeys(ctx, 0, 0)
98 require.NoError(t, err)
99 require.ElementsMatch(t, liveKeys, keys)
100 values, err := from.Export(ctx, keys)
101 require.NoError(t, err)
102 envelope := &protocol.ImportRequest{Keys: keys, Values: values}
103 wire, err := envelope.MarshalVT()
104 require.NoError(t, err)
105 // Both codecs must agree on explicit empty presence.
106 var standard protocol.ImportRequest
107 require.NoError(t, proto.Unmarshal(wire, &standard))
108 require.True(t, proto.Equal(envelope, &standard))
109 var decoded protocol.ImportRequest
110 require.NoError(t, decoded.UnmarshalVT(wire))
111 require.NoError(t, to.Import(ctx, decoded.Keys, decoded.Values))
112 
113 verify := func() {
114 t.Helper()
115 for _, r := range records {
116 value, err := to.Get(ctx, r.key)
117 require.NoError(t, err)
118 require.Equal(t, r.value, value, "simple value for %s", r.key)
119 children, err := to.PrefixList(ctx, r.key)
120 require.NoError(t, err)
121 if r.prefix {
122 require.Equal(t, [][]byte{[]byte("child")}, children)
123 } else {
124 require.Empty(t, children)
125 }
126 }
127 transferredKeys, err := to.RangeKeys(ctx, 0, 0)
128 require.NoError(t, err)
129 require.ElementsMatch(t, liveKeys, transferredKeys)
130 listed, err := to.ListKeys(ctx, nil)
131 require.NoError(t, err)
132 require.ElementsMatch(t, composites, listed)
133 exported, err := to.Export(ctx, keys)
134 require.NoError(t, err)
135 require.True(t, proto.Equal(envelope, &protocol.ImportRequest{Keys: keys, Values: exported}))
136 }
137 verify()
138 if destination != "memory" {
139 // Exercise SQLite persistence and AOF's nested Import WAL codec.
140 closeDestination()
141 to, _ = transferProvider(t, destination, destinationDir)
142 verify()
143 }
144 for _, r := range records {
145 if r.lease != 0 {
146 // The transferred token, not a replacement lease, remains valid.
147 require.NoError(t, to.Release(ctx, r.key, r.lease))
148 }
149 }
150 })
151 }
152 }
153}