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