Skip to content
File

Blob: kv/aof/kv_test.go

go402 lines
1package aof
2 
3import (
4 "context"
5 "crypto/rand"
6 "io"
7 "io/fs"
8 "os"
9 "path/filepath"
10 "testing"
11 "time"
12 
13 "go.miragespace.co/specter/spec/chord"
14 
15 "github.com/stretchr/testify/require"
16 "go.uber.org/zap"
17 "go.uber.org/zap/zaptest"
18)
19 
20func TestStartStop(t *testing.T) {
21 as := require.New(t)
22 logger := zaptest.NewLogger(t, zaptest.WrapOptions(zap.AddCaller()))
23 
24 dir, err := os.MkdirTemp("", "aof")
25 as.NoError(err)
26 defer os.RemoveAll(dir)
27 
28 cfg := Config{
29 Logger: logger,
30 HasnFn: chord.Hash,
31 DataDir: dir,
32 FlushInterval: time.Millisecond * 500,
33 }
34 
35 kv, err := New(cfg)
36 as.NoError(err)
37 go kv.Start()
38 
39 key := make([]byte, 8)
40 value := make([]byte, 16)
41 
42 rand.Read(key)
43 rand.Read(value)
44 
45 err = kv.Put(context.Background(), key, value)
46 as.NoError(err)
47 
48 kv.Stop()
49 
50 kv, err = New(cfg)
51 as.NoError(err)
52 go kv.Start()
53 
54 val, err := kv.Get(context.Background(), key)
55 as.NoError(err)
56 as.Equal(value, val)
57 
58 rand.Read(key)
59 rand.Read(value)
60 
61 err = kv.Put(context.Background(), key, value)
62 as.NoError(err)
63 
64 kv.Stop()
65 
66 kv, err = New(cfg)
67 as.NoError(err)
68 
69 val, err = kv.Get(context.Background(), key)
70 as.NoError(err)
71 as.Equal(value, val)
72}
73 
74func TestEverything(t *testing.T) {
75 as := require.New(t)
76 logger := zaptest.NewLogger(t, zaptest.WrapOptions(zap.AddCaller()))
77 
78 dir, err := os.MkdirTemp("", "aof")
79 as.NoError(err)
80 defer os.RemoveAll(dir)
81 
82 cfg := Config{
83 Logger: logger,
84 HasnFn: chord.Hash,
85 DataDir: dir,
86 FlushInterval: time.Millisecond * 500,
87 }
88 
89 kv, err := New(cfg)
90 as.NoError(err)
91 go kv.Start()
92 
93 key := make([]byte, 8)
94 value := make([]byte, 16)
95 
96 rand.Read(key)
97 rand.Read(value)
98 
99 err = kv.Put(context.Background(), key, value)
100 as.NoError(err)
101 
102 err = kv.Put(context.Background(), value, key)
103 as.NoError(err)
104 
105 err = kv.Delete(context.Background(), value)
106 as.NoError(err)
107 
108 err = kv.PrefixAppend(context.Background(), key, value)
109 as.NoError(err)
110 
111 err = kv.PrefixAppend(context.Background(), value, key)
112 as.NoError(err)
113 
114 err = kv.PrefixRemove(context.Background(), value, key)
115 as.NoError(err)
116 
117 time.Sleep(cfg.FlushInterval * 2)
118 
119 keys, err := kv.RangeKeys(context.Background(), 0, 0)
120 as.NoError(err)
121 snapshot1, err := kv.Export(context.Background(), keys)
122 as.NoError(err)
123 
124 kv.Stop()
125 
126 // mutations after close should error
127 err = kv.Put(context.Background(), key, value)
128 as.ErrorIs(err, fs.ErrClosed)
129 err = kv.Delete(context.Background(), key)
130 as.ErrorIs(err, fs.ErrClosed)
131 err = kv.PrefixAppend(context.Background(), key, value)
132 as.ErrorIs(err, fs.ErrClosed)
133 err = kv.PrefixRemove(context.Background(), key, value)
134 as.ErrorIs(err, fs.ErrClosed)
135 err = kv.Import(context.Background(), keys, snapshot1)
136 as.ErrorIs(err, fs.ErrClosed)
137 err = kv.RemoveKeys(context.Background(), keys) // no-op
138 as.NoError(err)
139 
140 kv.Stop() // no-op
141 
142 kv, err = New(cfg)
143 as.NoError(err)
144 go kv.Start()
145 
146 val, err := kv.Get(context.Background(), key)
147 as.NoError(err)
148 as.Equal(value, val)
149 
150 ll, err := kv.PrefixList(context.Background(), key)
151 as.NoError(err)
152 as.Len(ll, 1)
153 as.EqualValues(value, ll[0])
154 
155 has, err := kv.PrefixContains(context.Background(), key, value)
156 as.NoError(err)
157 as.True(has)
158 
159 has, err = kv.PrefixContains(context.Background(), value, key)
160 as.NoError(err)
161 as.False(has)
162 
163 val, err = kv.Get(context.Background(), value)
164 as.NoError(err)
165 as.Nil(val)
166 
167 keys, err = kv.RangeKeys(context.Background(), 0, 0)
168 as.NoError(err)
169 snapshot2, err := kv.Export(context.Background(), keys)
170 as.NoError(err)
171 
172 as.EqualValues(snapshot1, snapshot2)
173 
174 kv.Stop()
175}
176 
177func TestImportAndRemoveKeys(t *testing.T) {
178 as := require.New(t)
179 logger := zaptest.NewLogger(t, zaptest.WrapOptions(zap.AddCaller()))
180 
181 dir1, err := os.MkdirTemp("", "aof")
182 as.NoError(err)
183 defer os.RemoveAll(dir1)
184 
185 cfg1 := Config{
186 Logger: logger,
187 HasnFn: chord.Hash,
188 DataDir: dir1,
189 FlushInterval: time.Millisecond * 500,
190 }
191 
192 kv, err := New(cfg1)
193 as.NoError(err)
194 go kv.Start()
195 
196 keys := make([][]byte, 16)
197 values := make([][]byte, 16)
198 for i := range keys {
199 keys[i] = make([]byte, 8)
200 values[i] = make([]byte, 16)
201 rand.Read(keys[i])
202 rand.Read(values[i])
203 }
204 
205 for i := range keys {
206 err := kv.Put(context.Background(), keys[i], values[i])
207 as.NoError(err)
208 }
209 
210 expKeys, err := kv.RangeKeys(context.Background(), 0, 0)
211 as.NoError(err)
212 snapshot, err := kv.Export(context.Background(), expKeys)
213 as.NoError(err)
214 
215 kv.Stop()
216 
217 dir2, err := os.MkdirTemp("", "aof")
218 as.NoError(err)
219 defer os.RemoveAll(dir2)
220 
221 cfg2 := Config{
222 Logger: logger,
223 HasnFn: chord.Hash,
224 DataDir: dir2,
225 FlushInterval: time.Millisecond * 500,
226 }
227 
228 kv, err = New(cfg2)
229 as.NoError(err)
230 go kv.Start()
231 
232 for i := range keys {
233 val, err := kv.Get(context.Background(), keys[i])
234 as.NoError(err)
235 as.Nil(val)
236 }
237 
238 err = kv.Import(context.Background(), expKeys, snapshot)
239 as.NoError(err)
240 
241 kv.Stop()
242 
243 kv, err = New(cfg2)
244 as.NoError(err)
245 go kv.Start()
246 
247 for i := range keys {
248 val, err := kv.Get(context.Background(), keys[i])
249 as.NoError(err)
250 as.EqualValues(values[i], val)
251 }
252 
253 err = kv.RemoveKeys(context.Background(), expKeys)
254 as.NoError(err)
255 
256 kv.Stop()
257 
258 kv, err = New(cfg2)
259 as.NoError(err)
260 go kv.Start()
261 
262 for i := range keys {
263 val, err := kv.Get(context.Background(), keys[i])
264 as.NoError(err)
265 as.Nil(val)
266 }
267 
268 kv.Stop()
269}
270 
271// Lease KV operations are volatile
272func TestVolatile(t *testing.T) {
273 as := require.New(t)
274 logger := zaptest.NewLogger(t, zaptest.WrapOptions(zap.AddCaller()))
275 
276 dir1, err := os.MkdirTemp("", "aof")
277 as.NoError(err)
278 defer os.RemoveAll(dir1)
279 
280 cfg1 := Config{
281 Logger: logger,
282 HasnFn: chord.Hash,
283 DataDir: dir1,
284 FlushInterval: time.Millisecond * 500,
285 }
286 
287 kv, err := New(cfg1)
288 as.NoError(err)
289 go kv.Start()
290 
291 token, err := kv.Acquire(context.Background(), []byte("lease"), time.Second)
292 as.NoError(err)
293 
294 token, err = kv.Renew(context.Background(), []byte("lease"), time.Second, token)
295 as.NoError(err)
296 
297 kv.Stop()
298 
299 kv, err = New(cfg1)
300 as.NoError(err)
301 go kv.Start()
302 
303 err = kv.Release(context.Background(), []byte("lease"), token)
304 as.ErrorIs(err, chord.ErrKVLeaseExpired)
305 
306 kv.Stop()
307}
308 
309func TestConflictRollback(t *testing.T) {
310 as := require.New(t)
311 logger := zaptest.NewLogger(t, zaptest.WrapOptions(zap.AddCaller()))
312 
313 dir1, err := os.MkdirTemp("", "aof")
314 as.NoError(err)
315 defer os.RemoveAll(dir1)
316 
317 cfg1 := Config{
318 Logger: logger,
319 HasnFn: chord.Hash,
320 DataDir: dir1,
321 FlushInterval: time.Millisecond * 500,
322 }
323 
324 kv, err := New(cfg1)
325 as.NoError(err)
326 go kv.Start()
327 
328 err = kv.PrefixAppend(context.Background(), []byte("prefix"), []byte("child"))
329 as.NoError(err)
330 err = kv.PrefixAppend(context.Background(), []byte("prefix"), []byte("child"))
331 as.ErrorIs(err, chord.ErrKVPrefixConflict)
332 
333 err = kv.PrefixAppend(context.Background(), []byte("prefix"), []byte("grandchild"))
334 as.NoError(err)
335 
336 kv.Stop()
337 
338 kv, err = New(cfg1)
339 as.NoError(err)
340 go kv.Start()
341 
342 has, err := kv.PrefixContains(context.Background(), []byte("prefix"), []byte("child"))
343 as.NoError(err)
344 as.True(has)
345 
346 has, err = kv.PrefixContains(context.Background(), []byte("prefix"), []byte("grandchild"))
347 as.NoError(err)
348 as.True(has)
349 
350 kv.Stop()
351}
352 
353func TestCorruptedLog(t *testing.T) {
354 as := require.New(t)
355 logger := zaptest.NewLogger(t, zaptest.WrapOptions(zap.AddCaller()))
356 
357 dir1, err := os.MkdirTemp("", "aof")
358 as.NoError(err)
359 defer os.RemoveAll(dir1)
360 
361 cfg1 := Config{
362 Logger: logger,
363 HasnFn: chord.Hash,
364 DataDir: dir1,
365 FlushInterval: time.Millisecond * 500,
366 }
367 
368 kv, err := New(cfg1)
369 as.NoError(err)
370 go kv.Start()
371 
372 err = kv.Put(context.Background(), []byte("key"), []byte("value"))
373 as.NoError(err)
374 
375 kv.Stop()
376 
377 files, err := os.ReadDir(logPath(dir1))
378 as.NoError(err)
379 
380 logFile := files[0]
381 f, err := os.OpenFile(filepath.Join(logPath(dir1), logFile.Name()), os.O_RDWR, logFile.Type())
382 as.NoError(err)
383 
384 info, err := logFile.Info()
385 as.NoError(err)
386 
387 buf := make([]byte, info.Size())
388 f.Read(buf)
389 t.Log(buf)
390 buf[10] = 0
391 buf[15] = 0
392 t.Log(buf)
393 f.Seek(0, io.SeekStart)
394 f.Write(buf)
395 f.Sync()
396 
397 f.Close()
398 
399 _, err = New(cfg1)
400 as.Error(err)
401}