wal_test.mx raw
1 package wal
2
3 import (
4 "bytes"
5 "os"
6 "testing"
7 )
8
9 type walSink struct {
10 serials []uint64
11 data [][]byte
12 n int32
13 }
14
15 func newWalSink() (s *walSink) {
16 s = &walSink{
17 serials: []uint64{:32},
18 data: [][]byte{:32},
19 }
20 return
21 }
22
23 func (s *walSink) add(ser uint64, data []byte) (ok bool) {
24 if s.n >= int32(len(s.serials)) {
25 return false
26 }
27 cp := []byte{:len(data)}
28 copy(cp, data)
29 s.serials[s.n] = ser
30 s.data[s.n] = cp
31 s.n++
32 return true
33 }
34
35 func walTmp(t *testing.T) (dir string, ok bool) {
36 t.Helper()
37 d, err := os.MkdirTemp("", "wal-*")
38 if err != nil {
39 t.Fatal(err)
40 return "", false
41 }
42 return d, true
43 }
44
45 // TestAppendReadRoundTrip pins the serial encoding (segment<<32 | offset) and
46 // the length-prefixed entry format that every reader depends on.
47 func TestAppendReadRoundTrip(t *testing.T) {
48 dir, ok := walTmp(t)
49 if !ok {
50 return
51 }
52 defer os.RemoveAll(dir)
53
54 w, err := Open(dir)
55 if err != nil {
56 t.Fatal(err)
57 return
58 }
59 defer w.Close()
60
61 e1 := []byte("alpha")
62 e2 := []byte("beta")
63 e3 := []byte{:1000}
64 for i := range e3 {
65 e3[i] = byte(i)
66 }
67
68 s1, err1 := w.Append(e1)
69 if err1 != nil {
70 t.Fatal(err1)
71 return
72 }
73 s2, err2 := w.Append(e2)
74 if err2 != nil {
75 t.Fatal(err2)
76 return
77 }
78 s3, err3 := w.Append(e3)
79 if err3 != nil {
80 t.Fatal(err3)
81 return
82 }
83
84 // Segment 0, first byte at offset 0.
85 if s1 != 0 {
86 t.Fatalf("first serial = %d", s1)
87 }
88 if s2 != uint64(HdrSize+len(e1)) {
89 t.Fatalf("second serial = %d", s2)
90 }
91 if s3 != uint64(HdrSize+len(e1)+HdrSize+len(e2)) {
92 t.Fatalf("third serial = %d", s3)
93 }
94 if w.cur != 0 {
95 t.Fatalf("cur = %d", w.cur)
96 }
97 if w.off != int64(HdrSize+len(e1)+HdrSize+len(e2)+HdrSize+len(e3)) {
98 t.Fatalf("off = %d", w.off)
99 }
100
101 g1, gerr1 := w.Read(s1)
102 if gerr1 != nil {
103 t.Fatal(gerr1)
104 return
105 }
106 if !bytes.Equal(g1, e1) {
107 t.Fatal("entry 1 round trip")
108 }
109 g2, gerr2 := w.Read(s2)
110 if gerr2 != nil {
111 t.Fatal(gerr2)
112 return
113 }
114 if !bytes.Equal(g2, e2) {
115 t.Fatal("entry 2 round trip")
116 }
117 g3, gerr3 := w.Read(s3)
118 if gerr3 != nil {
119 t.Fatal(gerr3)
120 return
121 }
122 if !bytes.Equal(g3, e3) {
123 t.Fatal("entry 3 round trip")
124 }
125 }
126
127 func TestReadOutOfRange(t *testing.T) {
128 dir, ok := walTmp(t)
129 if !ok {
130 return
131 }
132 defer os.RemoveAll(dir)
133
134 w, err := Open(dir)
135 if err != nil {
136 t.Fatal(err)
137 return
138 }
139 defer w.Close()
140
141 // Empty segment: offset 0 has no header.
142 if _, e1 := w.Read(0); e1 == nil {
143 t.Fatal("read from an empty segment must fail")
144 }
145 // Segment index past the end.
146 if _, e2 := w.Read(uint64(1) << 32); e2 == nil {
147 t.Fatal("read from a missing segment must fail")
148 }
149 // Negative segment index after the int32 cast.
150 if _, e3 := w.Read(uint64(0xFFFFFFFFFFFFFFFF)); e3 == nil {
151 t.Fatal("read from a negative segment must fail")
152 }
153 }
154
155 // TestReadImpossibleLength plants a length prefix larger than MaxEntrySize and
156 // checks the guard rejects it instead of allocating the value.
157 func TestReadImpossibleLength(t *testing.T) {
158 dir, ok := walTmp(t)
159 if !ok {
160 return
161 }
162 defer os.RemoveAll(dir)
163
164 bad := []byte{0xFF, 0xFF, 0xFF, 0xFF}
165 if werr := os.WriteFile(dir|"/seg-000.dat", bad, 0644); werr != nil {
166 t.Fatal(werr)
167 return
168 }
169 w, err := Open(dir)
170 if err != nil {
171 t.Fatal(err)
172 return
173 }
174 defer w.Close()
175
176 if _, rerr := w.Read(0); rerr == nil {
177 t.Fatal("an impossible length must fail the read")
178 }
179 sink := newWalSink()
180 if ferr := w.ForEach(sink.add); ferr != nil {
181 t.Fatal(ferr)
182 return
183 }
184 if sink.n != 0 {
185 t.Fatalf("ForEach yielded %d entries from a bad header", sink.n)
186 }
187 }
188
189 func TestForEachOrder(t *testing.T) {
190 dir, ok := walTmp(t)
191 if !ok {
192 return
193 }
194 defer os.RemoveAll(dir)
195
196 w, err := Open(dir)
197 if err != nil {
198 t.Fatal(err)
199 return
200 }
201 defer w.Close()
202
203 names := [][]byte{[]byte("one"), []byte("two"), []byte("three")}
204 for _, n := range names {
205 if _, aerr := w.Append(n); aerr != nil {
206 t.Fatal(aerr)
207 return
208 }
209 }
210
211 sink := newWalSink()
212 if ferr := w.ForEach(sink.add); ferr != nil {
213 t.Fatal(ferr)
214 return
215 }
216 if sink.n != 3 {
217 t.Fatalf("ForEach count = %d", sink.n)
218 }
219 for i := 0; i < 3; i++ {
220 if !bytes.Equal(sink.data[i], names[i]) {
221 t.Fatalf("entry %d = %s", i, string(sink.data[i]))
222 }
223 }
224 if !(sink.serials[0] < sink.serials[1] && sink.serials[1] < sink.serials[2]) {
225 t.Fatal("ForEach serials not ascending")
226 }
227 }
228
229 func TestForEachFrom(t *testing.T) {
230 dir, ok := walTmp(t)
231 if !ok {
232 return
233 }
234 defer os.RemoveAll(dir)
235
236 w, err := Open(dir)
237 if err != nil {
238 t.Fatal(err)
239 return
240 }
241 defer w.Close()
242
243 a := []byte("aaa")
244 b := []byte("bbb")
245 c := []byte("ccc")
246 sa, ea := w.Append(a)
247 if ea != nil {
248 t.Fatal(ea)
249 return
250 }
251 sb, eb := w.Append(b)
252 if eb != nil {
253 t.Fatal(eb)
254 return
255 }
256 sc, ec := w.Append(c)
257 if ec != nil {
258 t.Fatal(ec)
259 return
260 }
261
262 // startSer == 0 iterates everything.
263 sink := newWalSink()
264 if ferr := w.ForEachFrom(0, sink.add); ferr != nil {
265 t.Fatal(ferr)
266 return
267 }
268 if sink.n != 3 {
269 t.Fatalf("ForEachFrom(0) count = %d", sink.n)
270 }
271
272 // startSer is exclusive: from b, only c follows.
273 sink2 := newWalSink()
274 if ferr2 := w.ForEachFrom(sb, sink2.add); ferr2 != nil {
275 t.Fatal(ferr2)
276 return
277 }
278 if sink2.n != 1 || !bytes.Equal(sink2.data[0], c) {
279 t.Fatalf("ForEachFrom(b) count = %d", sink2.n)
280 }
281 if sink2.serials[0] != sc {
282 t.Fatal("ForEachFrom(b) serial")
283 }
284
285 // From the last entry: nothing follows.
286 sink3 := newWalSink()
287 if ferr3 := w.ForEachFrom(sc, sink3.add); ferr3 != nil {
288 t.Fatal(ferr3)
289 return
290 }
291 if sink3.n != 0 {
292 t.Fatalf("ForEachFrom(c) count = %d", sink3.n)
293 }
294
295 // Serial 0 is the first entry's serial and also the "iterate all"
296 // sentinel: ForEachFrom(sa) must behave like ForEachFrom(0) and include a.
297 if sa != 0 {
298 t.Fatalf("first serial = %d, expected the serial-0 sentinel", sa)
299 }
300 sink4 := newWalSink()
301 if ferr4 := w.ForEachFrom(sa, sink4.add); ferr4 != nil {
302 t.Fatal(ferr4)
303 return
304 }
305 if sink4.n != 3 || !bytes.Equal(sink4.data[0], a) {
306 t.Fatalf("ForEachFrom(0 sentinel) count = %d", sink4.n)
307 }
308 }
309
310 func TestForEachFromStaleAndMissingSegment(t *testing.T) {
311 dir, ok := walTmp(t)
312 if !ok {
313 return
314 }
315 defer os.RemoveAll(dir)
316
317 w, err := Open(dir)
318 if err != nil {
319 t.Fatal(err)
320 return
321 }
322 defer w.Close()
323
324 if _, aerr := w.Append([]byte("payload")); aerr != nil {
325 t.Fatal(aerr)
326 return
327 }
328
329 sink := newWalSink()
330 // Offset well past the segment size means the checkpoint refers to a
331 // segment that was truncated: callers must see ErrCheckpointStale.
332 if ferr := w.ForEachFrom(1000, sink.add); ferr != ErrCheckpointStale {
333 t.Fatalf("stale offset error = %v", ferr)
334 }
335 // A segment index that does not exist yet is not stale, just empty.
336 sink2 := newWalSink()
337 if ferr2 := w.ForEachFrom(uint64(5)<<32, sink2.add); ferr2 != nil {
338 t.Fatal(ferr2)
339 return
340 }
341 if sink2.n != 0 {
342 t.Fatalf("missing segment yielded %d entries", sink2.n)
343 }
344 // An offset exactly at end of file skips nothing and yields nothing.
345 sink3 := newWalSink()
346 if ferr3 := w.ForEachFrom(w.off, sink3.add); ferr3 != nil {
347 t.Fatal(ferr3)
348 return
349 }
350 if sink3.n != 0 {
351 t.Fatalf("end-of-file offset yielded %d entries", sink3.n)
352 }
353 }
354
355 // TestForEachTruncatedEntry cuts the payload off the second entry and checks
356 // ForEach stops at the torn record instead of reading past it.
357 func TestForEachTruncatedEntry(t *testing.T) {
358 dir, ok := walTmp(t)
359 if !ok {
360 return
361 }
362 defer os.RemoveAll(dir)
363
364 w, err := Open(dir)
365 if err != nil {
366 t.Fatal(err)
367 return
368 }
369
370 first := []byte("first-payload")
371 second := []byte("second-payload")
372 if _, aerr := w.Append(first); aerr != nil {
373 t.Fatal(aerr)
374 return
375 }
376 full := w.off
377 if _, berr := w.Append(second); berr != nil {
378 t.Fatal(berr)
379 return
380 }
381 if terr := w.segs[0].Truncate(full + int64(HdrSize) + 2); terr != nil {
382 t.Fatal(terr)
383 return
384 }
385 w.Close()
386
387 // Reopen so the segment size is re-stat'ed from the truncated file.
388 w2, err2 := Open(dir)
389 if err2 != nil {
390 t.Fatal(err2)
391 return
392 }
393 defer w2.Close()
394
395 sink := newWalSink()
396 if ferr := w2.ForEach(sink.add); ferr != nil {
397 t.Fatal(ferr)
398 return
399 }
400 if sink.n != 1 || !bytes.Equal(sink.data[0], first) {
401 t.Fatalf("truncated ForEach count = %d", sink.n)
402 }
403 }
404
405 func TestReopenAppends(t *testing.T) {
406 dir, ok := walTmp(t)
407 if !ok {
408 return
409 }
410 defer os.RemoveAll(dir)
411
412 w1, err := Open(dir)
413 if err != nil {
414 t.Fatal(err)
415 return
416 }
417 a := []byte("entry-a")
418 sa, aerr := w1.Append(a)
419 if aerr != nil {
420 t.Fatal(aerr)
421 return
422 }
423 if serr := w1.Sync(); serr != nil {
424 t.Fatal(serr)
425 return
426 }
427 if cerr := w1.Close(); cerr != nil {
428 t.Fatal(cerr)
429 return
430 }
431
432 w2, err2 := Open(dir)
433 if err2 != nil {
434 t.Fatal(err2)
435 return
436 }
437 if len(w2.segs) != 1 {
438 t.Fatalf("reopened segment count = %d", int32(len(w2.segs)))
439 }
440 if w2.off != int64(HdrSize+len(a)) {
441 t.Fatalf("reopened off = %d", w2.off)
442 }
443
444 sink := newWalSink()
445 if ferr := w2.ForEach(sink.add); ferr != nil {
446 t.Fatal(ferr)
447 return
448 }
449 if sink.n != 1 || !bytes.Equal(sink.data[0], a) {
450 t.Fatalf("reopened ForEach count = %d", sink.n)
451 }
452
453 // The next append continues where the old one stopped, not over it.
454 b := []byte("entry-bb")
455 sb, berr := w2.Append(b)
456 if berr != nil {
457 t.Fatal(berr)
458 return
459 }
460 if sb != sa+uint64(HdrSize+len(a)) {
461 t.Fatalf("append after reopen serial = %d", sb)
462 }
463 gb, gerr := w2.Read(sb)
464 if gerr != nil {
465 t.Fatal(gerr)
466 return
467 }
468 if !bytes.Equal(gb, b) {
469 t.Fatal("append after reopen round trip")
470 }
471 if cerr2 := w2.Close(); cerr2 != nil {
472 t.Fatal(cerr2)
473 return
474 }
475 }
476
477 func TestEmptyWAL(t *testing.T) {
478 dir, ok := walTmp(t)
479 if !ok {
480 return
481 }
482 defer os.RemoveAll(dir)
483
484 w, err := Open(dir)
485 if err != nil {
486 t.Fatal(err)
487 return
488 }
489 defer w.Close()
490
491 if len(w.segs) != 1 {
492 t.Fatalf("fresh WAL segment count = %d", int32(len(w.segs)))
493 }
494 if w.off != 0 {
495 t.Fatalf("fresh WAL off = %d", w.off)
496 }
497 sink := newWalSink()
498 if ferr := w.ForEach(sink.add); ferr != nil {
499 t.Fatal(ferr)
500 return
501 }
502 if sink.n != 0 {
503 t.Fatalf("empty WAL yielded %d entries", sink.n)
504 }
505 }
506