// Package sorted provides a sorted flat-file index with binary search. // Records are fixed-width. Inserts batch into a write buffer that is // periodically merge-sorted into the main file on Flush. package sorted import ( "bytes" "os" "sort" "git.smesh.lol/moxie/pkg/mxutil" ) const scanChunk = 256 // records per buffered read // File is a sorted flat-file index. type File struct { path string recLen int32 // total record length in bytes cmpLen int32 // bytes to compare for ordering/search f *os.File count int64 // records on disk buf []byte // write buffer (concatenated records) bufN int32 // records in buffer sorted bool del [][]byte // deleted keys (cmpLen bytes each) bufFile *os.File // .buf sidecar for durable quick flush bufCount int64 // records in .buf file lastFlushedSer uint64 // highest WAL serial covered by last QuickFlush mergeThresh int64 // .buf size that triggers full merge // sideBuf holds the .buf sidecar while it is being searched or merged. It // belongs to the File, which is a self-mutating type, so it is allocated // once and reused: see readBufFile. sideBuf []byte } // Open opens or creates a sorted index file. func Open(path string, recLen, cmpLen int32) (sf *File, derr error) { os.Remove(path | ".tmp") f, err := os.OpenFile(path, os.O_RDWR|os.O_CREATE, 0644) if err != nil { return nil, err } info, err := f.Stat() if err != nil { f.Close() return nil, err } mainSize := info.Size() bufPath := path | ".buf" bf, err := os.OpenFile(bufPath, os.O_RDWR|os.O_CREATE, 0644) if err != nil { f.Close() return nil, err } binfo, err := bf.Stat() if err != nil { f.Close() bf.Close() return nil, err } bufCount := binfo.Size() / int64(recLen) // Seek to end for appending. bf.Seek(0, 2) thresh := mainSize / 10 if thresh < 1<<20 { thresh = 1 << 20 } return &File{ path: path, recLen: recLen, cmpLen: cmpLen, f: f, count: mainSize / int64(recLen), sorted: true, bufFile: bf, bufCount: bufCount, mergeThresh: thresh, }, nil } // refreshSidecar re-reads the .buf sidecar size. The sidecar is the durable // channel between the domain that writes (the ingest worker) and the one that // reads (the relay parent): the in-memory buf belongs to the writer alone, so a // reader's cached bufCount stays at whatever it saw when the file was opened // and Scan would find nothing the writer had flushed. func (s *File) refreshSidecar() { if s.bufFile == nil { return } info, err := s.bufFile.Stat() if err != nil { return } s.bufCount = info.Size() / int64(s.recLen) } // Count returns total records (disk + .buf sidecar + memory buffer). func (s *File) Count() (n int64) { return s.count + s.bufCount + int64(s.bufN) } // Clear removes all records (disk, buffer, and .buf sidecar). func (s *File) Clear() (err error) { s.buf = s.buf[:0] s.bufN = 0 s.del = nil s.sorted = true if err = s.f.Truncate(0); err != nil { return err } if s.bufFile != nil { s.bufFile.Truncate(0) s.bufFile.Seek(0, 0) } s.bufCount = 0 s.count = 0 return nil } // Put adds a record to the write buffer. func (s *File) Put(rec []byte) { // Grow amortized, not once per record. `s.buf = s.buf | rec` allocated a // fresh buffer on every call (concat with a non-empty left side), so each // Put abandoned the previous buffer in the sovereign arena: a 2000-event // seed across 21 indexes left a 67MB data arena and thousands of // compactions. Ensure doubles instead, so only growth steps allocate. s.buf = mxutil.Ensure(s.buf, len(rec)) s.buf = push(s.buf, rec...) s.bufN++ s.sorted = false } // QuickFlush sorts the in-memory buffer and appends it to the .buf sidecar. // Returns true if .buf has exceeded mergeThresh and a full Flush is needed. // sync controls fsync: a flush that only has to make records visible to another // reader writes and returns, because fsync across every index is what makes a // checkpoint expensive. func (s *File) QuickFlush(ser uint64, sync bool) (done bool, derr error) { if s.bufN == 0 { s.lastFlushedSer = ser return false, nil } s.ensureSorted() if _, werr := s.bufFile.Write(s.buf[:s.bufN*s.recLen]); werr != nil { return false, werr } if sync { if serr := s.bufFile.Sync(); serr != nil { return false, serr } } s.bufCount += int64(s.bufN) s.buf = s.buf[:0] s.bufN = 0 s.sorted = true s.lastFlushedSer = ser info, err := s.bufFile.Stat() if err != nil { return false, nil } return info.Size() >= s.mergeThresh, nil } // LastFlushedSer returns the WAL serial recorded by the last QuickFlush. func (s *File) LastFlushedSer() (n uint64) { return s.lastFlushedSer } // Delete marks a key for deletion. The record is excluded from Scan results // and removed on the next Flush. func (s *File) Delete(key []byte) { k := []byte{:s.cmpLen} copy(k, key[:s.cmpLen]) s.del = push(s.del, k) } func (s *File) isDeleted(rec []byte) (ok bool) { for _, d := range s.del { if bytes.Equal(rec[:s.cmpLen], d) { return true } } return false } // readBufFile returns the .buf sidecar's records in a buffer the File owns and // reuses. It used to allocate bufCount*recLen on every call, and Get, GetPrefix // and Scan all call it - a point lookup copied the whole sidecar of every index // it touched, which on a store with thousands of events is megabytes of garbage // per query, all of it landing in the caller's arena. // // The buffer is shared: callers must not retain the slice across another call // on the same File, and Scan sorts it in place (it re-reads on the next call). func (s *File) readBufFile() (buf []byte) { if s.bufCount == 0 { return nil } n := int32(s.bufCount * int64(s.recLen)) s.sideBuf = mxutil.Ensure(s.sideBuf, n) s.sideBuf = s.sideBuf[:n] // A short read leaves the tail zeroed, as it did when the buffer was // allocated fresh: bufCount is the authority on how many records there are. s.bufFile.ReadAt(s.sideBuf, 0) return s.sideBuf } func (s *File) ensureSorted() { if s.sorted || s.bufN <= 1 { s.sorted = true return } tmp := []byte{:s.recLen} rs := &recSorter{s.buf, s.recLen, s.cmpLen, tmp} sort.Sort(rs) s.sorted = true } func (s *File) readAt(dst []byte, idx int64) { s.f.ReadAt(dst[:s.recLen], idx*int64(s.recLen)) } // Get returns the first record matching key (cmpLen bytes). func (s *File) Get(key []byte) (res []byte, okay bool) { // The .buf sidecar is written by another File handle's flush, so its size // has to be re-read here exactly as Scan does; a stale bufCount makes Get // miss records that only live in the sidecar. s.refreshSidecar() rec := []byte{:s.recLen} lo, hi := int64(0), s.count-1 for lo <= hi { mid := lo + (hi-lo)/2 s.readAt(rec, mid) cmp := bytes.Compare(rec[:s.cmpLen], key[:s.cmpLen]) if cmp < 0 { lo = mid + 1 } else if cmp > 0 { hi = mid - 1 } else { // Same key: skip tombstones. Scan honours s.del and this did not, // so a deleted event still came back through a point lookup (an // ids filter, or any GetByID-based path) while a full scan // correctly hid it. for j := mid; j < s.count; j++ { s.readAt(rec, j) if bytes.Compare(rec[:s.cmpLen], key[:s.cmpLen]) != 0 { break } if !s.isDeleted(rec) { return rec, true } } return nil, false } } // Search .buf sidecar. if s.bufCount > 0 { bufData := s.readBufFile() for i := int64(0); i < s.bufCount; i++ { off := i * int64(s.recLen) if bytes.Equal(bufData[off:off+int64(s.cmpLen)], key[:s.cmpLen]) { result := []byte{:s.recLen} copy(result, bufData[off:off+int64(s.recLen)]) if s.isDeleted(result) { continue // tombstone: try the next record with this key } return result, true } } } // Search in-memory buffer. s.ensureSorted() for i := 0; i < s.bufN; i++ { off := i * s.recLen if bytes.Equal(s.buf[off:off+s.cmpLen], key[:s.cmpLen]) { result := []byte{:s.recLen} copy(result, s.buf[off:off+s.recLen]) if s.isDeleted(result) { continue // tombstone: try the next record with this key } return result, true } } return nil, false } // GetPrefix returns the first record whose key begins with prefix. Records are // kept sorted, so this is a lower-bound lookup, not a scan. func (s *File) GetPrefix(prefix []byte) (res []byte, okay bool) { s.refreshSidecar() n := s.lowerBound(prefix) if n < s.count { rec := []byte{:s.recLen} s.readAt(rec, n) if bytes.HasPrefix(rec[:s.cmpLen], prefix) { return rec, true } } if s.bufCount > 0 { bufData := s.readBufFile() for i := int64(0); i < s.bufCount; i++ { off := i * int64(s.recLen) if bytes.HasPrefix(bufData[off:off+int64(s.cmpLen)], prefix) { result := []byte{:s.recLen} copy(result, bufData[off:off+int64(s.recLen)]) return result, true } } } s.ensureSorted() for i := 0; i < s.bufN; i++ { off := i * s.recLen if bytes.HasPrefix(s.buf[off:off+s.cmpLen], prefix) { result := []byte{:s.recLen} copy(result, s.buf[off:off+s.recLen]) return result, true } } return nil, false } // Last returns the record with the largest key. func (s *File) Last() (res []byte, okay bool) { s.ensureSorted() var best []byte if s.count > 0 { best = []byte{:s.recLen} s.readAt(best, s.count-1) } // Check .buf sidecar - scan all records since it's multiple sorted runs. if s.bufCount > 0 { bufData := s.readBufFile() for i := int64(0); i < s.bufCount; i++ { off := i * int64(s.recLen) rec := bufData[off : off+int64(s.recLen)] if best == nil || bytes.Compare(rec[:s.cmpLen], best[:s.cmpLen]) > 0 { if best == nil { best = []byte{:s.recLen} } copy(best, rec) } } } // Check in-memory buffer. if s.bufN > 0 { off := (s.bufN - 1) * s.recLen rec := s.buf[off : off+s.recLen] if best == nil || bytes.Compare(rec[:s.cmpLen], best[:s.cmpLen]) > 0 { if best == nil { best = []byte{:s.recLen} } copy(best, rec) } } if best == nil { return nil, false } return best, true } func (s *File) lowerBound(key []byte) (n int64) { lo, hi := int64(0), s.count rec := []byte{:s.recLen} for lo < hi { mid := lo + (hi-lo)/2 s.readAt(rec, mid) if bytes.Compare(rec[:s.cmpLen], key) < 0 { lo = mid + 1 } else { hi = mid } } return lo } func (s *File) lowerBoundBuf(key []byte) (n int32) { lo, hi := 0, s.bufN for lo < hi { mid := lo + (hi-lo)/2 off := mid * s.recLen if bytes.Compare(s.buf[off:off+s.cmpLen], key) < 0 { lo = mid + 1 } else { hi = mid } } return lo } // Scan iterates records where start <= key <= end. // fn receives each record; return false to stop. // fileChunker reads file records in fixed-size chunks. The cursor is real // state on its own type rather than a closure capture: Moxie closures do not // capture the enclosing scope, so a func literal cannot own chunkStart/chunkEnd. type fileChunker struct { f *File chunk []byte start int64 end int64 } func (c *fileChunker) get(idx int64) (rec []byte) { if idx < c.start || idx >= c.end { c.start = idx n := c.f.count - idx if n > scanChunk { n = scanChunk } c.f.f.ReadAt(c.chunk[:n*int64(c.f.recLen)], idx*int64(c.f.recLen)) c.end = idx + n } off := (idx - c.start) * int64(c.f.recLen) return c.chunk[off : off+int64(c.f.recLen)] } func (s *File) Scan(start, end []byte, fn func(rec []byte) bool) { s.refreshSidecar() s.ensureSorted() fi := s.lowerBound(start) bi := s.lowerBoundBuf(start) // Sort and prepare .buf sidecar data for merge. var sideData []byte var sideN int64 if s.bufCount > 0 { sideData = s.readBufFile() sideN = s.bufCount scratch := []byte{:s.recLen} rs := &recSorter{sideData, s.recLen, s.cmpLen, scratch} sort.Sort(rs) } si := s.lowerBoundSlice(sideData, sideN, start) // Buffered reads from the on-disk file. src := &fileChunker{f: s, chunk: []byte{:scanChunk * s.recLen}, start: -1, end: -1} // Declared outside the loop: a declaration in the body of a self-mutating // method allocates in the sovereign arena on every iteration. var haveFile, haveBuf, haveSide bool var fileRec, bufRec, sideRec []byte var boff int32 var soff int64 var rec []byte for { haveFile, haveBuf, haveSide = false, false, false if fi < s.count { fileRec = src.get(fi) haveFile = bytes.Compare(fileRec[:s.cmpLen], end) <= 0 } if bi < s.bufN { boff = bi * s.recLen bufRec = s.buf[boff : boff+s.recLen] haveBuf = bytes.Compare(bufRec[:s.cmpLen], end) <= 0 } if si < sideN { soff = si * int64(s.recLen) sideRec = sideData[soff : soff+int64(s.recLen)] haveSide = bytes.Compare(sideRec[:s.cmpLen], end) <= 0 } if !haveFile && !haveBuf && !haveSide { break } // Pick the smallest key across all three sources. rec = s.pickSmallest( fileRec, haveFile, bufRec, haveBuf, sideRec, haveSide, &fi, &bi, &si, ) if len(s.del) > 0 && s.isDeleted(rec) { continue } if !fn(rec) { break } } } func (s *File) lowerBoundSlice(data []byte, n int64, key []byte) (nv int64) { lo, hi := int64(0), n for lo < hi { mid := lo + (hi-lo)/2 off := mid * int64(s.recLen) if bytes.Compare(data[off:off+int64(s.cmpLen)], key) < 0 { lo = mid + 1 } else { hi = mid } } return lo } func (s *File) pickSmallest( fileRec []byte, haveFile bool, bufRec []byte, haveBuf bool, sideRec []byte, haveSide bool, fi *int64, bi *int32, si *int64, ) (best []byte) { cl := s.cmpLen var bestKey, rec []byte if haveFile { bestKey = fileRec[:cl] rec = fileRec } if haveBuf { if bestKey == nil || bytes.Compare(bufRec[:cl], bestKey) < 0 { bestKey = bufRec[:cl] rec = bufRec } } if haveSide { if bestKey == nil || bytes.Compare(sideRec[:cl], bestKey) < 0 { bestKey = sideRec[:cl] rec = sideRec } } // Advance all sources that match bestKey (dedup). if haveFile && bytes.Equal(fileRec[:cl], bestKey) { *fi++ } if haveBuf && bytes.Equal(bufRec[:cl], bestKey) { *bi++ } if haveSide && bytes.Equal(sideRec[:cl], bestKey) { *si++ } return rec } // Flush merge-sorts the write buffer and .buf sidecar into the on-disk file. func (s *File) Flush() (err error) { if s.bufN == 0 && s.bufCount == 0 && len(s.del) == 0 { return nil } s.ensureSorted() // Read entire main file into memory. var fileData []byte if s.count > 0 { fileData = []byte{:s.count*int64(s.recLen)} if _, err = s.f.ReadAt(fileData, 0); err != nil { return err } } // Read and sort .buf sidecar. var sideData []byte var sideN int64 if s.bufCount > 0 { sideData = s.readBufFile() sideN = s.bufCount scratch := []byte{:s.recLen} rs := &recSorter{sideData, s.recLen, s.cmpLen, scratch} sort.Sort(rs) } tmpPath := s.path | ".tmp" tmp, err := os.Create(tmpPath) if err != nil { return err } // 3-way merge-sort: main file + .buf sidecar + memory buffer. fi, si, bi := int64(0), int64(0), 0 rl := int64(s.recLen) cl := int64(s.cmpLen) // Declared outside the loop: a declaration in the body of a self-mutating // method allocates in the sovereign arena on every iteration. var fileRec, sideRec, bufRec []byte var haveFile, haveSide, haveBuf bool var off int64 var boff int32 var bestKey, rec []byte for fi < s.count || si < sideN || bi < s.bufN { haveFile = fi < s.count haveSide = si < sideN haveBuf = bi < s.bufN if haveFile { off = fi * rl fileRec = fileData[off : off+rl] } if haveSide { off = si * rl sideRec = sideData[off : off+rl] } if haveBuf { boff = bi * s.recLen bufRec = s.buf[boff : boff+s.recLen] } // Find smallest key, advance all ties (buffer/side win over file). bestKey, rec = s.mergeStep( fileRec, haveFile, sideRec, haveSide, bufRec, haveBuf, ) if haveFile && bytes.Equal(fileRec[:cl], bestKey) { fi++ } if haveSide && bytes.Equal(sideRec[:cl], bestKey) { si++ } if haveBuf && bytes.Equal(bufRec[:s.cmpLen], bestKey) { bi++ } if len(s.del) > 0 && s.isDeleted(rec) { continue } if _, err = tmp.Write(rec); err != nil { tmp.Close() os.Remove(tmpPath) return err } } s.del = s.del[:0] info, err := tmp.Stat() if err != nil { tmp.Close() os.Remove(tmpPath) return err } newCount := info.Size() / rl tmp.Close() s.f.Close() if err = os.Rename(tmpPath, s.path); err != nil { return err } s.f, err = os.OpenFile(s.path, os.O_RDWR, 0644) if err != nil { return err } // Truncate .buf sidecar. s.bufFile.Truncate(0) s.bufFile.Seek(0, 0) s.bufCount = 0 s.count = newCount s.buf = s.buf[:0] s.bufN = 0 s.sorted = true // Recompute merge threshold. mainSize := newCount * rl s.mergeThresh = mainSize / 10 if s.mergeThresh < 1<<20 { s.mergeThresh = 1 << 20 } return nil } // mergeStep picks the record with the smallest key. Buffer/side win ties over file. func (s *File) mergeStep( fileRec []byte, haveFile bool, sideRec []byte, haveSide bool, bufRec []byte, haveBuf bool, ) (key []byte, out []byte) { cl := s.cmpLen var bestKey, rec []byte // Priority: bufRec > sideRec > fileRec on tie. if haveFile { bestKey = fileRec[:cl] rec = fileRec } if haveSide { if bestKey == nil || bytes.Compare(sideRec[:cl], bestKey) <= 0 { bestKey = sideRec[:cl] rec = sideRec } } if haveBuf { if bestKey == nil || bytes.Compare(bufRec[:cl], bestKey) <= 0 { bestKey = bufRec[:cl] rec = bufRec } } return bestKey, rec } // Dirty returns true if there are unflushed records. func (s *File) Dirty() (ok bool) { return s.bufN > 0 || s.bufCount > 0 || len(s.del) > 0 } // SkipFlush marks the buffer as empty so Close won't re-flush. func (s *File) SkipFlush() { s.bufN = 0; s.bufCount = 0; s.del = nil } // Close flushes and closes the file. func (s *File) Close() (err error) { if err = s.Flush(); err != nil { return err } if s.bufFile != nil { s.bufFile.Close() } return s.f.Close() } // recSorter sorts fixed-width records in a flat byte slice. type recSorter struct { data []byte recLen int32 cmpLen int32 tmp []byte } func (r *recSorter) Len() (n int32) { return len(r.data) / r.recLen } func (r *recSorter) Less(i, j int32) (ok bool) { return bytes.Compare( r.data[i*r.recLen:i*r.recLen+r.cmpLen], r.data[j*r.recLen:j*r.recLen+r.cmpLen], ) < 0 } func (r *recSorter) Swap(i, j int32) { a := r.data[i*r.recLen : (i+1)*r.recLen] b := r.data[j*r.recLen : (j+1)*r.recLen] copy(r.tmp, a) copy(a, b) copy(b, r.tmp) }