sorted.mx raw
1 // Package sorted provides a sorted flat-file index with binary search.
2 // Records are fixed-width. Inserts batch into a write buffer that is
3 // periodically merge-sorted into the main file on Flush.
4 package sorted
5
6 import (
7 "bytes"
8 "os"
9 "sort"
10 "git.smesh.lol/moxie/pkg/mxutil"
11 )
12
13 const scanChunk = 256 // records per buffered read
14
15 // File is a sorted flat-file index.
16 type File struct {
17 path string
18 recLen int32 // total record length in bytes
19 cmpLen int32 // bytes to compare for ordering/search
20 f *os.File
21 count int64 // records on disk
22 buf []byte // write buffer (concatenated records)
23 bufN int32 // records in buffer
24 sorted bool
25 del [][]byte // deleted keys (cmpLen bytes each)
26
27 bufFile *os.File // .buf sidecar for durable quick flush
28 bufCount int64 // records in .buf file
29 lastFlushedSer uint64 // highest WAL serial covered by last QuickFlush
30 mergeThresh int64 // .buf size that triggers full merge
31 // sideBuf holds the .buf sidecar while it is being searched or merged. It
32 // belongs to the File, which is a self-mutating type, so it is allocated
33 // once and reused: see readBufFile.
34 sideBuf []byte
35 }
36
37 // Open opens or creates a sorted index file.
38 func Open(path string, recLen, cmpLen int32) (sf *File, derr error) {
39 os.Remove(path | ".tmp")
40 f, err := os.OpenFile(path, os.O_RDWR|os.O_CREATE, 0644)
41 if err != nil {
42 return nil, err
43 }
44 info, err := f.Stat()
45 if err != nil {
46 f.Close()
47 return nil, err
48 }
49 mainSize := info.Size()
50
51 bufPath := path | ".buf"
52 bf, err := os.OpenFile(bufPath, os.O_RDWR|os.O_CREATE, 0644)
53 if err != nil {
54 f.Close()
55 return nil, err
56 }
57 binfo, err := bf.Stat()
58 if err != nil {
59 f.Close()
60 bf.Close()
61 return nil, err
62 }
63 bufCount := binfo.Size() / int64(recLen)
64 // Seek to end for appending.
65 bf.Seek(0, 2)
66
67 thresh := mainSize / 10
68 if thresh < 1<<20 {
69 thresh = 1 << 20
70 }
71
72 return &File{
73 path: path,
74 recLen: recLen,
75 cmpLen: cmpLen,
76 f: f,
77 count: mainSize / int64(recLen),
78 sorted: true,
79 bufFile: bf,
80 bufCount: bufCount,
81 mergeThresh: thresh,
82 }, nil
83 }
84
85 // refreshSidecar re-reads the .buf sidecar size. The sidecar is the durable
86 // channel between the domain that writes (the ingest worker) and the one that
87 // reads (the relay parent): the in-memory buf belongs to the writer alone, so a
88 // reader's cached bufCount stays at whatever it saw when the file was opened
89 // and Scan would find nothing the writer had flushed.
90 func (s *File) refreshSidecar() {
91 if s.bufFile == nil {
92 return
93 }
94 info, err := s.bufFile.Stat()
95 if err != nil {
96 return
97 }
98 s.bufCount = info.Size() / int64(s.recLen)
99 }
100
101 // Count returns total records (disk + .buf sidecar + memory buffer).
102 func (s *File) Count() (n int64) { return s.count + s.bufCount + int64(s.bufN) }
103
104 // Clear removes all records (disk, buffer, and .buf sidecar).
105 func (s *File) Clear() (err error) {
106 s.buf = s.buf[:0]
107 s.bufN = 0
108 s.del = nil
109 s.sorted = true
110 if err = s.f.Truncate(0); err != nil {
111 return err
112 }
113 if s.bufFile != nil {
114 s.bufFile.Truncate(0)
115 s.bufFile.Seek(0, 0)
116 }
117 s.bufCount = 0
118 s.count = 0
119 return nil
120 }
121
122 // Put adds a record to the write buffer.
123 func (s *File) Put(rec []byte) {
124 // Grow amortized, not once per record. `s.buf = s.buf | rec` allocated a
125 // fresh buffer on every call (concat with a non-empty left side), so each
126 // Put abandoned the previous buffer in the sovereign arena: a 2000-event
127 // seed across 21 indexes left a 67MB data arena and thousands of
128 // compactions. Ensure doubles instead, so only growth steps allocate.
129 s.buf = mxutil.Ensure(s.buf, len(rec))
130 s.buf = push(s.buf, rec...)
131 s.bufN++
132 s.sorted = false
133 }
134
135 // QuickFlush sorts the in-memory buffer and appends it to the .buf sidecar.
136 // Returns true if .buf has exceeded mergeThresh and a full Flush is needed.
137 // sync controls fsync: a flush that only has to make records visible to another
138 // reader writes and returns, because fsync across every index is what makes a
139 // checkpoint expensive.
140 func (s *File) QuickFlush(ser uint64, sync bool) (done bool, derr error) {
141 if s.bufN == 0 {
142 s.lastFlushedSer = ser
143 return false, nil
144 }
145 s.ensureSorted()
146 if _, werr := s.bufFile.Write(s.buf[:s.bufN*s.recLen]); werr != nil {
147 return false, werr
148 }
149 if sync {
150 if serr := s.bufFile.Sync(); serr != nil {
151 return false, serr
152 }
153 }
154 s.bufCount += int64(s.bufN)
155 s.buf = s.buf[:0]
156 s.bufN = 0
157 s.sorted = true
158 s.lastFlushedSer = ser
159
160 info, err := s.bufFile.Stat()
161 if err != nil {
162 return false, nil
163 }
164 return info.Size() >= s.mergeThresh, nil
165 }
166
167 // LastFlushedSer returns the WAL serial recorded by the last QuickFlush.
168 func (s *File) LastFlushedSer() (n uint64) { return s.lastFlushedSer }
169
170 // Delete marks a key for deletion. The record is excluded from Scan results
171 // and removed on the next Flush.
172 func (s *File) Delete(key []byte) {
173 k := []byte{:s.cmpLen}
174 copy(k, key[:s.cmpLen])
175 s.del = push(s.del, k)
176 }
177
178 func (s *File) isDeleted(rec []byte) (ok bool) {
179 for _, d := range s.del {
180 if bytes.Equal(rec[:s.cmpLen], d) {
181 return true
182 }
183 }
184 return false
185 }
186
187 // readBufFile returns the .buf sidecar's records in a buffer the File owns and
188 // reuses. It used to allocate bufCount*recLen on every call, and Get, GetPrefix
189 // and Scan all call it - a point lookup copied the whole sidecar of every index
190 // it touched, which on a store with thousands of events is megabytes of garbage
191 // per query, all of it landing in the caller's arena.
192 //
193 // The buffer is shared: callers must not retain the slice across another call
194 // on the same File, and Scan sorts it in place (it re-reads on the next call).
195 func (s *File) readBufFile() (buf []byte) {
196 if s.bufCount == 0 {
197 return nil
198 }
199 n := int32(s.bufCount * int64(s.recLen))
200 s.sideBuf = mxutil.Ensure(s.sideBuf, n)
201 s.sideBuf = s.sideBuf[:n]
202 // A short read leaves the tail zeroed, as it did when the buffer was
203 // allocated fresh: bufCount is the authority on how many records there are.
204 s.bufFile.ReadAt(s.sideBuf, 0)
205 return s.sideBuf
206 }
207
208 func (s *File) ensureSorted() {
209 if s.sorted || s.bufN <= 1 {
210 s.sorted = true
211 return
212 }
213 tmp := []byte{:s.recLen}
214 rs := &recSorter{s.buf, s.recLen, s.cmpLen, tmp}
215 sort.Sort(rs)
216 s.sorted = true
217 }
218
219 func (s *File) readAt(dst []byte, idx int64) {
220 s.f.ReadAt(dst[:s.recLen], idx*int64(s.recLen))
221 }
222
223 // Get returns the first record matching key (cmpLen bytes).
224 func (s *File) Get(key []byte) (res []byte, okay bool) {
225 // The .buf sidecar is written by another File handle's flush, so its size
226 // has to be re-read here exactly as Scan does; a stale bufCount makes Get
227 // miss records that only live in the sidecar.
228 s.refreshSidecar()
229 rec := []byte{:s.recLen}
230 lo, hi := int64(0), s.count-1
231 for lo <= hi {
232 mid := lo + (hi-lo)/2
233 s.readAt(rec, mid)
234 cmp := bytes.Compare(rec[:s.cmpLen], key[:s.cmpLen])
235 if cmp < 0 {
236 lo = mid + 1
237 } else if cmp > 0 {
238 hi = mid - 1
239 } else {
240 // Same key: skip tombstones. Scan honours s.del and this did not,
241 // so a deleted event still came back through a point lookup (an
242 // ids filter, or any GetByID-based path) while a full scan
243 // correctly hid it.
244 for j := mid; j < s.count; j++ {
245 s.readAt(rec, j)
246 if bytes.Compare(rec[:s.cmpLen], key[:s.cmpLen]) != 0 {
247 break
248 }
249 if !s.isDeleted(rec) {
250 return rec, true
251 }
252 }
253 return nil, false
254 }
255 }
256 // Search .buf sidecar.
257 if s.bufCount > 0 {
258 bufData := s.readBufFile()
259 for i := int64(0); i < s.bufCount; i++ {
260 off := i * int64(s.recLen)
261 if bytes.Equal(bufData[off:off+int64(s.cmpLen)], key[:s.cmpLen]) {
262 result := []byte{:s.recLen}
263 copy(result, bufData[off:off+int64(s.recLen)])
264 if s.isDeleted(result) {
265 continue // tombstone: try the next record with this key
266 }
267 return result, true
268 }
269 }
270 }
271 // Search in-memory buffer.
272 s.ensureSorted()
273 for i := 0; i < s.bufN; i++ {
274 off := i * s.recLen
275 if bytes.Equal(s.buf[off:off+s.cmpLen], key[:s.cmpLen]) {
276 result := []byte{:s.recLen}
277 copy(result, s.buf[off:off+s.recLen])
278 if s.isDeleted(result) {
279 continue // tombstone: try the next record with this key
280 }
281 return result, true
282 }
283 }
284 return nil, false
285 }
286
287 // GetPrefix returns the first record whose key begins with prefix. Records are
288 // kept sorted, so this is a lower-bound lookup, not a scan.
289 func (s *File) GetPrefix(prefix []byte) (res []byte, okay bool) {
290 s.refreshSidecar()
291 n := s.lowerBound(prefix)
292 if n < s.count {
293 rec := []byte{:s.recLen}
294 s.readAt(rec, n)
295 if bytes.HasPrefix(rec[:s.cmpLen], prefix) {
296 return rec, true
297 }
298 }
299 if s.bufCount > 0 {
300 bufData := s.readBufFile()
301 for i := int64(0); i < s.bufCount; i++ {
302 off := i * int64(s.recLen)
303 if bytes.HasPrefix(bufData[off:off+int64(s.cmpLen)], prefix) {
304 result := []byte{:s.recLen}
305 copy(result, bufData[off:off+int64(s.recLen)])
306 return result, true
307 }
308 }
309 }
310 s.ensureSorted()
311 for i := 0; i < s.bufN; i++ {
312 off := i * s.recLen
313 if bytes.HasPrefix(s.buf[off:off+s.cmpLen], prefix) {
314 result := []byte{:s.recLen}
315 copy(result, s.buf[off:off+s.recLen])
316 return result, true
317 }
318 }
319 return nil, false
320 }
321
322 // Last returns the record with the largest key.
323 func (s *File) Last() (res []byte, okay bool) {
324 s.ensureSorted()
325 var best []byte
326 if s.count > 0 {
327 best = []byte{:s.recLen}
328 s.readAt(best, s.count-1)
329 }
330 // Check .buf sidecar - scan all records since it's multiple sorted runs.
331 if s.bufCount > 0 {
332 bufData := s.readBufFile()
333 for i := int64(0); i < s.bufCount; i++ {
334 off := i * int64(s.recLen)
335 rec := bufData[off : off+int64(s.recLen)]
336 if best == nil || bytes.Compare(rec[:s.cmpLen], best[:s.cmpLen]) > 0 {
337 if best == nil {
338 best = []byte{:s.recLen}
339 }
340 copy(best, rec)
341 }
342 }
343 }
344 // Check in-memory buffer.
345 if s.bufN > 0 {
346 off := (s.bufN - 1) * s.recLen
347 rec := s.buf[off : off+s.recLen]
348 if best == nil || bytes.Compare(rec[:s.cmpLen], best[:s.cmpLen]) > 0 {
349 if best == nil {
350 best = []byte{:s.recLen}
351 }
352 copy(best, rec)
353 }
354 }
355 if best == nil {
356 return nil, false
357 }
358 return best, true
359 }
360
361 func (s *File) lowerBound(key []byte) (n int64) {
362 lo, hi := int64(0), s.count
363 rec := []byte{:s.recLen}
364 for lo < hi {
365 mid := lo + (hi-lo)/2
366 s.readAt(rec, mid)
367 if bytes.Compare(rec[:s.cmpLen], key) < 0 {
368 lo = mid + 1
369 } else {
370 hi = mid
371 }
372 }
373 return lo
374 }
375
376 func (s *File) lowerBoundBuf(key []byte) (n int32) {
377 lo, hi := 0, s.bufN
378 for lo < hi {
379 mid := lo + (hi-lo)/2
380 off := mid * s.recLen
381 if bytes.Compare(s.buf[off:off+s.cmpLen], key) < 0 {
382 lo = mid + 1
383 } else {
384 hi = mid
385 }
386 }
387 return lo
388 }
389
390 // Scan iterates records where start <= key <= end.
391 // fn receives each record; return false to stop.
392 // fileChunker reads file records in fixed-size chunks. The cursor is real
393 // state on its own type rather than a closure capture: Moxie closures do not
394 // capture the enclosing scope, so a func literal cannot own chunkStart/chunkEnd.
395 type fileChunker struct {
396 f *File
397 chunk []byte
398 start int64
399 end int64
400 }
401
402 func (c *fileChunker) get(idx int64) (rec []byte) {
403 if idx < c.start || idx >= c.end {
404 c.start = idx
405 n := c.f.count - idx
406 if n > scanChunk {
407 n = scanChunk
408 }
409 c.f.f.ReadAt(c.chunk[:n*int64(c.f.recLen)], idx*int64(c.f.recLen))
410 c.end = idx + n
411 }
412 off := (idx - c.start) * int64(c.f.recLen)
413 return c.chunk[off : off+int64(c.f.recLen)]
414 }
415
416 func (s *File) Scan(start, end []byte, fn func(rec []byte) bool) {
417 s.refreshSidecar()
418 s.ensureSorted()
419 fi := s.lowerBound(start)
420 bi := s.lowerBoundBuf(start)
421
422 // Sort and prepare .buf sidecar data for merge.
423 var sideData []byte
424 var sideN int64
425 if s.bufCount > 0 {
426 sideData = s.readBufFile()
427 sideN = s.bufCount
428 scratch := []byte{:s.recLen}
429 rs := &recSorter{sideData, s.recLen, s.cmpLen, scratch}
430 sort.Sort(rs)
431 }
432 si := s.lowerBoundSlice(sideData, sideN, start)
433
434 // Buffered reads from the on-disk file.
435 src := &fileChunker{f: s, chunk: []byte{:scanChunk * s.recLen}, start: -1, end: -1}
436
437 // Declared outside the loop: a declaration in the body of a self-mutating
438 // method allocates in the sovereign arena on every iteration.
439 var haveFile, haveBuf, haveSide bool
440 var fileRec, bufRec, sideRec []byte
441 var boff int32
442 var soff int64
443 var rec []byte
444 for {
445 haveFile, haveBuf, haveSide = false, false, false
446
447 if fi < s.count {
448 fileRec = src.get(fi)
449 haveFile = bytes.Compare(fileRec[:s.cmpLen], end) <= 0
450 }
451 if bi < s.bufN {
452 boff = bi * s.recLen
453 bufRec = s.buf[boff : boff+s.recLen]
454 haveBuf = bytes.Compare(bufRec[:s.cmpLen], end) <= 0
455 }
456 if si < sideN {
457 soff = si * int64(s.recLen)
458 sideRec = sideData[soff : soff+int64(s.recLen)]
459 haveSide = bytes.Compare(sideRec[:s.cmpLen], end) <= 0
460 }
461 if !haveFile && !haveBuf && !haveSide {
462 break
463 }
464
465 // Pick the smallest key across all three sources.
466 rec = s.pickSmallest(
467 fileRec, haveFile, bufRec, haveBuf, sideRec, haveSide,
468 &fi, &bi, &si,
469 )
470
471 if len(s.del) > 0 && s.isDeleted(rec) {
472 continue
473 }
474 if !fn(rec) {
475 break
476 }
477 }
478 }
479
480 func (s *File) lowerBoundSlice(data []byte, n int64, key []byte) (nv int64) {
481 lo, hi := int64(0), n
482 for lo < hi {
483 mid := lo + (hi-lo)/2
484 off := mid * int64(s.recLen)
485 if bytes.Compare(data[off:off+int64(s.cmpLen)], key) < 0 {
486 lo = mid + 1
487 } else {
488 hi = mid
489 }
490 }
491 return lo
492 }
493
494 func (s *File) pickSmallest(
495 fileRec []byte, haveFile bool,
496 bufRec []byte, haveBuf bool,
497 sideRec []byte, haveSide bool,
498 fi *int64, bi *int32, si *int64,
499 ) (best []byte) {
500 cl := s.cmpLen
501 var bestKey, rec []byte
502
503 if haveFile {
504 bestKey = fileRec[:cl]
505 rec = fileRec
506 }
507 if haveBuf {
508 if bestKey == nil || bytes.Compare(bufRec[:cl], bestKey) < 0 {
509 bestKey = bufRec[:cl]
510 rec = bufRec
511 }
512 }
513 if haveSide {
514 if bestKey == nil || bytes.Compare(sideRec[:cl], bestKey) < 0 {
515 bestKey = sideRec[:cl]
516 rec = sideRec
517 }
518 }
519
520 // Advance all sources that match bestKey (dedup).
521 if haveFile && bytes.Equal(fileRec[:cl], bestKey) {
522 *fi++
523 }
524 if haveBuf && bytes.Equal(bufRec[:cl], bestKey) {
525 *bi++
526 }
527 if haveSide && bytes.Equal(sideRec[:cl], bestKey) {
528 *si++
529 }
530 return rec
531 }
532
533 // Flush merge-sorts the write buffer and .buf sidecar into the on-disk file.
534 func (s *File) Flush() (err error) {
535 if s.bufN == 0 && s.bufCount == 0 && len(s.del) == 0 {
536 return nil
537 }
538 s.ensureSorted()
539
540 // Read entire main file into memory.
541 var fileData []byte
542 if s.count > 0 {
543 fileData = []byte{:s.count*int64(s.recLen)}
544 if _, err = s.f.ReadAt(fileData, 0); err != nil {
545 return err
546 }
547 }
548
549 // Read and sort .buf sidecar.
550 var sideData []byte
551 var sideN int64
552 if s.bufCount > 0 {
553 sideData = s.readBufFile()
554 sideN = s.bufCount
555 scratch := []byte{:s.recLen}
556 rs := &recSorter{sideData, s.recLen, s.cmpLen, scratch}
557 sort.Sort(rs)
558 }
559
560 tmpPath := s.path | ".tmp"
561 tmp, err := os.Create(tmpPath)
562 if err != nil {
563 return err
564 }
565
566 // 3-way merge-sort: main file + .buf sidecar + memory buffer.
567 fi, si, bi := int64(0), int64(0), 0
568 rl := int64(s.recLen)
569 cl := int64(s.cmpLen)
570 // Declared outside the loop: a declaration in the body of a self-mutating
571 // method allocates in the sovereign arena on every iteration.
572 var fileRec, sideRec, bufRec []byte
573 var haveFile, haveSide, haveBuf bool
574 var off int64
575 var boff int32
576 var bestKey, rec []byte
577 for fi < s.count || si < sideN || bi < s.bufN {
578 haveFile = fi < s.count
579 haveSide = si < sideN
580 haveBuf = bi < s.bufN
581 if haveFile {
582 off = fi * rl
583 fileRec = fileData[off : off+rl]
584 }
585 if haveSide {
586 off = si * rl
587 sideRec = sideData[off : off+rl]
588 }
589 if haveBuf {
590 boff = bi * s.recLen
591 bufRec = s.buf[boff : boff+s.recLen]
592 }
593
594 // Find smallest key, advance all ties (buffer/side win over file).
595 bestKey, rec = s.mergeStep(
596 fileRec, haveFile, sideRec, haveSide, bufRec, haveBuf,
597 )
598 if haveFile && bytes.Equal(fileRec[:cl], bestKey) {
599 fi++
600 }
601 if haveSide && bytes.Equal(sideRec[:cl], bestKey) {
602 si++
603 }
604 if haveBuf && bytes.Equal(bufRec[:s.cmpLen], bestKey) {
605 bi++
606 }
607
608 if len(s.del) > 0 && s.isDeleted(rec) {
609 continue
610 }
611 if _, err = tmp.Write(rec); err != nil {
612 tmp.Close()
613 os.Remove(tmpPath)
614 return err
615 }
616 }
617 s.del = s.del[:0]
618
619 info, err := tmp.Stat()
620 if err != nil {
621 tmp.Close()
622 os.Remove(tmpPath)
623 return err
624 }
625 newCount := info.Size() / rl
626 tmp.Close()
627 s.f.Close()
628
629 if err = os.Rename(tmpPath, s.path); err != nil {
630 return err
631 }
632 s.f, err = os.OpenFile(s.path, os.O_RDWR, 0644)
633 if err != nil {
634 return err
635 }
636
637 // Truncate .buf sidecar.
638 s.bufFile.Truncate(0)
639 s.bufFile.Seek(0, 0)
640 s.bufCount = 0
641
642 s.count = newCount
643 s.buf = s.buf[:0]
644 s.bufN = 0
645 s.sorted = true
646
647 // Recompute merge threshold.
648 mainSize := newCount * rl
649 s.mergeThresh = mainSize / 10
650 if s.mergeThresh < 1<<20 {
651 s.mergeThresh = 1 << 20
652 }
653 return nil
654 }
655
656 // mergeStep picks the record with the smallest key. Buffer/side win ties over file.
657 func (s *File) mergeStep(
658 fileRec []byte, haveFile bool,
659 sideRec []byte, haveSide bool,
660 bufRec []byte, haveBuf bool,
661 ) (key []byte, out []byte) {
662 cl := s.cmpLen
663 var bestKey, rec []byte
664 // Priority: bufRec > sideRec > fileRec on tie.
665 if haveFile {
666 bestKey = fileRec[:cl]
667 rec = fileRec
668 }
669 if haveSide {
670 if bestKey == nil || bytes.Compare(sideRec[:cl], bestKey) <= 0 {
671 bestKey = sideRec[:cl]
672 rec = sideRec
673 }
674 }
675 if haveBuf {
676 if bestKey == nil || bytes.Compare(bufRec[:cl], bestKey) <= 0 {
677 bestKey = bufRec[:cl]
678 rec = bufRec
679 }
680 }
681 return bestKey, rec
682 }
683
684 // Dirty returns true if there are unflushed records.
685 func (s *File) Dirty() (ok bool) { return s.bufN > 0 || s.bufCount > 0 || len(s.del) > 0 }
686
687 // SkipFlush marks the buffer as empty so Close won't re-flush.
688 func (s *File) SkipFlush() { s.bufN = 0; s.bufCount = 0; s.del = nil }
689
690 // Close flushes and closes the file.
691 func (s *File) Close() (err error) {
692 if err = s.Flush(); err != nil {
693 return err
694 }
695 if s.bufFile != nil {
696 s.bufFile.Close()
697 }
698 return s.f.Close()
699 }
700
701 // recSorter sorts fixed-width records in a flat byte slice.
702 type recSorter struct {
703 data []byte
704 recLen int32
705 cmpLen int32
706 tmp []byte
707 }
708
709 func (r *recSorter) Len() (n int32) { return len(r.data) / r.recLen }
710 func (r *recSorter) Less(i, j int32) (ok bool) {
711 return bytes.Compare(
712 r.data[i*r.recLen:i*r.recLen+r.cmpLen],
713 r.data[j*r.recLen:j*r.recLen+r.cmpLen],
714 ) < 0
715 }
716 func (r *recSorter) Swap(i, j int32) {
717 a := r.data[i*r.recLen : (i+1)*r.recLen]
718 b := r.data[j*r.recLen : (j+1)*r.recLen]
719 copy(r.tmp, a)
720 copy(a, b)
721 copy(b, r.tmp)
722 }
723