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