diff --git a/db/etl/buffers.go b/db/etl/buffers.go index 844e15770ef..80487604ff5 100644 --- a/db/etl/buffers.go +++ b/db/etl/buffers.go @@ -33,15 +33,15 @@ import ( ) const ( - //SliceBuffer - just simple slice w + // SliceBuffer - just simple slice w SortableSliceBuffer = iota - //SortableAppendBuffer - map[k] [v1 v2 v3] + // SortableAppendBuffer - map[k] [v1 v2 v3] SortableAppendBuffer // SortableOldestAppearedBuffer - buffer that keeps only the oldest entries. // if first v1 was added under key K, then v2; only v1 will stay SortableOldestAppearedBuffer - //BufIOSize - 128 pages | default is 1 page | increasing over `64 * 4096` doesn't show speedup on SSD/NVMe, but show speedup in cloud drives + // BufIOSize - 128 pages | default is 1 page | increasing over `64 * 4096` doesn't show speedup on SSD/NVMe, but show speedup in cloud drives BufIOSize = 128 * 4096 entryLocSize = 16 // sizeof(entryLoc): insertionOrder(4) + offset(4) + keyLen(4) + valLen(4) @@ -79,23 +79,64 @@ func writeSortedEntries(w io.Writer, entries []sortableBufferEntry) error { var BufferOptimalSize = dbg.EnvDataSize("ETL_OPTIMAL", 256*datasize.MB) /* var because we want to sometimes change it from tests or command-line flags */ -// etlSmallBufRAM (BufferOptimalSize/8) bounds the flush threshold so a full -// set of domain/history/index flush collectors (~17 per batch writer) stays -// around 512 MB when all run full. Pooled buffers start empty and grow with -// the data they actually see; grown capacity survives reuse (Reset preserves -// cap), so hot collectors amortize growth while never-full ones stay small. -var etlSmallBufRAM = dbg.EnvDataSize("ETL_SMALL", BufferOptimalSize/8) -var SmallSortableBuffers = NewAllocator(&sync.Pool{ - New: func() any { - return NewSortableBuffer(etlSmallBufRAM) - }, -}) -var etlLargeBufRAM = BufferOptimalSize -var LargeSortableBuffers = NewAllocator(&sync.Pool{ - New: func() any { - return NewSortableBuffer(etlLargeBufRAM) - }, -}) +// etlSmallBufRAM (BufferOptimalSize/8) bounds the flush threshold: +// 3_domains * 2 + 3_history * 1 + 4_indices * 2 = 17 etl collectors, +// 17*(256Mb/8) = 544Mb for all collectors combined. Buffers pool their +// chunks — see dataChunks below. +var ( + etlSmallBufRAM = dbg.EnvDataSize("ETL_SMALL", BufferOptimalSize/8) + SmallSortableBuffers = NewAllocator(&sync.Pool{ + New: func() any { + return NewSortableBuffer(etlSmallBufRAM) + }, + }) +) + +var ( + etlLargeBufRAM = BufferOptimalSize + LargeSortableBuffers = NewAllocator(&sync.Pool{ + New: func() any { + return NewSortableBuffer(etlLargeBufRAM) + }, + }) +) + +const ( + // sortableBuffer stores key/value bytes in chunks of a power-of-two size, so + // entryLoc.offset can pack the chunk index with the offset inside the chunk + // and splitting the two is a shift and a mask. 1MB is also the least a + // collector can hold once it takes a chunk at all. + dataChunkBits = 20 + dataChunkSize = 1 << dataChunkBits // 1MB + + // The chunk index takes what is left of a positive int32, so one buffer + // addresses at most maxDataChunks*dataChunkSize bytes (~2GB); nextChunk + // panics past that. NewSortableBuffer's MaxInt32 bound on optimalSize + // does not fully rule this out, since Put can grow the buffer past + // optimalSize before CheckFlushSize is checked. + maxDataChunks = math.MaxInt32>>dataChunkBits + 1 +) + +// dataChunks are shared by all sortableBuffer instances: a buffer takes chunks as +// it fills and gives them back on Reset, instead of pinning its peak size forever. +var dataChunks = sync.Pool{New: func() any { + c := make([]byte, dataChunkSize) + return &c +}} + +func getDataChunk() []byte { return *dataChunks.Get().(*[]byte) } + +// isPooledChunk reports whether c came from the pool. An oversized entry gets a +// private chunk instead, and handing that one back would let a later +// getDataChunk give an unrelated buffer a chunk of the wrong size. +func isPooledChunk(c []byte) bool { return len(c) == dataChunkSize } + +func putDataChunk(c []byte) { + if !isPooledChunk(c) { + return + } + dataChunks.Put(&c) +} type Buffer interface { // Put does copy `k` and `v` @@ -123,9 +164,10 @@ var ( _ Buffer = &oldestEntrySortableBuffer{} ) -// entryLoc stores the location of a key/value pair within sortableBuffer.data. -// Key occupies data[offset : offset+keyLen], value follows at data[offset+max(0,keyLen) : ...+valLen]. -// keyLen/valLen of -1 indicates nil. +// entryLoc stores the location of a key/value pair inside sortableBuffer. +// offset packs the chunk index and the offset inside that chunk: +// idx<= maxDataChunks { + panic(fmt.Sprintf("etl: sortableBuffer exceeded %d chunks", maxDataChunks)) + } + if n > dataChunkSize { + b.cur = make([]byte, n) + } else { + b.cur = getDataChunk() + } + b.chunks = append(b.chunks, b.cur) + b.curBase = int32(len(b.chunks)-1) << dataChunkBits //nolint:gosec + b.curOff = 0 + b.chunkBytes += len(b.cur) +} + +// entryData returns e's bytes: the key, immediately followed by the value. +func (b *sortableBuffer) entryData(e *entryLoc) []byte { + return b.chunks[e.offset>>dataChunkBits][e.offset&(dataChunkSize-1):] +} + // Put adds key and value to the buffer. These slices will not be accessed later, // so no copying is necessary func (b *sortableBuffer) Put(k, v []byte) { e := entryLoc{ - offset: int32(len(b.data)), //nolint:gosec keyLen: int32(len(k)), //nolint:gosec valLen: int32(len(v)), //nolint:gosec insertionOrder: int32(len(b.entries)), //nolint:gosec @@ -163,11 +234,26 @@ func (b *sortableBuffer) Put(k, v []byte) { if v == nil { e.valLen = -1 } + if n := len(k) + len(v); n > 0 { + off := b.curOff + if int(off)+n > len(b.cur) { + b.nextChunk(n) + off = 0 + } + data := b.cur[off:] + copy(data, k) + copy(data[len(k):], v) + e.offset = b.curBase | off + b.curOff = off + int32(n) //nolint:gosec + } b.entries = append(b.entries, e) - b.data = append(append(b.data, k...), v...) } -func (b *sortableBuffer) Size() int { return len(b.data) + len(b.entries)*entryLocSize } +// Size counts the stored bytes, the tails wasted by the chunks already filled, +// and entryLocSize bytes of metadata per entry. +func (b *sortableBuffer) Size() int { + return b.chunkBytes - (len(b.cur) - int(b.curOff)) + len(b.entries)*entryLocSize +} func (b *sortableBuffer) Len() int { return len(b.entries) @@ -176,49 +262,69 @@ func (b *sortableBuffer) Len() int { func (b *sortableBuffer) Get(i int) ([]byte, []byte) { e := &b.entries[i] kLen, vLen := int(e.keyLen), int(e.valLen) - keyOffset := int(e.offset) - valOffset := keyOffset - if kLen > 0 { - valOffset += kLen - } var key, val []byte - if kLen > 0 { - key = b.data[keyOffset : keyOffset+kLen] - } else if kLen == 0 { + if kLen == 0 { key = []byte{} } - if vLen > 0 { - val = b.data[valOffset : valOffset+vLen] - } else if vLen == 0 { + if vLen == 0 { val = []byte{} } + if kLen <= 0 && vLen <= 0 { + return key, val + } + data := b.entryData(e) + if kLen > 0 { + key = data[:kLen:kLen] + data = data[kLen:] + } + if vLen > 0 { + val = data[:vLen:vLen] + } return key, val } +// Prealloc sizes the entries slice. predictDataSize only reserves room in the +// chunks slice for the chunk pointers; the chunks themselves are still taken +// one at a time, which is what keeps an idle buffer from holding its peak. func (b *sortableBuffer) Prealloc(predictKeysAmount, predictDataSize int) Buffer { if cap(b.entries) < predictKeysAmount { b.entries = make([]entryLoc, 0, predictKeysAmount) } - if cap(b.data) < predictDataSize { - b.data = make([]byte, 0, predictDataSize) + if n := predictDataSize/dataChunkSize + 1; cap(b.chunks) < n { + b.chunks = slices.Grow(b.chunks, n) } return b } func (b *sortableBuffer) Reset() { b.entries = b.entries[:0] - b.data = b.data[:0] + for i, c := range b.chunks { + putDataChunk(c) + b.chunks[i] = nil + } + b.chunks = b.chunks[:0] + b.cur, b.curBase, b.curOff = nil, 0, 0 + b.chunkBytes = 0 } func (b *sortableBuffer) SizeLimit() int { return b.optimalSize } func (b *sortableBuffer) Sort() { - data := b.data - cmp := func(a, b entryLoc) int { - aKey := data[a.offset : a.offset+max(a.keyLen, 0)] - bKey := data[b.offset : b.offset+max(b.keyLen, 0)] - if c := bytes.Compare(aKey, bKey); c != 0 { + chunks := b.chunks + // Key extraction stays inside cmp: pdqsortCmpFunc calls the comparator + // indirectly, so a separate closure never inlines and costs a call per key. + cmp := func(x, y entryLoc) int { + var xk, yk []byte + if x.keyLen > 0 { + off := x.offset & (dataChunkSize - 1) + xk = chunks[x.offset>>dataChunkBits][off : off+x.keyLen] + } + if y.keyLen > 0 { + off := y.offset & (dataChunkSize - 1) + yk = chunks[y.offset>>dataChunkBits][off : off+y.keyLen] + } + if c := bytes.Compare(xk, yk); c != 0 { return c } - return int(a.insertionOrder - b.insertionOrder) // StableSort: preserve insertion order for duplicate keys + return int(x.insertionOrder - y.insertionOrder) // StableSort: preserve insertion order for duplicate keys } if slices.IsSortedFunc(b.entries, cmp) { return @@ -235,10 +341,9 @@ func (b *sortableBuffer) Write(w io.Writer) error { for i := range b.entries { e := &b.entries[i] kLen, vLen := int(e.keyLen), int(e.valLen) - keyOffset := int(e.offset) - valOffset := keyOffset - if kLen > 0 { - valOffset += kLen + var data []byte + if kLen > 0 || vLen > 0 { + data = b.entryData(e) } // write key n := binary.PutVarint(numBuf[:], int64(e.keyLen)) @@ -246,9 +351,10 @@ func (b *sortableBuffer) Write(w io.Writer) error { return err } if kLen > 0 { - if _, err := w.Write(b.data[keyOffset : keyOffset+kLen]); err != nil { + if _, err := w.Write(data[:kLen]); err != nil { return err } + data = data[kLen:] } // write value n = binary.PutVarint(numBuf[:], int64(e.valLen)) @@ -256,7 +362,7 @@ func (b *sortableBuffer) Write(w io.Writer) error { return err } if vLen > 0 { - if _, err := w.Write(b.data[valOffset : valOffset+vLen]); err != nil { + if _, err := w.Write(data[:vLen]); err != nil { return err } } @@ -294,6 +400,7 @@ func (b *appendSortableBuffer) SizeLimit() int { return b.optimalSize } func (b *appendSortableBuffer) Len() int { return len(b.entries) } + func (b *appendSortableBuffer) Sort() { b.sortedBuf = b.sortedBuf[:0] if cap(b.sortedBuf) < len(b.entries) { @@ -316,11 +423,13 @@ func (b *appendSortableBuffer) Swap(i, j int) { func (b *appendSortableBuffer) Get(i int) ([]byte, []byte) { return b.sortedBuf[i].key, b.sortedBuf[i].value } + func (b *appendSortableBuffer) Reset() { b.sortedBuf = nil b.entries = make(map[string][]byte) b.size = 0 } + func (b *appendSortableBuffer) Prealloc(predictKeysAmount, predictDataSize int) Buffer { b.entries = make(map[string][]byte, predictKeysAmount) // maps have no cap(), always recreate if cap(b.sortedBuf) < predictKeysAmount { @@ -392,11 +501,13 @@ func (b *oldestEntrySortableBuffer) Swap(i, j int) { func (b *oldestEntrySortableBuffer) Get(i int) ([]byte, []byte) { return b.sortedBuf[i].key, b.sortedBuf[i].value } + func (b *oldestEntrySortableBuffer) Reset() { b.sortedBuf = nil b.entries = make(map[string][]byte) b.size = 0 } + func (b *oldestEntrySortableBuffer) Prealloc(predictKeysAmount, predictDataSize int) Buffer { b.entries = make(map[string][]byte, predictKeysAmount) // maps have no cap(), always recreate if cap(b.sortedBuf) < predictKeysAmount { @@ -408,6 +519,7 @@ func (b *oldestEntrySortableBuffer) Prealloc(predictKeysAmount, predictDataSize func (b *oldestEntrySortableBuffer) Write(w io.Writer) error { return writeSortedEntries(w, b.sortedBuf) } + func (b *oldestEntrySortableBuffer) CheckFlushSize() bool { return b.size >= b.optimalSize } diff --git a/db/etl/collector.go b/db/etl/collector.go index f1ceab82cd0..6371dce78ae 100644 --- a/db/etl/collector.go +++ b/db/etl/collector.go @@ -45,9 +45,7 @@ func (a *Allocator) Put(b Buffer) { if b == nil { return } - //if cast, ok := b.(*sortableBuffer); ok { - // log.Warn("[dbg] return buf", "cap(cast.data)", cap(cast.data), "cap(cast.lens)", cap(cast.lens)) - //} + b.Reset() // return the buffer's chunks to the pool now — see dataChunks in buffers.go a.p.Put(b) } func (a *Allocator) Get() Buffer { @@ -263,6 +261,14 @@ func (c *Collector) Load(db kv.RwTx, toBucket string, loadFunc LoadFunc, args Tr } func (c *Collector) Close() { + // Providers first: a KeepInRAM one reads straight from `buf`, whose chunks + // Reset hands to a pool that other collectors draw from. + if c.dataProviders != nil { //idempotency + for _, p := range c.dataProviders { + p.Dispose() + } + c.dataProviders = nil + } if c.buf != nil { //idempotency if c.allocator != nil { c.allocator.Put(c.buf) @@ -271,12 +277,6 @@ func (c *Collector) Close() { c.buf.Reset() } } - if c.dataProviders != nil { //idempotency - for _, p := range c.dataProviders { - p.Dispose() - } - c.dataProviders = nil - } c.allFlushed = false } diff --git a/db/etl/dataprovider.go b/db/etl/dataprovider.go index c4a35b2fc14..2dd50763782 100644 --- a/db/etl/dataprovider.go +++ b/db/etl/dataprovider.go @@ -189,12 +189,13 @@ func readField(m *mmapBytesReader) ([]byte, error) { func (p *fileDataProvider) Wait() error { return p.wg.Wait() } func (p *fileDataProvider) Dispose() { + // Wait first: the async flush assigns p.file from its own goroutine, so + // reading it before joining both races and can leak a file created after. + p.Wait() if p.file == nil { return } - p.Wait() - if p.mmapData != nil { _ = p.mmapData.Unmap() p.mmapData = nil diff --git a/db/etl/etl_test.go b/db/etl/etl_test.go index 3887b1a1662..07c8844e282 100644 --- a/db/etl/etl_test.go +++ b/db/etl/etl_test.go @@ -495,10 +495,10 @@ func TestReuseCollectorAfterLoad(t *testing.T) { require.Equal(t, 1, see) c.Close() - // buffers are not lost - require.Empty(t, buf.data) + // buffer state resets for reuse: entries keep their cap, chunks are cleared + require.Empty(t, buf.chunks) require.Empty(t, buf.entries) - require.NotZero(t, cap(buf.data)) + require.Zero(t, buf.Size()) require.NotZero(t, cap(buf.entries)) // teset that no data visible @@ -624,9 +624,8 @@ func TestAppendAcrossProviders(t *testing.T) { } // TestAppendAcrossMemProviders tests that value concatenation works correctly -// when multiple memoryDataProviders have the same key. GetRef returns zero-copy -// slices into sortableBuffer.data — appending to prevV without copying would -// corrupt adjacent entries in the buffer. +// when multiple memoryDataProviders have the same key, across providers backed +// by different buffer types (file-flushed and in-memory). func TestAppendAcrossMemProviders(t *testing.T) { tmpdir := t.TempDir() @@ -917,7 +916,7 @@ func TestMixedProvidersInterleavedKeys(t *testing.T) { } // TestMixedProvidersZeroCopyIntegrity verifies that zero-copy slices from -// memoryDataProvider (GetRef) are not corrupted by subsequent Next() calls. +// memoryDataProvider (Get) are not corrupted by subsequent Next() calls. func TestMixedProvidersZeroCopyIntegrity(t *testing.T) { tmpdir := t.TempDir() @@ -927,7 +926,7 @@ func TestMixedProvidersZeroCopyIntegrity(t *testing.T) { fileProvider, err := FlushToDisk("test", fileBuf, tmpdir, log.LvlInfo) require.NoError(t, err) - // Memory provider with multiple keys - GetRef returns slices into sortableBuffer.data + // Memory provider with multiple keys - Get returns slices into sortableBuffer.chunks memBuf := NewSortableBuffer(BufferOptimalSize) memBuf.Put([]byte("bbb"), []byte("mem-bbb")) memBuf.Put([]byte("ccc"), []byte("mem-ccc")) @@ -1574,3 +1573,156 @@ func TestCollectorWithAllocatorDrawsBufferLazily(t *testing.T) { require.NoError(err) require.Equal([]byte{1}, v) } + +// TestSortableBufferChunks pins the chunked layout: key/value bytes live in +// fixed-size chunks, so a growing buffer never re-allocates and copies the +// bytes it already holds. +func TestSortableBufferChunks(t *testing.T) { + buf := NewSortableBuffer(256 * datasize.MB) + + const entries = 512 + val := bytes.Repeat([]byte{0xAB}, 16*1024) // 512*16KB = 8MB of values + key := make([]byte, 8) + for i := range entries { + binary.BigEndian.PutUint64(key, uint64(i)) + buf.Put(key, val) + } + + require.Equal(t, entries, buf.Len()) + require.Greater(t, len(buf.chunks), 1, "data must be split into chunks") + for i, c := range buf.chunks { + require.Equal(t, dataChunkSize, cap(c), "chunk %d", i) + } + + for i := range entries { + binary.BigEndian.PutUint64(key, uint64(i)) + k, v := buf.Get(i) + require.Equal(t, key, k, "entry %d", i) + require.Equal(t, val, v, "entry %d", i) + } +} + +// TestSortableBufferSortAcrossChunks: the sort comparator has to split a +// packed offset back into a chunk index and an offset inside it, so entries +// must still order correctly once they live past chunk 0. +func TestSortableBufferSortAcrossChunks(t *testing.T) { + buf := NewSortableBuffer(256 * datasize.MB) + + const entries = 512 + val := bytes.Repeat([]byte{0xCD}, 16*1024) // 512*16KB = 8MB of values + key := make([]byte, 8) + for i := range entries { + // Scrambled, so IsSortedFunc cannot short-circuit and pdqsort really runs. + // 313 is odd, so it permutes a power-of-two range. + binary.BigEndian.PutUint64(key, uint64(i*313%entries)) + buf.Put(key, val) + } + require.Greater(t, len(buf.chunks), 1, "data must be split into chunks") + + buf.Sort() + for i := range entries { + binary.BigEndian.PutUint64(key, uint64(i)) + k, v := buf.Get(i) + require.Equal(t, key, k, "entry %d", i) + require.Equal(t, val, v, "entry %d", i) + } +} + +// TestSortableBufferOversizedEntry: an entry bigger than one chunk gets a chunk +// of its own - Get must still return one contiguous slice per key and value. +func TestSortableBufferOversizedEntry(t *testing.T) { + buf := NewSortableBuffer(256 * datasize.MB) + + big := bytes.Repeat([]byte{0xCD}, dataChunkSize+7) + buf.Put([]byte{0x01}, []byte("small")) + buf.Put([]byte{0x02}, big) + buf.Put([]byte{0x03}, []byte("after")) + + k, v := buf.Get(1) + require.Equal(t, []byte{0x02}, k) + require.Equal(t, big, v) + k, v = buf.Get(2) + require.Equal(t, []byte{0x03}, k) + require.Equal(t, []byte("after"), v) + + w := bytes.NewBuffer(nil) + require.NoError(t, buf.Write(w)) + m := &mmapBytesReader{data: w.Bytes()} + for i := range buf.Len() { + wantK, wantV := buf.Get(i) + gotK, err := readField(m) + require.NoError(t, err) + gotV, err := readField(m) + require.NoError(t, err) + require.Equal(t, wantK, gotK) + require.Equal(t, wantV, gotV) + } +} + +// TestSortableBufferResetReleasesChunks: Reset drops the buffer's own chunk +// slice and size bookkeeping so it can be reused immediately. +func TestSortableBufferResetReleasesChunks(t *testing.T) { + buf := NewSortableBuffer(256 * datasize.MB) + val := bytes.Repeat([]byte{0xEF}, 16*1024) + for i := range 512 { + buf.Put(binary.BigEndian.AppendUint64(nil, uint64(i)), val) + } + require.NotEmpty(t, buf.chunks) + + buf.Reset() + require.Empty(t, buf.chunks) + require.Zero(t, buf.Size()) + require.Zero(t, buf.Len()) + + buf.Put([]byte{0x01}, []byte("reused")) + k, v := buf.Get(0) + require.Equal(t, []byte{0x01}, k) + require.Equal(t, []byte("reused"), v) +} + +// TestPutDataChunkRejectsOversized: an entry's private chunk (bigger than +// dataChunkSize) must never enter the shared pool — a later getDataChunk +// handing it out under a normal chunk index would corrupt an unrelated buffer. +func TestPutDataChunkRejectsOversized(t *testing.T) { + for _, tc := range []struct { + name string + length int + pooled bool + }{ + {"short", dataChunkSize - 1, false}, + {"exact", dataChunkSize, true}, + {"oversized", dataChunkSize + 7, false}, + } { + t.Run(tc.name, func(t *testing.T) { + require.Equal(t, tc.pooled, isPooledChunk(make([]byte, tc.length))) + }) + } +} + +// disposeProbe records whether the collector still owned its data chunks when +// the provider was disposed. +type disposeProbe struct { + buf *sortableBuffer + sawOwnChunks bool +} + +func (p *disposeProbe) Next() ([]byte, []byte, error) { return nil, nil, io.EOF } +func (p *disposeProbe) Wait() error { return nil } +func (p *disposeProbe) String() string { return "disposeProbe" } +func (p *disposeProbe) Dispose() { p.sawOwnChunks = len(p.buf.chunks) > 0 } + +// TestCloseDisposesProvidersBeforeBuffer: KeepInRAM hands out a provider backed +// by the collector's own buffer, and Reset gives that buffer's chunks to a pool +// other collectors draw from. So Close must be done with every provider before +// it recycles the buffer. +func TestCloseDisposesProvidersBeforeBuffer(t *testing.T) { + allocator := NewAllocator(&sync.Pool{New: func() any { return NewSortableBuffer(BufferOptimalSize) }}) + c := NewCollectorWithAllocator(t.Name(), t.TempDir(), allocator, log.New()) + require.NoError(t, c.Collect([]byte{1}, []byte{2})) + + probe := &disposeProbe{buf: c.buf.(*sortableBuffer)} + c.dataProviders = append(c.dataProviders, probe) + c.Close() + + require.True(t, probe.sawOwnChunks, "buffer was recycled before its providers were disposed") +}