wal.mx raw

   1  // Package wal provides an append-only segmented value log.
   2  // Event data is written sequentially. The serial encoding
   3  // (segment_id << 32 | offset) provides O(1) reads.
   4  package wal
   5  
   6  import (
   7  	"git.smesh.lol/moxie/pkg/mxutil"
   8  	"encoding/binary"
   9  	"errors"
  10  	"fmt"
  11  	"os"
  12  	"path/filepath"
  13  	"sort"
  14  	"bytes"
  15  
  16  	"git.smesh.lol/morly/pkg/metrics"
  17  )
  18  
  19  const (
  20  	MaxSegSize int64 = 1 << 32 // 4GB per segment
  21  	HdrSize          = 4       // uint32 length prefix
  22  )
  23  
  24  // WAL is an append-only segmented value log.
  25  type WAL struct {
  26  	dir      string
  27  	segs     []*os.File
  28  	cur      int32   // current segment index
  29  	off      int64 // write offset in current segment
  30  	writes   int32   // appends since last sync
  31  	syncEvery int32  // sync after this many appends (0 = manual only)
  32  }
  33  
  34  // Open opens or creates a WAL in dir.
  35  func Open(dir string) (w *WAL, oerr error) {
  36  	if derr := os.MkdirAll(dir, 0755); derr != nil {
  37  		return nil, derr
  38  	}
  39  	w := &WAL{dir: dir, syncEvery: 100}
  40  	entries, err := os.ReadDir(dir)
  41  	if err != nil {
  42  		return nil, err
  43  	}
  44  	var names []string
  45  	for _, e := range entries {
  46  		if bytes.HasPrefix(e.Name(), "seg-") && bytes.HasSuffix(e.Name(), ".dat") {
  47  			names = mxutil.Ensure(names, 1)
  48  			names = push(names, e.Name())
  49  		}
  50  	}
  51  	sort.Strings(names)
  52  	for _, name := range names {
  53  		f, ferr := os.OpenFile(filepath.Join(dir, name), os.O_RDWR, 0644)
  54  		if ferr != nil {
  55  			return nil, ferr
  56  		}
  57  		w.segs = push(w.segs, f)
  58  	}
  59  	if len(w.segs) == 0 {
  60  		if nerr := w.newSeg(); nerr != nil {
  61  			return nil, nerr
  62  		}
  63  	} else {
  64  		w.cur = len(w.segs) - 1
  65  		info, ierr := w.segs[w.cur].Stat()
  66  		if ierr != nil {
  67  			return nil, ierr
  68  		}
  69  		w.off = info.Size()
  70  		// Seek to end so Write() appends rather than overwriting from position 0.
  71  		if _, serr := w.segs[w.cur].Seek(0, 2); serr != nil {
  72  			return nil, serr
  73  		}
  74  	}
  75  	return w, nil
  76  }
  77  
  78  func (w *WAL) segPath(n int32) (s string) {
  79  	return filepath.Join(w.dir, string([]byte(nil) | fmt.Sprintf("seg-%03d.dat", n)))
  80  }
  81  
  82  func (w *WAL) newSeg() (err error) {
  83  	n := len(w.segs)
  84  	f, err := os.Create(w.segPath(n))
  85  	if err != nil {
  86  		return err
  87  	}
  88  	w.segs = push(w.segs, f)
  89  	w.cur = n
  90  	w.off = 0
  91  	return nil
  92  }
  93  
  94  // Append writes data and returns the serial (segment<<32 | offset).
  95  func (w *WAL) Append(data []byte) (ser uint64, derr error) {
  96  	appendStart := metrics.Now()
  97  	defer func() { metrics.WALAppendNs.Observe(metrics.Since(appendStart)) }()
  98  	needed := int64(HdrSize) + int64(len(data))
  99  	if w.off+needed > MaxSegSize {
 100  		if err := w.newSeg(); err != nil {
 101  			return 0, err
 102  		}
 103  	}
 104  	off := w.off
 105  	var hdr [HdrSize]byte
 106  	binary.BigEndian().PutUint32(hdr[:], uint32(len(data)))
 107  	if _, err := w.segs[w.cur].Write(hdr[:]); err != nil {
 108  		return 0, err
 109  	}
 110  	if _, err := w.segs[w.cur].Write(data); err != nil {
 111  		return 0, err
 112  	}
 113  	w.off += needed
 114  	w.writes++
 115  	if w.syncEvery > 0 && w.writes >= w.syncEvery {
 116  		fsyncStart := metrics.Now()
 117  		w.segs[w.cur].Sync()
 118  		metrics.WALFsyncNs.Observe(metrics.Since(fsyncStart))
 119  		w.writes = 0
 120  	}
 121  	return (uint64(w.cur) << 32) | uint64(off), nil
 122  }
 123  
 124  // MaxEntrySize is the largest valid WAL entry (16 MB).
 125  const MaxEntrySize = 16 << 20
 126  
 127  // Read returns the data at the given serial.
 128  func (w *WAL) Read(ser uint64) (out []byte, derr error) {
 129  	seg := int32(ser >> 32)
 130  	off := int64(ser & 0xFFFFFFFF)
 131  	if seg < 0 || seg >= len(w.segs) {
 132  		return nil, fmt.Errorf("wal: segment %d out of range (have %d)", seg, len(w.segs))
 133  	}
 134  	var hdr [HdrSize]byte
 135  	if _, err := w.segs[seg].ReadAt(hdr[:], off); err != nil {
 136  		return nil, fmt.Errorf("wal: read header at seg %d off %d: %w", seg, off, err)
 137  	}
 138  	length := binary.BigEndian().Uint32(hdr[:])
 139  	if length > MaxEntrySize {
 140  		return nil, fmt.Errorf("wal: entry at seg %d off %d has impossible length %d", seg, off, length)
 141  	}
 142  	data := []byte{:length}
 143  	if _, err := w.segs[seg].ReadAt(data, off+int64(HdrSize)); err != nil {
 144  		return nil, fmt.Errorf("wal: read data at seg %d off %d len %d: %w", seg, off, length, err)
 145  	}
 146  	return data, nil
 147  }
 148  
 149  // ForEach iterates all entries in the WAL in order.
 150  // The callback receives the serial and raw data. Return false to stop.
 151  func (w *WAL) ForEach(fn func(ser uint64, data []byte) bool) (err error) {
 152  	for seg := 0; seg < len(w.segs); seg++ {
 153  		info, ierr := w.segs[seg].Stat()
 154  		if ierr != nil {
 155  			return ierr
 156  		}
 157  		var off int64
 158  		for off < info.Size() {
 159  			var hdr [HdrSize]byte
 160  			if _, rerr := w.segs[seg].ReadAt(hdr[:], off); rerr != nil {
 161  				break // truncated entry at end of segment
 162  			}
 163  			length := binary.BigEndian().Uint32(hdr[:])
 164  			if length > MaxEntrySize || off+int64(HdrSize)+int64(length) > info.Size() {
 165  				break // truncated entry
 166  			}
 167  			data := []byte{:length}
 168  			if _, rerr := w.segs[seg].ReadAt(data, off+int64(HdrSize)); rerr != nil {
 169  				break
 170  			}
 171  			ser := (uint64(seg) << 32) | uint64(off)
 172  			if !fn(ser, data) {
 173  				return nil
 174  			}
 175  			off += int64(HdrSize) + int64(length)
 176  		}
 177  	}
 178  	return nil
 179  }
 180  
 181  // ErrCheckpointStale is returned by ForEachFrom when the checkpoint offset
 182  // is beyond the segment's size (segment was truncated/rotated).
 183  var ErrCheckpointStale error
 184  
 185  // ForEachFrom iterates entries after startSer (exclusive).
 186  // If startSer == 0, iterates all entries (equivalent to ForEach).
 187  func (w *WAL) ForEachFrom(startSer uint64, fn func(ser uint64, data []byte) bool) (err error) {
 188  	if startSer == 0 {
 189  		return w.ForEach(fn)
 190  	}
 191  	seg := int32(startSer >> 32)
 192  	off := int64(startSer & 0xFFFFFFFF)
 193  
 194  	if seg >= len(w.segs) {
 195  		return nil
 196  	}
 197  
 198  	info, ierr := w.segs[seg].Stat()
 199  	if ierr != nil {
 200  		return ierr
 201  	}
 202  	if off > info.Size() {
 203  		return ErrCheckpointStale
 204  	}
 205  
 206  	// Skip the entry AT startSer to yield entries AFTER it.
 207  	if off < info.Size() {
 208  		var hdr [HdrSize]byte
 209  		if _, rerr := w.segs[seg].ReadAt(hdr[:], off); rerr != nil {
 210  			return ErrCheckpointStale
 211  		}
 212  		length := binary.BigEndian().Uint32(hdr[:])
 213  		if length > MaxEntrySize {
 214  			return ErrCheckpointStale
 215  		}
 216  		off += int64(HdrSize) + int64(length)
 217  	}
 218  
 219  	// Iterate from current position in this segment, then remaining segments.
 220  	for s := seg; s < len(w.segs); s++ {
 221  		sinfo, serr := w.segs[s].Stat()
 222  		if serr != nil {
 223  			return serr
 224  		}
 225  		pos := int64(0)
 226  		if s == seg {
 227  			pos = off
 228  		}
 229  		for pos < sinfo.Size() {
 230  			var hdr [HdrSize]byte
 231  			if _, rerr := w.segs[s].ReadAt(hdr[:], pos); rerr != nil {
 232  				break
 233  			}
 234  			length := binary.BigEndian().Uint32(hdr[:])
 235  			if length > MaxEntrySize || pos+int64(HdrSize)+int64(length) > sinfo.Size() {
 236  				break
 237  			}
 238  			data := []byte{:length}
 239  			if _, rerr := w.segs[s].ReadAt(data, pos+int64(HdrSize)); rerr != nil {
 240  				break
 241  			}
 242  			ser := (uint64(s) << 32) | uint64(pos)
 243  			if !fn(ser, data) {
 244  				return nil
 245  			}
 246  			pos += int64(HdrSize) + int64(length)
 247  		}
 248  	}
 249  	return nil
 250  }
 251  
 252  // Sync flushes all segment files.
 253  func (w *WAL) Sync() (err error) {
 254  	for _, f := range w.segs {
 255  		if serr := f.Sync(); serr != nil {
 256  			return serr
 257  		}
 258  	}
 259  	return nil
 260  }
 261  
 262  // Close syncs and closes all segments.
 263  func (w *WAL) Close() (err error) {
 264  	var firstErr error
 265  	for _, f := range w.segs {
 266  		f.Sync()
 267  		if cerr := f.Close(); cerr != nil && firstErr == nil {
 268  			firstErr = cerr
 269  		}
 270  	}
 271  	return firstErr
 272  }
 273  
 274  func init() {
 275  	ErrCheckpointStale = errors.New("wal: checkpoint stale")
 276  }
 277