mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-27 15:16:43 +00:00
common, core/filtermaps: use Range type
This commit is contained in:
parent
9031f90af5
commit
6a6deb1bd8
8 changed files with 299 additions and 253 deletions
|
|
@ -23,6 +23,7 @@ import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"math/big"
|
"math/big"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
"reflect"
|
"reflect"
|
||||||
|
|
@ -30,6 +31,7 @@ import (
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common/hexutil"
|
"github.com/ethereum/go-ethereum/common/hexutil"
|
||||||
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
"golang.org/x/crypto/sha3"
|
"golang.org/x/crypto/sha3"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -486,3 +488,83 @@ func (b PrettyBytes) TerminalString() string {
|
||||||
}
|
}
|
||||||
return fmt.Sprintf("%#x...%x (%dB)", b[:3], b[len(b)-3:], len(b))
|
return fmt.Sprintf("%#x...%x (%dB)", b[:3], b[len(b)-3:], len(b))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type Range[T uint32 | uint64] struct {
|
||||||
|
first, afterLast T
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewRange[T uint32 | uint64](first, count T) Range[T] {
|
||||||
|
return Range[T]{first, first + count}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *Range[T]) EncodeRLP(w io.Writer) error {
|
||||||
|
if err := rlp.Encode(w, &r.first); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return rlp.Encode(w, &r.afterLast)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *Range[T]) DecodeRLP(s *rlp.Stream) error {
|
||||||
|
if err := s.Decode(&r.first); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return s.Decode(&r.afterLast)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r Range[T]) First() T {
|
||||||
|
return r.first
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r Range[T]) Last() T {
|
||||||
|
if r.first == r.afterLast {
|
||||||
|
panic("last item of zero length range is not allowed")
|
||||||
|
}
|
||||||
|
return r.afterLast - 1
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r Range[T]) AfterLast() T {
|
||||||
|
return r.afterLast
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r Range[T]) Count() T {
|
||||||
|
return r.afterLast - r.first
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r Range[T]) IsEmpty() bool {
|
||||||
|
return r.first == r.afterLast
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r Range[T]) Includes(v T) bool {
|
||||||
|
return v >= r.first && v < r.afterLast
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *Range[T]) SetFirst(v T) {
|
||||||
|
r.first = v
|
||||||
|
if r.afterLast < r.first {
|
||||||
|
r.afterLast = r.first
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *Range[T]) SetAfterLast(v T) {
|
||||||
|
r.afterLast = v
|
||||||
|
if r.afterLast < r.first {
|
||||||
|
r.first = r.afterLast
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *Range[T]) SetLast(v T) {
|
||||||
|
r.SetAfterLast(v + 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r Range[T]) Intersection(q Range[T]) Range[T] {
|
||||||
|
if r.first > q.first {
|
||||||
|
q.first = r.first
|
||||||
|
}
|
||||||
|
if r.afterLast < q.afterLast {
|
||||||
|
q.afterLast = r.afterLast
|
||||||
|
}
|
||||||
|
if q.first > q.afterLast {
|
||||||
|
return Range[T]{}
|
||||||
|
}
|
||||||
|
return q
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -145,24 +145,23 @@ func (a FilterRow) Equal(b FilterRow) bool {
|
||||||
// filterMapsRange describes the rendered range of filter maps and the range
|
// filterMapsRange describes the rendered range of filter maps and the range
|
||||||
// of fully rendered blocks.
|
// of fully rendered blocks.
|
||||||
type filterMapsRange struct {
|
type filterMapsRange struct {
|
||||||
initialized bool
|
initialized bool
|
||||||
headBlockIndexed bool
|
headIndexed bool
|
||||||
headBlockDelimiter uint64 // zero if afterLastIndexedBlock != targetBlockNumber
|
headDelimiter uint64 // zero if headIndexed is false
|
||||||
// if initialized then all maps are rendered between firstRenderedMap and
|
// if initialized then all maps are rendered in the maps range
|
||||||
// afterLastRenderedMap-1
|
maps common.Range[uint32]
|
||||||
firstRenderedMap, afterLastRenderedMap uint32
|
|
||||||
// if tailPartialEpoch > 0 then maps between firstRenderedMap-mapsPerEpoch and
|
// if tailPartialEpoch > 0 then maps between firstRenderedMap-mapsPerEpoch and
|
||||||
// firstRenderedMap-mapsPerEpoch+tailPartialEpoch-1 are rendered
|
// firstRenderedMap-mapsPerEpoch+tailPartialEpoch-1 are rendered
|
||||||
tailPartialEpoch uint32
|
tailPartialEpoch uint32
|
||||||
// if initialized then all log values belonging to blocks between
|
// if initialized then all log values in the blocks range are fully
|
||||||
// firstIndexedBlock and afterLastIndexedBlock are fully rendered
|
// rendered
|
||||||
// blockLvPointers are available between firstIndexedBlock and afterLastIndexedBlock-1
|
// blockLvPointers are available in the blocks range
|
||||||
firstIndexedBlock, afterLastIndexedBlock uint64
|
blocks common.Range[uint64]
|
||||||
}
|
}
|
||||||
|
|
||||||
// hasIndexedBlocks returns true if the range has at least one fully indexed block.
|
// hasIndexedBlocks returns true if the range has at least one fully indexed block.
|
||||||
func (fmr *filterMapsRange) hasIndexedBlocks() bool {
|
func (fmr *filterMapsRange) hasIndexedBlocks() bool {
|
||||||
return fmr.initialized && fmr.afterLastIndexedBlock > fmr.firstIndexedBlock
|
return fmr.initialized && !fmr.blocks.IsEmpty() && !fmr.maps.IsEmpty()
|
||||||
}
|
}
|
||||||
|
|
||||||
// lastBlockOfMap is used for caching the (number, id) pairs belonging to the
|
// lastBlockOfMap is used for caching the (number, id) pairs belonging to the
|
||||||
|
|
@ -200,14 +199,12 @@ func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, historyCutoff, f
|
||||||
exportFileName: config.ExportFileName,
|
exportFileName: config.ExportFileName,
|
||||||
Params: params,
|
Params: params,
|
||||||
indexedRange: filterMapsRange{
|
indexedRange: filterMapsRange{
|
||||||
initialized: initialized,
|
initialized: initialized,
|
||||||
headBlockIndexed: rs.HeadBlockIndexed,
|
headIndexed: rs.HeadIndexed,
|
||||||
headBlockDelimiter: rs.HeadBlockDelimiter,
|
headDelimiter: rs.HeadDelimiter,
|
||||||
firstIndexedBlock: rs.FirstIndexedBlock,
|
blocks: rs.Blocks,
|
||||||
afterLastIndexedBlock: rs.AfterLastIndexedBlock,
|
maps: rs.Maps,
|
||||||
firstRenderedMap: rs.FirstRenderedMap,
|
tailPartialEpoch: rs.TailPartialEpoch,
|
||||||
afterLastRenderedMap: rs.AfterLastRenderedMap,
|
|
||||||
tailPartialEpoch: rs.TailPartialEpoch,
|
|
||||||
},
|
},
|
||||||
matcherSyncCh: make(chan *FilterMapsMatcherBackend),
|
matcherSyncCh: make(chan *FilterMapsMatcherBackend),
|
||||||
matchers: make(map[*FilterMapsMatcherBackend]struct{}),
|
matchers: make(map[*FilterMapsMatcherBackend]struct{}),
|
||||||
|
|
@ -222,24 +219,23 @@ func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, historyCutoff, f
|
||||||
f.targetView = initView
|
f.targetView = initView
|
||||||
if f.indexedRange.initialized {
|
if f.indexedRange.initialized {
|
||||||
f.indexedView = f.initChainView(f.targetView)
|
f.indexedView = f.initChainView(f.targetView)
|
||||||
f.indexedRange.headBlockIndexed = f.indexedRange.afterLastIndexedBlock == f.indexedView.headNumber+1
|
f.indexedRange.headIndexed = f.indexedRange.blocks.AfterLast() == f.indexedView.headNumber+1
|
||||||
if !f.indexedRange.headBlockIndexed {
|
if !f.indexedRange.headIndexed {
|
||||||
f.indexedRange.headBlockDelimiter = 0
|
f.indexedRange.headDelimiter = 0
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if f.indexedRange.hasIndexedBlocks() {
|
if f.indexedRange.hasIndexedBlocks() {
|
||||||
log.Info("Initialized log indexer",
|
log.Info("Initialized log indexer",
|
||||||
"first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
|
"first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(),
|
||||||
"first map", f.indexedRange.firstRenderedMap, "last map", f.indexedRange.afterLastRenderedMap-1,
|
"first map", f.indexedRange.maps.First(), "last map", f.indexedRange.maps.Last(),
|
||||||
"head indexed", f.indexedRange.headBlockIndexed)
|
"head indexed", f.indexedRange.headIndexed)
|
||||||
}
|
}
|
||||||
return f
|
return f
|
||||||
}
|
}
|
||||||
|
|
||||||
// Start starts the indexer.
|
// Start starts the indexer.
|
||||||
func (f *FilterMaps) Start() {
|
func (f *FilterMaps) Start() {
|
||||||
if !f.testDisableSnapshots && f.indexedRange.initialized && f.indexedRange.headBlockIndexed &&
|
if !f.testDisableSnapshots && f.indexedRange.hasIndexedBlocks() && f.indexedRange.headIndexed {
|
||||||
f.indexedRange.firstRenderedMap < f.indexedRange.afterLastRenderedMap {
|
|
||||||
// previous target head rendered; load last map as snapshot
|
// previous target head rendered; load last map as snapshot
|
||||||
if err := f.loadHeadSnapshot(); err != nil {
|
if err := f.loadHeadSnapshot(); err != nil {
|
||||||
log.Error("Could not load head filter map snapshot", "error", err)
|
log.Error("Could not load head filter map snapshot", "error", err)
|
||||||
|
|
@ -261,7 +257,7 @@ func (f *FilterMaps) Stop() {
|
||||||
// Note that the returned view might be shorter than the existing index if
|
// Note that the returned view might be shorter than the existing index if
|
||||||
// the latest maps are not consistent with targetView.
|
// the latest maps are not consistent with targetView.
|
||||||
func (f *FilterMaps) initChainView(chainView *ChainView) *ChainView {
|
func (f *FilterMaps) initChainView(chainView *ChainView) *ChainView {
|
||||||
mapIndex := f.indexedRange.afterLastRenderedMap
|
mapIndex := f.indexedRange.maps.AfterLast()
|
||||||
for {
|
for {
|
||||||
var ok bool
|
var ok bool
|
||||||
mapIndex, ok = f.lastMapBoundaryBefore(mapIndex)
|
mapIndex, ok = f.lastMapBoundaryBefore(mapIndex)
|
||||||
|
|
@ -331,10 +327,8 @@ func (f *FilterMaps) init() error {
|
||||||
}
|
}
|
||||||
if bestLen > 0 {
|
if bestLen > 0 {
|
||||||
cp := checkpoints[bestIdx][bestLen-1]
|
cp := checkpoints[bestIdx][bestLen-1]
|
||||||
fmr.firstIndexedBlock = cp.BlockNumber + 1
|
fmr.blocks = common.NewRange[uint64](cp.BlockNumber+1, 0)
|
||||||
fmr.afterLastIndexedBlock = cp.BlockNumber + 1
|
fmr.maps = common.NewRange[uint32](uint32(bestLen)<<f.logMapsPerEpoch, 0)
|
||||||
fmr.firstRenderedMap = uint32(bestLen) << f.logMapsPerEpoch
|
|
||||||
fmr.afterLastRenderedMap = uint32(bestLen) << f.logMapsPerEpoch
|
|
||||||
}
|
}
|
||||||
f.setRange(batch, f.targetView, fmr)
|
f.setRange(batch, f.targetView, fmr)
|
||||||
return batch.Write()
|
return batch.Write()
|
||||||
|
|
@ -385,13 +379,11 @@ func (f *FilterMaps) setRange(batch ethdb.KeyValueWriter, newView *ChainView, ne
|
||||||
f.updateMatchersValidRange()
|
f.updateMatchersValidRange()
|
||||||
if newRange.initialized {
|
if newRange.initialized {
|
||||||
rs := rawdb.FilterMapsRange{
|
rs := rawdb.FilterMapsRange{
|
||||||
HeadBlockIndexed: newRange.headBlockIndexed,
|
HeadIndexed: newRange.headIndexed,
|
||||||
HeadBlockDelimiter: newRange.headBlockDelimiter,
|
HeadDelimiter: newRange.headDelimiter,
|
||||||
FirstIndexedBlock: newRange.firstIndexedBlock,
|
Blocks: newRange.blocks,
|
||||||
AfterLastIndexedBlock: newRange.afterLastIndexedBlock,
|
Maps: newRange.maps,
|
||||||
FirstRenderedMap: newRange.firstRenderedMap,
|
TailPartialEpoch: newRange.tailPartialEpoch,
|
||||||
AfterLastRenderedMap: newRange.afterLastRenderedMap,
|
|
||||||
TailPartialEpoch: newRange.tailPartialEpoch,
|
|
||||||
}
|
}
|
||||||
rawdb.WriteFilterMapsRange(batch, rs)
|
rawdb.WriteFilterMapsRange(batch, rs)
|
||||||
} else {
|
} else {
|
||||||
|
|
@ -409,7 +401,7 @@ func (f *FilterMaps) setRange(batch ethdb.KeyValueWriter, newView *ChainView, ne
|
||||||
// called from outside the indexerLoop goroutine.
|
// called from outside the indexerLoop goroutine.
|
||||||
func (f *FilterMaps) getLogByLvIndex(lvIndex uint64) (*types.Log, error) {
|
func (f *FilterMaps) getLogByLvIndex(lvIndex uint64) (*types.Log, error) {
|
||||||
mapIndex := uint32(lvIndex >> f.logValuesPerMap)
|
mapIndex := uint32(lvIndex >> f.logValuesPerMap)
|
||||||
if mapIndex < f.indexedRange.firstRenderedMap || mapIndex >= f.indexedRange.afterLastRenderedMap {
|
if !f.indexedRange.maps.Includes(mapIndex) {
|
||||||
return nil, nil
|
return nil, nil
|
||||||
}
|
}
|
||||||
// find possible block range based on map to block pointers
|
// find possible block range based on map to block pointers
|
||||||
|
|
@ -424,8 +416,8 @@ func (f *FilterMaps) getLogByLvIndex(lvIndex uint64) (*types.Log, error) {
|
||||||
return nil, fmt.Errorf("failed to retrieve last block of map %d before searched log value index %d: %v", mapIndex, lvIndex, err)
|
return nil, fmt.Errorf("failed to retrieve last block of map %d before searched log value index %d: %v", mapIndex, lvIndex, err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if firstBlockNumber < f.indexedRange.firstIndexedBlock {
|
if firstBlockNumber < f.indexedRange.blocks.First() {
|
||||||
firstBlockNumber = f.indexedRange.firstIndexedBlock
|
firstBlockNumber = f.indexedRange.blocks.First()
|
||||||
}
|
}
|
||||||
// find block with binary search based on block to log value index pointers
|
// find block with binary search based on block to log value index pointers
|
||||||
for firstBlockNumber < lastBlockNumber {
|
for firstBlockNumber < lastBlockNumber {
|
||||||
|
|
@ -584,8 +576,8 @@ func (f *FilterMaps) mapRowIndex(mapIndex, rowIndex uint32) uint64 {
|
||||||
// Note that this function assumes that the indexer read lock is being held when
|
// Note that this function assumes that the indexer read lock is being held when
|
||||||
// called from outside the indexerLoop goroutine.
|
// called from outside the indexerLoop goroutine.
|
||||||
func (f *FilterMaps) getBlockLvPointer(blockNumber uint64) (uint64, error) {
|
func (f *FilterMaps) getBlockLvPointer(blockNumber uint64) (uint64, error) {
|
||||||
if blockNumber >= f.indexedRange.afterLastIndexedBlock && f.indexedRange.headBlockIndexed {
|
if blockNumber >= f.indexedRange.blocks.AfterLast() && f.indexedRange.headIndexed {
|
||||||
return f.indexedRange.headBlockDelimiter, nil
|
return f.indexedRange.headDelimiter, nil
|
||||||
}
|
}
|
||||||
if lvPointer, ok := f.lvPointerCache.Get(blockNumber); ok {
|
if lvPointer, ok := f.lvPointerCache.Get(blockNumber); ok {
|
||||||
return lvPointer, nil
|
return lvPointer, nil
|
||||||
|
|
@ -662,26 +654,28 @@ func (f *FilterMaps) deleteTailEpoch(epoch uint32) error {
|
||||||
firstBlock++
|
firstBlock++
|
||||||
}
|
}
|
||||||
fmr := f.indexedRange
|
fmr := f.indexedRange
|
||||||
if f.indexedRange.firstRenderedMap == firstMap &&
|
if f.indexedRange.maps.First() == firstMap &&
|
||||||
f.indexedRange.afterLastRenderedMap > firstMap+f.mapsPerEpoch &&
|
f.indexedRange.maps.AfterLast() > firstMap+f.mapsPerEpoch &&
|
||||||
f.indexedRange.tailPartialEpoch == 0 {
|
f.indexedRange.tailPartialEpoch == 0 {
|
||||||
fmr.firstRenderedMap = firstMap + f.mapsPerEpoch
|
fmr.maps.SetFirst(firstMap + f.mapsPerEpoch)
|
||||||
fmr.firstIndexedBlock = lastBlock + 1
|
fmr.blocks.SetFirst(lastBlock + 1)
|
||||||
} else if f.indexedRange.firstRenderedMap == firstMap+f.mapsPerEpoch {
|
} else if f.indexedRange.maps.First() == firstMap+f.mapsPerEpoch {
|
||||||
fmr.tailPartialEpoch = 0
|
fmr.tailPartialEpoch = 0
|
||||||
} else {
|
} else {
|
||||||
return errors.New("invalid tail epoch number")
|
return errors.New("invalid tail epoch number")
|
||||||
}
|
}
|
||||||
f.setRange(f.db, f.indexedView, fmr)
|
f.setRange(f.db, f.indexedView, fmr)
|
||||||
rawdb.DeleteFilterMapRows(f.db, f.mapRowIndex(firstMap, 0), f.mapRowIndex(firstMap+f.mapsPerEpoch, 0))
|
first := f.mapRowIndex(firstMap, 0)
|
||||||
|
count := f.mapRowIndex(firstMap+f.mapsPerEpoch, 0) - first
|
||||||
|
rawdb.DeleteFilterMapRows(f.db, common.NewRange[uint64](first, count))
|
||||||
for mapIndex := firstMap; mapIndex < firstMap+f.mapsPerEpoch; mapIndex++ {
|
for mapIndex := firstMap; mapIndex < firstMap+f.mapsPerEpoch; mapIndex++ {
|
||||||
f.filterMapCache.Remove(mapIndex)
|
f.filterMapCache.Remove(mapIndex)
|
||||||
}
|
}
|
||||||
rawdb.DeleteFilterMapLastBlocks(f.db, firstMap, firstMap+f.mapsPerEpoch-1) // keep last enrty
|
rawdb.DeleteFilterMapLastBlocks(f.db, common.NewRange[uint32](firstMap, f.mapsPerEpoch-1)) // keep last enrty
|
||||||
for mapIndex := firstMap; mapIndex < firstMap+f.mapsPerEpoch-1; mapIndex++ {
|
for mapIndex := firstMap; mapIndex < firstMap+f.mapsPerEpoch-1; mapIndex++ {
|
||||||
f.lastBlockCache.Remove(mapIndex)
|
f.lastBlockCache.Remove(mapIndex)
|
||||||
}
|
}
|
||||||
rawdb.DeleteBlockLvPointers(f.db, firstBlock, lastBlock) // keep last enrty
|
rawdb.DeleteBlockLvPointers(f.db, common.NewRange[uint64](firstBlock, lastBlock-firstBlock)) // keep last enrty
|
||||||
for blockNumber := firstBlock; blockNumber < lastBlock; blockNumber++ {
|
for blockNumber := firstBlock; blockNumber < lastBlock; blockNumber++ {
|
||||||
f.lvPointerCache.Remove(blockNumber)
|
f.lvPointerCache.Remove(blockNumber)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -193,20 +193,20 @@ func (f *FilterMaps) tryIndexHead() bool {
|
||||||
f.lastLogHeadIndex = time.Now()
|
f.lastLogHeadIndex = time.Now()
|
||||||
f.startedHeadIndexAt = f.lastLogHeadIndex
|
f.startedHeadIndexAt = f.lastLogHeadIndex
|
||||||
f.startedHeadIndex = true
|
f.startedHeadIndex = true
|
||||||
f.ptrHeadIndex = f.indexedRange.afterLastIndexedBlock
|
f.ptrHeadIndex = f.indexedRange.blocks.AfterLast()
|
||||||
}
|
}
|
||||||
if _, err := headRenderer.run(func() bool {
|
if _, err := headRenderer.run(func() bool {
|
||||||
f.processEvents()
|
f.processEvents()
|
||||||
return f.stop
|
return f.stop
|
||||||
}, func() {
|
}, func() {
|
||||||
f.tryUnindexTail()
|
f.tryUnindexTail()
|
||||||
if f.indexedRange.hasIndexedBlocks() && f.indexedRange.afterLastIndexedBlock >= f.ptrHeadIndex &&
|
if f.indexedRange.hasIndexedBlocks() && f.indexedRange.blocks.AfterLast() >= f.ptrHeadIndex &&
|
||||||
((!f.loggedHeadIndex && time.Since(f.startedHeadIndexAt) > headLogDelay) ||
|
((!f.loggedHeadIndex && time.Since(f.startedHeadIndexAt) > headLogDelay) ||
|
||||||
time.Since(f.lastLogHeadIndex) > logFrequency) {
|
time.Since(f.lastLogHeadIndex) > logFrequency) {
|
||||||
log.Info("Log index head rendering in progress",
|
log.Info("Log index head rendering in progress",
|
||||||
"first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
|
"first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(),
|
||||||
"processed", f.indexedRange.afterLastIndexedBlock-f.ptrHeadIndex,
|
"processed", f.indexedRange.blocks.AfterLast()-f.ptrHeadIndex,
|
||||||
"remaining", f.indexedView.headNumber+1-f.indexedRange.afterLastIndexedBlock,
|
"remaining", f.indexedView.headNumber-f.indexedRange.blocks.Last(),
|
||||||
"elapsed", common.PrettyDuration(time.Since(f.startedHeadIndexAt)))
|
"elapsed", common.PrettyDuration(time.Since(f.startedHeadIndexAt)))
|
||||||
f.loggedHeadIndex = true
|
f.loggedHeadIndex = true
|
||||||
f.lastLogHeadIndex = time.Now()
|
f.lastLogHeadIndex = time.Now()
|
||||||
|
|
@ -215,10 +215,10 @@ func (f *FilterMaps) tryIndexHead() bool {
|
||||||
log.Error("Log index head rendering failed", "error", err)
|
log.Error("Log index head rendering failed", "error", err)
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
if f.loggedHeadIndex {
|
if f.loggedHeadIndex && f.indexedRange.hasIndexedBlocks() {
|
||||||
log.Info("Log index head rendering finished",
|
log.Info("Log index head rendering finished",
|
||||||
"first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
|
"first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(),
|
||||||
"processed", f.indexedRange.afterLastIndexedBlock-f.ptrHeadIndex,
|
"processed", f.indexedRange.blocks.AfterLast()-f.ptrHeadIndex,
|
||||||
"elapsed", common.PrettyDuration(time.Since(f.startedHeadIndexAt)))
|
"elapsed", common.PrettyDuration(time.Since(f.startedHeadIndexAt)))
|
||||||
}
|
}
|
||||||
f.loggedHeadIndex, f.startedHeadIndex = false, false
|
f.loggedHeadIndex, f.startedHeadIndex = false, false
|
||||||
|
|
@ -232,7 +232,7 @@ func (f *FilterMaps) tryIndexHead() bool {
|
||||||
// is changed.
|
// is changed.
|
||||||
func (f *FilterMaps) tryIndexTail() bool {
|
func (f *FilterMaps) tryIndexTail() bool {
|
||||||
for {
|
for {
|
||||||
firstEpoch := f.indexedRange.firstRenderedMap >> f.logMapsPerEpoch
|
firstEpoch := f.indexedRange.maps.First() >> f.logMapsPerEpoch
|
||||||
if firstEpoch == 0 || !f.needTailEpoch(firstEpoch-1) {
|
if firstEpoch == 0 || !f.needTailEpoch(firstEpoch-1) {
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|
@ -243,12 +243,12 @@ func (f *FilterMaps) tryIndexTail() bool {
|
||||||
// resume process if tail rendering was interrupted because of head rendering
|
// resume process if tail rendering was interrupted because of head rendering
|
||||||
tailRenderer := f.tailRenderer
|
tailRenderer := f.tailRenderer
|
||||||
f.tailRenderer = nil
|
f.tailRenderer = nil
|
||||||
if tailRenderer != nil && tailRenderer.afterLastMap != f.indexedRange.firstRenderedMap {
|
if tailRenderer != nil && tailRenderer.renderBefore != f.indexedRange.maps.First() {
|
||||||
tailRenderer = nil
|
tailRenderer = nil
|
||||||
}
|
}
|
||||||
if tailRenderer == nil {
|
if tailRenderer == nil {
|
||||||
var err error
|
var err error
|
||||||
tailRenderer, err = f.renderMapsBefore(f.indexedRange.firstRenderedMap)
|
tailRenderer, err = f.renderMapsBefore(f.indexedRange.maps.First())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Error creating log index tail renderer", "error", err)
|
log.Error("Error creating log index tail renderer", "error", err)
|
||||||
return false
|
return false
|
||||||
|
|
@ -261,7 +261,7 @@ func (f *FilterMaps) tryIndexTail() bool {
|
||||||
f.lastLogTailIndex = time.Now()
|
f.lastLogTailIndex = time.Now()
|
||||||
f.startedTailIndexAt = f.lastLogTailIndex
|
f.startedTailIndexAt = f.lastLogTailIndex
|
||||||
f.startedTailIndex = true
|
f.startedTailIndex = true
|
||||||
f.ptrTailIndex = f.indexedRange.firstIndexedBlock - f.tailPartialBlocks()
|
f.ptrTailIndex = f.indexedRange.blocks.First() - f.tailPartialBlocks()
|
||||||
}
|
}
|
||||||
done, err := tailRenderer.run(func() bool {
|
done, err := tailRenderer.run(func() bool {
|
||||||
f.processEvents()
|
f.processEvents()
|
||||||
|
|
@ -269,14 +269,14 @@ func (f *FilterMaps) tryIndexTail() bool {
|
||||||
}, func() {
|
}, func() {
|
||||||
tpb, ttb := f.tailPartialBlocks(), f.tailTargetBlock()
|
tpb, ttb := f.tailPartialBlocks(), f.tailTargetBlock()
|
||||||
remaining := uint64(1)
|
remaining := uint64(1)
|
||||||
if f.indexedRange.firstIndexedBlock > ttb+tpb {
|
if f.indexedRange.blocks.First() > ttb+tpb {
|
||||||
remaining = f.indexedRange.firstIndexedBlock - ttb - tpb
|
remaining = f.indexedRange.blocks.First() - ttb - tpb
|
||||||
}
|
}
|
||||||
if f.indexedRange.hasIndexedBlocks() && f.ptrTailIndex >= f.indexedRange.firstIndexedBlock &&
|
if f.indexedRange.hasIndexedBlocks() && f.ptrTailIndex >= f.indexedRange.blocks.First() &&
|
||||||
(!f.loggedTailIndex || time.Since(f.lastLogTailIndex) > logFrequency) {
|
(!f.loggedTailIndex || time.Since(f.lastLogTailIndex) > logFrequency) {
|
||||||
log.Info("Log index tail rendering in progress",
|
log.Info("Log index tail rendering in progress",
|
||||||
"first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
|
"first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(),
|
||||||
"processed", f.ptrTailIndex-f.indexedRange.firstIndexedBlock+tpb,
|
"processed", f.ptrTailIndex-f.indexedRange.blocks.First()+tpb,
|
||||||
"remaining", remaining,
|
"remaining", remaining,
|
||||||
"next tail epoch percentage", f.indexedRange.tailPartialEpoch*100/f.mapsPerEpoch,
|
"next tail epoch percentage", f.indexedRange.tailPartialEpoch*100/f.mapsPerEpoch,
|
||||||
"elapsed", common.PrettyDuration(time.Since(f.startedTailIndexAt)))
|
"elapsed", common.PrettyDuration(time.Since(f.startedTailIndexAt)))
|
||||||
|
|
@ -293,10 +293,10 @@ func (f *FilterMaps) tryIndexTail() bool {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if f.loggedTailIndex {
|
if f.loggedTailIndex && f.indexedRange.hasIndexedBlocks() {
|
||||||
log.Info("Log index tail rendering finished",
|
log.Info("Log index tail rendering finished",
|
||||||
"first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
|
"first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(),
|
||||||
"processed", f.ptrTailIndex-f.indexedRange.firstIndexedBlock,
|
"processed", f.ptrTailIndex-f.indexedRange.blocks.First(),
|
||||||
"elapsed", common.PrettyDuration(time.Since(f.startedTailIndexAt)))
|
"elapsed", common.PrettyDuration(time.Since(f.startedTailIndexAt)))
|
||||||
f.loggedTailIndex = false
|
f.loggedTailIndex = false
|
||||||
}
|
}
|
||||||
|
|
@ -309,7 +309,7 @@ func (f *FilterMaps) tryIndexTail() bool {
|
||||||
// data from the database and is also called while running head indexing.
|
// data from the database and is also called while running head indexing.
|
||||||
func (f *FilterMaps) tryUnindexTail() bool {
|
func (f *FilterMaps) tryUnindexTail() bool {
|
||||||
for {
|
for {
|
||||||
firstEpoch := (f.indexedRange.firstRenderedMap - f.indexedRange.tailPartialEpoch) >> f.logMapsPerEpoch
|
firstEpoch := (f.indexedRange.maps.First() - f.indexedRange.tailPartialEpoch) >> f.logMapsPerEpoch
|
||||||
if f.needTailEpoch(firstEpoch) {
|
if f.needTailEpoch(firstEpoch) {
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|
@ -320,19 +320,19 @@ func (f *FilterMaps) tryUnindexTail() bool {
|
||||||
if !f.startedTailUnindex {
|
if !f.startedTailUnindex {
|
||||||
f.startedTailUnindexAt = time.Now()
|
f.startedTailUnindexAt = time.Now()
|
||||||
f.startedTailUnindex = true
|
f.startedTailUnindex = true
|
||||||
f.ptrTailUnindexMap = f.indexedRange.firstRenderedMap - f.indexedRange.tailPartialEpoch
|
f.ptrTailUnindexMap = f.indexedRange.maps.First() - f.indexedRange.tailPartialEpoch
|
||||||
f.ptrTailUnindexBlock = f.indexedRange.firstIndexedBlock - f.tailPartialBlocks()
|
f.ptrTailUnindexBlock = f.indexedRange.blocks.First() - f.tailPartialBlocks()
|
||||||
}
|
}
|
||||||
if err := f.deleteTailEpoch(firstEpoch); err != nil {
|
if err := f.deleteTailEpoch(firstEpoch); err != nil {
|
||||||
log.Error("Log index tail epoch unindexing failed", "error", err)
|
log.Error("Log index tail epoch unindexing failed", "error", err)
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if f.startedTailUnindex {
|
if f.startedTailUnindex && f.indexedRange.hasIndexedBlocks() {
|
||||||
log.Info("Log index tail unindexing finished",
|
log.Info("Log index tail unindexing finished",
|
||||||
"first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
|
"first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(),
|
||||||
"removed maps", f.indexedRange.firstRenderedMap-f.ptrTailUnindexMap,
|
"removed maps", f.indexedRange.maps.First()-f.ptrTailUnindexMap,
|
||||||
"removed blocks", f.indexedRange.firstIndexedBlock-f.tailPartialBlocks()-f.ptrTailUnindexBlock,
|
"removed blocks", f.indexedRange.blocks.First()-f.tailPartialBlocks()-f.ptrTailUnindexBlock,
|
||||||
"elapsed", common.PrettyDuration(time.Since(f.startedTailUnindexAt)))
|
"elapsed", common.PrettyDuration(time.Since(f.startedTailUnindexAt)))
|
||||||
f.startedTailUnindex = false
|
f.startedTailUnindex = false
|
||||||
}
|
}
|
||||||
|
|
@ -342,7 +342,7 @@ func (f *FilterMaps) tryUnindexTail() bool {
|
||||||
// needTailEpoch returns true if the given tail epoch needs to be kept
|
// needTailEpoch returns true if the given tail epoch needs to be kept
|
||||||
// according to the current tail target, false if it can be removed.
|
// according to the current tail target, false if it can be removed.
|
||||||
func (f *FilterMaps) needTailEpoch(epoch uint32) bool {
|
func (f *FilterMaps) needTailEpoch(epoch uint32) bool {
|
||||||
firstEpoch := f.indexedRange.firstRenderedMap >> f.logMapsPerEpoch
|
firstEpoch := f.indexedRange.maps.First() >> f.logMapsPerEpoch
|
||||||
if epoch > firstEpoch {
|
if epoch > firstEpoch {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
@ -382,15 +382,15 @@ func (f *FilterMaps) tailPartialBlocks() uint64 {
|
||||||
if f.indexedRange.tailPartialEpoch == 0 {
|
if f.indexedRange.tailPartialEpoch == 0 {
|
||||||
return 0
|
return 0
|
||||||
}
|
}
|
||||||
end, _, err := f.getLastBlockOfMap(f.indexedRange.firstRenderedMap - f.mapsPerEpoch + f.indexedRange.tailPartialEpoch - 1)
|
end, _, err := f.getLastBlockOfMap(f.indexedRange.maps.First() - f.mapsPerEpoch + f.indexedRange.tailPartialEpoch - 1)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Error fetching last block of map", "mapIndex", f.indexedRange.firstRenderedMap-f.mapsPerEpoch+f.indexedRange.tailPartialEpoch-1, "error", err)
|
log.Error("Error fetching last block of map", "mapIndex", f.indexedRange.maps.First()-f.mapsPerEpoch+f.indexedRange.tailPartialEpoch-1, "error", err)
|
||||||
}
|
}
|
||||||
var start uint64
|
var start uint64
|
||||||
if f.indexedRange.firstRenderedMap-f.mapsPerEpoch > 0 {
|
if f.indexedRange.maps.First()-f.mapsPerEpoch > 0 {
|
||||||
start, _, err = f.getLastBlockOfMap(f.indexedRange.firstRenderedMap - f.mapsPerEpoch - 1)
|
start, _, err = f.getLastBlockOfMap(f.indexedRange.maps.First() - f.mapsPerEpoch - 1)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Error fetching last block of map", "mapIndex", f.indexedRange.firstRenderedMap-f.mapsPerEpoch-1, "error", err)
|
log.Error("Error fetching last block of map", "mapIndex", f.indexedRange.maps.First()-f.mapsPerEpoch-1, "error", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return end - start
|
return end - start
|
||||||
|
|
@ -399,5 +399,5 @@ func (f *FilterMaps) tailPartialBlocks() uint64 {
|
||||||
// targetHeadIndexed returns true if the current log index is consistent with
|
// targetHeadIndexed returns true if the current log index is consistent with
|
||||||
// targetView with its head block fully rendered.
|
// targetView with its head block fully rendered.
|
||||||
func (f *FilterMaps) targetHeadIndexed() bool {
|
func (f *FilterMaps) targetHeadIndexed() bool {
|
||||||
return equalViews(f.targetView, f.indexedView) && f.indexedRange.headBlockIndexed
|
return equalViews(f.targetView, f.indexedView) && f.indexedRange.headIndexed
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -144,16 +144,16 @@ func TestIndexerRandomRange(t *testing.T) {
|
||||||
expTailBlock++
|
expTailBlock++
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if ts.fm.indexedRange.afterLastIndexedBlock != uint64(head+1) {
|
if ts.fm.indexedRange.blocks.Last() != uint64(head) {
|
||||||
ts.t.Fatalf("Invalid index head (expected #%d, got #%d)", head, ts.fm.indexedRange.afterLastIndexedBlock-1)
|
ts.t.Fatalf("Invalid index head (expected #%d, got #%d)", head, ts.fm.indexedRange.blocks.Last())
|
||||||
}
|
}
|
||||||
expHeadDelimiter := expdpos(uint64(head))
|
expHeadDelimiter := expdpos(uint64(head))
|
||||||
if ts.fm.indexedRange.headBlockDelimiter != expHeadDelimiter {
|
if ts.fm.indexedRange.headDelimiter != expHeadDelimiter {
|
||||||
ts.t.Fatalf("Invalid index head delimiter pointer (expected %d, got %d)", expHeadDelimiter, ts.fm.indexedRange.headBlockDelimiter)
|
ts.t.Fatalf("Invalid index head delimiter pointer (expected %d, got %d)", expHeadDelimiter, ts.fm.indexedRange.headDelimiter)
|
||||||
}
|
}
|
||||||
|
|
||||||
if ts.fm.indexedRange.firstIndexedBlock != expTailBlock {
|
if ts.fm.indexedRange.blocks.First() != expTailBlock {
|
||||||
ts.t.Fatalf("Invalid index tail block (expected #%d, got #%d)", expTailBlock, ts.fm.indexedRange.firstIndexedBlock)
|
ts.t.Fatalf("Invalid index tail block (expected #%d, got #%d)", expTailBlock, ts.fm.indexedRange.blocks.First())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -46,12 +46,12 @@ var (
|
||||||
// mapRenderer represents a process that renders filter maps in a specified
|
// mapRenderer represents a process that renders filter maps in a specified
|
||||||
// range according to the actual targetView.
|
// range according to the actual targetView.
|
||||||
type mapRenderer struct {
|
type mapRenderer struct {
|
||||||
f *FilterMaps
|
f *FilterMaps
|
||||||
afterLastMap uint32
|
renderBefore uint32
|
||||||
currentMap *renderedMap
|
currentMap *renderedMap
|
||||||
finishedMaps map[uint32]*renderedMap
|
finishedMaps map[uint32]*renderedMap
|
||||||
firstFinished, afterLastFinished uint32
|
finished common.Range[uint32]
|
||||||
iterator *logIterator
|
iterator *logIterator
|
||||||
}
|
}
|
||||||
|
|
||||||
// renderedMap represents a single filter map that is being rendered in memory.
|
// renderedMap represents a single filter map that is being rendered in memory.
|
||||||
|
|
@ -74,22 +74,22 @@ func (r *renderedMap) firstBlock() uint64 {
|
||||||
// specified map index boundary, starting from the latest available starting
|
// specified map index boundary, starting from the latest available starting
|
||||||
// point that is consistent with the current targetView.
|
// point that is consistent with the current targetView.
|
||||||
// The renderer ensures that filterMapsRange, indexedView and the actual map
|
// The renderer ensures that filterMapsRange, indexedView and the actual map
|
||||||
// data are always consistent with each other. If afterLastMap is greater than
|
// data are always consistent with each other. If renderBefore is greater than
|
||||||
// the latest existing rendered map then indexedView is updated to targetView,
|
// the latest existing rendered map then indexedView is updated to targetView,
|
||||||
// otherwise it is checked that the rendered range is consistent with both
|
// otherwise it is checked that the rendered range is consistent with both
|
||||||
// views.
|
// views.
|
||||||
func (f *FilterMaps) renderMapsBefore(afterLastMap uint32) (*mapRenderer, error) {
|
func (f *FilterMaps) renderMapsBefore(renderBefore uint32) (*mapRenderer, error) {
|
||||||
nextMap, startBlock, startLvPtr, err := f.lastCanonicalMapBoundaryBefore(afterLastMap)
|
nextMap, startBlock, startLvPtr, err := f.lastCanonicalMapBoundaryBefore(renderBefore)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
if snapshot := f.lastCanonicalSnapshotBefore(afterLastMap); snapshot != nil && snapshot.mapIndex >= nextMap {
|
if snapshot := f.lastCanonicalSnapshotBefore(renderBefore); snapshot != nil && snapshot.mapIndex >= nextMap {
|
||||||
return f.renderMapsFromSnapshot(snapshot)
|
return f.renderMapsFromSnapshot(snapshot)
|
||||||
}
|
}
|
||||||
if nextMap >= afterLastMap {
|
if nextMap >= renderBefore {
|
||||||
return nil, nil
|
return nil, nil
|
||||||
}
|
}
|
||||||
return f.renderMapsFromMapBoundary(nextMap, afterLastMap, startBlock, startLvPtr)
|
return f.renderMapsFromMapBoundary(nextMap, renderBefore, startBlock, startLvPtr)
|
||||||
}
|
}
|
||||||
|
|
||||||
// renderMapsFromSnapshot creates a mapRenderer that starts rendering from a
|
// renderMapsFromSnapshot creates a mapRenderer that starts rendering from a
|
||||||
|
|
@ -108,17 +108,16 @@ func (f *FilterMaps) renderMapsFromSnapshot(cp *renderedMap) (*mapRenderer, erro
|
||||||
lastBlock: cp.lastBlock,
|
lastBlock: cp.lastBlock,
|
||||||
blockLvPtrs: cp.blockLvPtrs,
|
blockLvPtrs: cp.blockLvPtrs,
|
||||||
},
|
},
|
||||||
finishedMaps: make(map[uint32]*renderedMap),
|
finishedMaps: make(map[uint32]*renderedMap),
|
||||||
firstFinished: cp.mapIndex,
|
finished: common.NewRange(cp.mapIndex, 0),
|
||||||
afterLastFinished: cp.mapIndex,
|
renderBefore: math.MaxUint32,
|
||||||
afterLastMap: math.MaxUint32,
|
iterator: iter,
|
||||||
iterator: iter,
|
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// renderMapsFromMapBoundary creates a mapRenderer that starts rendering at a
|
// renderMapsFromMapBoundary creates a mapRenderer that starts rendering at a
|
||||||
// map boundary.
|
// map boundary.
|
||||||
func (f *FilterMaps) renderMapsFromMapBoundary(firstMap, afterLastMap uint32, startBlock, startLvPtr uint64) (*mapRenderer, error) {
|
func (f *FilterMaps) renderMapsFromMapBoundary(firstMap, renderBefore uint32, startBlock, startLvPtr uint64) (*mapRenderer, error) {
|
||||||
iter, err := f.newLogIteratorFromMapBoundary(firstMap, startBlock, startLvPtr)
|
iter, err := f.newLogIteratorFromMapBoundary(firstMap, startBlock, startLvPtr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("failed to create log iterator from map boundary %d: %v", firstMap, err)
|
return nil, fmt.Errorf("failed to create log iterator from map boundary %d: %v", firstMap, err)
|
||||||
|
|
@ -130,22 +129,21 @@ func (f *FilterMaps) renderMapsFromMapBoundary(firstMap, afterLastMap uint32, st
|
||||||
mapIndex: firstMap,
|
mapIndex: firstMap,
|
||||||
lastBlock: iter.blockNumber,
|
lastBlock: iter.blockNumber,
|
||||||
},
|
},
|
||||||
finishedMaps: make(map[uint32]*renderedMap),
|
finishedMaps: make(map[uint32]*renderedMap),
|
||||||
firstFinished: firstMap,
|
finished: common.NewRange(firstMap, 0),
|
||||||
afterLastFinished: firstMap,
|
renderBefore: renderBefore,
|
||||||
afterLastMap: afterLastMap,
|
iterator: iter,
|
||||||
iterator: iter,
|
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// lastCanonicalSnapshotBefore returns the latest cached snapshot that matches
|
// lastCanonicalSnapshotBefore returns the latest cached snapshot that matches
|
||||||
// the current targetView.
|
// the current targetView.
|
||||||
func (f *FilterMaps) lastCanonicalSnapshotBefore(afterLastMap uint32) *renderedMap {
|
func (f *FilterMaps) lastCanonicalSnapshotBefore(renderBefore uint32) *renderedMap {
|
||||||
var best *renderedMap
|
var best *renderedMap
|
||||||
for _, blockNumber := range f.renderSnapshots.Keys() {
|
for _, blockNumber := range f.renderSnapshots.Keys() {
|
||||||
if cp, _ := f.renderSnapshots.Get(blockNumber); cp != nil && blockNumber < f.indexedRange.afterLastIndexedBlock &&
|
if cp, _ := f.renderSnapshots.Get(blockNumber); cp != nil && blockNumber < f.indexedRange.blocks.AfterLast() &&
|
||||||
blockNumber <= f.targetView.headNumber && f.targetView.getBlockId(blockNumber) == cp.lastBlockId &&
|
blockNumber <= f.targetView.headNumber && f.targetView.getBlockId(blockNumber) == cp.lastBlockId &&
|
||||||
cp.mapIndex < afterLastMap && (best == nil || blockNumber > best.lastBlock) {
|
cp.mapIndex < renderBefore && (best == nil || blockNumber > best.lastBlock) {
|
||||||
best = cp
|
best = cp
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -158,11 +156,11 @@ func (f *FilterMaps) lastCanonicalSnapshotBefore(afterLastMap uint32) *renderedM
|
||||||
// or the boundary of a currently rendered map.
|
// or the boundary of a currently rendered map.
|
||||||
// Along with the next map index where the rendering can be started, the number
|
// Along with the next map index where the rendering can be started, the number
|
||||||
// and starting log value pointer of the last block is also returned.
|
// and starting log value pointer of the last block is also returned.
|
||||||
func (f *FilterMaps) lastCanonicalMapBoundaryBefore(afterLastMap uint32) (nextMap uint32, startBlock, startLvPtr uint64, err error) {
|
func (f *FilterMaps) lastCanonicalMapBoundaryBefore(renderBefore uint32) (nextMap uint32, startBlock, startLvPtr uint64, err error) {
|
||||||
if !f.indexedRange.initialized {
|
if !f.indexedRange.initialized {
|
||||||
return 0, 0, 0, nil
|
return 0, 0, 0, nil
|
||||||
}
|
}
|
||||||
mapIndex := afterLastMap
|
mapIndex := renderBefore
|
||||||
for {
|
for {
|
||||||
var ok bool
|
var ok bool
|
||||||
if mapIndex, ok = f.lastMapBoundaryBefore(mapIndex); !ok {
|
if mapIndex, ok = f.lastMapBoundaryBefore(mapIndex); !ok {
|
||||||
|
|
@ -188,18 +186,18 @@ func (f *FilterMaps) lastCanonicalMapBoundaryBefore(afterLastMap uint32) (nextMa
|
||||||
// lastMapBoundaryBefore returns the latest map boundary before the specified
|
// lastMapBoundaryBefore returns the latest map boundary before the specified
|
||||||
// map index.
|
// map index.
|
||||||
func (f *FilterMaps) lastMapBoundaryBefore(mapIndex uint32) (uint32, bool) {
|
func (f *FilterMaps) lastMapBoundaryBefore(mapIndex uint32) (uint32, bool) {
|
||||||
if !f.indexedRange.initialized || f.indexedRange.afterLastRenderedMap == 0 {
|
if !f.indexedRange.initialized || f.indexedRange.maps.AfterLast() == 0 {
|
||||||
return 0, false
|
return 0, false
|
||||||
}
|
}
|
||||||
if mapIndex > f.indexedRange.afterLastRenderedMap {
|
if mapIndex > f.indexedRange.maps.AfterLast() {
|
||||||
mapIndex = f.indexedRange.afterLastRenderedMap
|
mapIndex = f.indexedRange.maps.AfterLast()
|
||||||
}
|
}
|
||||||
if mapIndex > f.indexedRange.firstRenderedMap {
|
if mapIndex > f.indexedRange.maps.First() {
|
||||||
return mapIndex - 1, true
|
return mapIndex - 1, true
|
||||||
}
|
}
|
||||||
if mapIndex+f.mapsPerEpoch > f.indexedRange.firstRenderedMap {
|
if mapIndex+f.mapsPerEpoch > f.indexedRange.maps.First() {
|
||||||
if mapIndex > f.indexedRange.firstRenderedMap-f.mapsPerEpoch+f.indexedRange.tailPartialEpoch {
|
if mapIndex > f.indexedRange.maps.First()-f.mapsPerEpoch+f.indexedRange.tailPartialEpoch {
|
||||||
mapIndex = f.indexedRange.firstRenderedMap - f.mapsPerEpoch + f.indexedRange.tailPartialEpoch
|
mapIndex = f.indexedRange.maps.First() - f.mapsPerEpoch + f.indexedRange.tailPartialEpoch
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
mapIndex = (mapIndex >> f.logMapsPerEpoch) << f.logMapsPerEpoch
|
mapIndex = (mapIndex >> f.logMapsPerEpoch) << f.logMapsPerEpoch
|
||||||
|
|
@ -218,19 +216,19 @@ func (f *FilterMaps) emptyFilterMap() filterMap {
|
||||||
// loadHeadSnapshot loads the last rendered map from the database and creates
|
// loadHeadSnapshot loads the last rendered map from the database and creates
|
||||||
// a snapshot.
|
// a snapshot.
|
||||||
func (f *FilterMaps) loadHeadSnapshot() error {
|
func (f *FilterMaps) loadHeadSnapshot() error {
|
||||||
fm, err := f.getFilterMap(f.indexedRange.afterLastRenderedMap - 1)
|
fm, err := f.getFilterMap(f.indexedRange.maps.Last())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("failed to load head snapshot map %d: %v", f.indexedRange.afterLastRenderedMap-1, err)
|
return fmt.Errorf("failed to load head snapshot map %d: %v", f.indexedRange.maps.Last(), err)
|
||||||
}
|
}
|
||||||
lastBlock, _, err := f.getLastBlockOfMap(f.indexedRange.afterLastRenderedMap - 1)
|
lastBlock, _, err := f.getLastBlockOfMap(f.indexedRange.maps.Last())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("failed to retrieve last block of head snapshot map %d: %v", f.indexedRange.afterLastRenderedMap-1, err)
|
return fmt.Errorf("failed to retrieve last block of head snapshot map %d: %v", f.indexedRange.maps.Last(), err)
|
||||||
}
|
}
|
||||||
var firstBlock uint64
|
var firstBlock uint64
|
||||||
if f.indexedRange.afterLastRenderedMap > 1 {
|
if f.indexedRange.maps.AfterLast() > 1 {
|
||||||
prevLastBlock, _, err := f.getLastBlockOfMap(f.indexedRange.afterLastRenderedMap - 2)
|
prevLastBlock, _, err := f.getLastBlockOfMap(f.indexedRange.maps.Last() - 1)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("failed to retrieve last block of map %d before head snapshot: %v", f.indexedRange.afterLastRenderedMap-2, err)
|
return fmt.Errorf("failed to retrieve last block of map %d before head snapshot: %v", f.indexedRange.maps.Last()-1, err)
|
||||||
}
|
}
|
||||||
firstBlock = prevLastBlock + 1
|
firstBlock = prevLastBlock + 1
|
||||||
}
|
}
|
||||||
|
|
@ -241,14 +239,14 @@ func (f *FilterMaps) loadHeadSnapshot() error {
|
||||||
return fmt.Errorf("failed to retrieve log value pointer of head snapshot block %d: %v", firstBlock+uint64(i), err)
|
return fmt.Errorf("failed to retrieve log value pointer of head snapshot block %d: %v", firstBlock+uint64(i), err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
f.renderSnapshots.Add(f.indexedRange.afterLastIndexedBlock-1, &renderedMap{
|
f.renderSnapshots.Add(f.indexedRange.blocks.Last(), &renderedMap{
|
||||||
filterMap: fm,
|
filterMap: fm,
|
||||||
mapIndex: f.indexedRange.afterLastRenderedMap - 1,
|
mapIndex: f.indexedRange.maps.Last(),
|
||||||
lastBlock: f.indexedRange.afterLastIndexedBlock - 1,
|
lastBlock: f.indexedRange.blocks.Last(),
|
||||||
lastBlockId: f.indexedView.getBlockId(f.indexedRange.afterLastIndexedBlock - 1),
|
lastBlockId: f.indexedView.getBlockId(f.indexedRange.blocks.Last()),
|
||||||
blockLvPtrs: lvPtrs,
|
blockLvPtrs: lvPtrs,
|
||||||
finished: true,
|
finished: true,
|
||||||
headDelimiter: f.indexedRange.headBlockDelimiter,
|
headDelimiter: f.indexedRange.headDelimiter,
|
||||||
})
|
})
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
@ -277,14 +275,14 @@ func (r *mapRenderer) run(stopCb func() bool, writeCb func()) (bool, error) {
|
||||||
}
|
}
|
||||||
// map finished
|
// map finished
|
||||||
r.finishedMaps[r.currentMap.mapIndex] = r.currentMap
|
r.finishedMaps[r.currentMap.mapIndex] = r.currentMap
|
||||||
r.afterLastFinished++
|
r.finished.SetLast(r.finished.AfterLast())
|
||||||
if len(r.finishedMaps) >= maxMapsPerBatch || r.afterLastFinished&(r.f.baseRowGroupLength-1) == 0 {
|
if len(r.finishedMaps) >= maxMapsPerBatch || r.finished.AfterLast()&(r.f.baseRowGroupLength-1) == 0 {
|
||||||
if err := r.writeFinishedMaps(stopCb); err != nil {
|
if err := r.writeFinishedMaps(stopCb); err != nil {
|
||||||
return false, err
|
return false, err
|
||||||
}
|
}
|
||||||
writeCb()
|
writeCb()
|
||||||
}
|
}
|
||||||
if r.afterLastFinished == r.afterLastMap || r.iterator.finished {
|
if r.finished.AfterLast() == r.renderBefore || r.iterator.finished {
|
||||||
if err := r.writeFinishedMaps(stopCb); err != nil {
|
if err := r.writeFinishedMaps(stopCb); err != nil {
|
||||||
return false, err
|
return false, err
|
||||||
}
|
}
|
||||||
|
|
@ -293,7 +291,7 @@ func (r *mapRenderer) run(stopCb func() bool, writeCb func()) (bool, error) {
|
||||||
}
|
}
|
||||||
r.currentMap = &renderedMap{
|
r.currentMap = &renderedMap{
|
||||||
filterMap: r.f.emptyFilterMap(),
|
filterMap: r.f.emptyFilterMap(),
|
||||||
mapIndex: r.afterLastFinished,
|
mapIndex: r.finished.AfterLast(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -346,7 +344,7 @@ func (r *mapRenderer) renderCurrentMap(stopCb func() bool) (bool, error) {
|
||||||
if r.iterator.blockStart {
|
if r.iterator.blockStart {
|
||||||
r.currentMap.blockLvPtrs = append(r.currentMap.blockLvPtrs, r.iterator.lvIndex)
|
r.currentMap.blockLvPtrs = append(r.currentMap.blockLvPtrs, r.iterator.lvIndex)
|
||||||
}
|
}
|
||||||
if !r.f.testDisableSnapshots && r.afterLastMap >= r.f.indexedRange.afterLastRenderedMap &&
|
if !r.f.testDisableSnapshots && r.renderBefore >= r.f.indexedRange.maps.AfterLast() &&
|
||||||
(r.iterator.delimiter || r.iterator.finished) {
|
(r.iterator.delimiter || r.iterator.finished) {
|
||||||
r.makeSnapshot()
|
r.makeSnapshot()
|
||||||
}
|
}
|
||||||
|
|
@ -403,7 +401,7 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error {
|
||||||
mapIndices []uint32
|
mapIndices []uint32
|
||||||
rows []FilterRow
|
rows []FilterRow
|
||||||
)
|
)
|
||||||
for mapIndex := r.firstFinished; mapIndex < r.afterLastFinished; mapIndex++ {
|
for mapIndex := r.finished.First(); mapIndex < r.finished.AfterLast(); mapIndex++ {
|
||||||
row := r.finishedMaps[mapIndex].filterMap[rowIndex]
|
row := r.finishedMaps[mapIndex].filterMap[rowIndex]
|
||||||
if fm, _ := r.f.filterMapCache.Get(mapIndex); fm != nil && row.Equal(fm[rowIndex]) {
|
if fm, _ := r.f.filterMapCache.Get(mapIndex); fm != nil && row.Equal(fm[rowIndex]) {
|
||||||
continue
|
continue
|
||||||
|
|
@ -411,8 +409,8 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error {
|
||||||
mapIndices = append(mapIndices, mapIndex)
|
mapIndices = append(mapIndices, mapIndex)
|
||||||
rows = append(rows, row)
|
rows = append(rows, row)
|
||||||
}
|
}
|
||||||
if newRange.afterLastRenderedMap == r.afterLastFinished { // head updated; remove future entries
|
if newRange.maps.AfterLast() == r.finished.AfterLast() { // head updated; remove future entries
|
||||||
for mapIndex := r.afterLastFinished; mapIndex < oldRange.afterLastRenderedMap; mapIndex++ {
|
for mapIndex := r.finished.AfterLast(); mapIndex < oldRange.maps.AfterLast(); mapIndex++ {
|
||||||
if fm, _ := r.f.filterMapCache.Get(mapIndex); fm != nil && len(fm[rowIndex]) == 0 {
|
if fm, _ := r.f.filterMapCache.Get(mapIndex); fm != nil && len(fm[rowIndex]) == 0 {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
@ -426,24 +424,24 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error {
|
||||||
checkWriteCnt()
|
checkWriteCnt()
|
||||||
}
|
}
|
||||||
// update filter map cache
|
// update filter map cache
|
||||||
if newRange.afterLastRenderedMap == r.afterLastFinished {
|
if newRange.maps.AfterLast() == r.finished.AfterLast() {
|
||||||
// head updated; cache new head maps and remove future entries
|
// head updated; cache new head maps and remove future entries
|
||||||
for mapIndex := r.firstFinished; mapIndex < r.afterLastFinished; mapIndex++ {
|
for mapIndex := r.finished.First(); mapIndex < r.finished.AfterLast(); mapIndex++ {
|
||||||
r.f.filterMapCache.Add(mapIndex, r.finishedMaps[mapIndex].filterMap)
|
r.f.filterMapCache.Add(mapIndex, r.finishedMaps[mapIndex].filterMap)
|
||||||
}
|
}
|
||||||
for mapIndex := r.afterLastFinished; mapIndex < oldRange.afterLastRenderedMap; mapIndex++ {
|
for mapIndex := r.finished.AfterLast(); mapIndex < oldRange.maps.AfterLast(); mapIndex++ {
|
||||||
r.f.filterMapCache.Remove(mapIndex)
|
r.f.filterMapCache.Remove(mapIndex)
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
// head not updated; do not cache maps during tail rendering because we
|
// head not updated; do not cache maps during tail rendering because we
|
||||||
// need head maps to be available in the cache
|
// need head maps to be available in the cache
|
||||||
for mapIndex := r.firstFinished; mapIndex < r.afterLastFinished; mapIndex++ {
|
for mapIndex := r.finished.First(); mapIndex < r.finished.AfterLast(); mapIndex++ {
|
||||||
r.f.filterMapCache.Remove(mapIndex)
|
r.f.filterMapCache.Remove(mapIndex)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// add or update block pointers
|
// add or update block pointers
|
||||||
blockNumber := r.finishedMaps[r.firstFinished].firstBlock()
|
blockNumber := r.finishedMaps[r.finished.First()].firstBlock()
|
||||||
for mapIndex := r.firstFinished; mapIndex < r.afterLastFinished; mapIndex++ {
|
for mapIndex := r.finished.First(); mapIndex < r.finished.AfterLast(); mapIndex++ {
|
||||||
renderedMap := r.finishedMaps[mapIndex]
|
renderedMap := r.finishedMaps[mapIndex]
|
||||||
r.f.storeLastBlockOfMap(batch, mapIndex, renderedMap.lastBlock, renderedMap.lastBlockId)
|
r.f.storeLastBlockOfMap(batch, mapIndex, renderedMap.lastBlock, renderedMap.lastBlockId)
|
||||||
checkWriteCnt()
|
checkWriteCnt()
|
||||||
|
|
@ -456,18 +454,18 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error {
|
||||||
blockNumber++
|
blockNumber++
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if newRange.afterLastRenderedMap == r.afterLastFinished { // head updated; remove future entries
|
if newRange.maps.AfterLast() == r.finished.AfterLast() { // head updated; remove future entries
|
||||||
for mapIndex := r.afterLastFinished; mapIndex < oldRange.afterLastRenderedMap; mapIndex++ {
|
for mapIndex := r.finished.AfterLast(); mapIndex < oldRange.maps.AfterLast(); mapIndex++ {
|
||||||
r.f.deleteLastBlockOfMap(batch, mapIndex)
|
r.f.deleteLastBlockOfMap(batch, mapIndex)
|
||||||
checkWriteCnt()
|
checkWriteCnt()
|
||||||
}
|
}
|
||||||
for ; blockNumber < oldRange.afterLastIndexedBlock; blockNumber++ {
|
for ; blockNumber < oldRange.blocks.AfterLast(); blockNumber++ {
|
||||||
r.f.deleteBlockLvPointer(batch, blockNumber)
|
r.f.deleteBlockLvPointer(batch, blockNumber)
|
||||||
checkWriteCnt()
|
checkWriteCnt()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
r.finishedMaps = make(map[uint32]*renderedMap)
|
r.finishedMaps = make(map[uint32]*renderedMap)
|
||||||
r.firstFinished = r.afterLastFinished
|
r.finished.SetFirst(r.finished.AfterLast())
|
||||||
r.f.setRange(batch, renderedView, newRange)
|
r.f.setRange(batch, renderedView, newRange)
|
||||||
if err := batch.Write(); err != nil {
|
if err := batch.Write(); err != nil {
|
||||||
log.Crit("Error writing log index update batch", "error", err)
|
log.Crit("Error writing log index update batch", "error", err)
|
||||||
|
|
@ -482,33 +480,33 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error {
|
||||||
// range to the unchanged region until all new map data is committed.
|
// range to the unchanged region until all new map data is committed.
|
||||||
func (r *mapRenderer) getTempRange() (filterMapsRange, error) {
|
func (r *mapRenderer) getTempRange() (filterMapsRange, error) {
|
||||||
tempRange := r.f.indexedRange
|
tempRange := r.f.indexedRange
|
||||||
if err := tempRange.addRenderedRange(r.firstFinished, r.firstFinished, r.afterLastMap, r.f.mapsPerEpoch); err != nil {
|
if err := tempRange.addRenderedRange(r.finished.First(), r.finished.First(), r.renderBefore, r.f.mapsPerEpoch); err != nil {
|
||||||
return filterMapsRange{}, fmt.Errorf("failed to update temporary rendered range: %v", err)
|
return filterMapsRange{}, fmt.Errorf("failed to update temporary rendered range: %v", err)
|
||||||
}
|
}
|
||||||
if tempRange.firstRenderedMap != r.f.indexedRange.firstRenderedMap {
|
if tempRange.maps.First() != r.f.indexedRange.maps.First() {
|
||||||
// first rendered map changed; update first indexed block
|
// first rendered map changed; update first indexed block
|
||||||
if tempRange.firstRenderedMap > 0 {
|
if tempRange.maps.First() > 0 {
|
||||||
lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.firstRenderedMap - 1)
|
firstBlock, _, err := r.f.getLastBlockOfMap(tempRange.maps.First() - 1)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return filterMapsRange{}, fmt.Errorf("failed to retrieve last block of map %d before temporary range: %v", tempRange.firstRenderedMap-1, err)
|
return filterMapsRange{}, fmt.Errorf("failed to retrieve last block of map %d before temporary range: %v", tempRange.maps.First()-1, err)
|
||||||
}
|
}
|
||||||
tempRange.firstIndexedBlock = lastBlock + 1
|
tempRange.blocks.SetFirst(firstBlock + 1) // firstBlock is probably partially rendered
|
||||||
} else {
|
} else {
|
||||||
tempRange.firstIndexedBlock = 0
|
tempRange.blocks.SetFirst(0)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if tempRange.afterLastRenderedMap != r.f.indexedRange.afterLastRenderedMap {
|
if tempRange.maps.AfterLast() != r.f.indexedRange.maps.AfterLast() {
|
||||||
// first rendered map changed; update first indexed block
|
// last rendered map changed; update last indexed block
|
||||||
if tempRange.afterLastRenderedMap > 0 {
|
if !tempRange.maps.IsEmpty() {
|
||||||
lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.afterLastRenderedMap - 1)
|
lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.maps.Last())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return filterMapsRange{}, fmt.Errorf("failed to retrieve last block of map %d at the end of temporary range: %v", tempRange.afterLastRenderedMap-1, err)
|
return filterMapsRange{}, fmt.Errorf("failed to retrieve last block of map %d at the end of temporary range: %v", tempRange.maps.Last(), err)
|
||||||
}
|
}
|
||||||
tempRange.afterLastIndexedBlock = lastBlock
|
tempRange.blocks.SetAfterLast(lastBlock) // lastBlock is probably partially rendered
|
||||||
} else {
|
} else {
|
||||||
tempRange.afterLastIndexedBlock = 0
|
tempRange.blocks.SetAfterLast(0)
|
||||||
}
|
}
|
||||||
tempRange.headBlockDelimiter = 0
|
tempRange.headDelimiter = 0
|
||||||
}
|
}
|
||||||
return tempRange, nil
|
return tempRange, nil
|
||||||
}
|
}
|
||||||
|
|
@ -518,39 +516,39 @@ func (r *mapRenderer) getTempRange() (filterMapsRange, error) {
|
||||||
func (r *mapRenderer) getUpdatedRange() (filterMapsRange, error) {
|
func (r *mapRenderer) getUpdatedRange() (filterMapsRange, error) {
|
||||||
// update filterMapsRange
|
// update filterMapsRange
|
||||||
newRange := r.f.indexedRange
|
newRange := r.f.indexedRange
|
||||||
if err := newRange.addRenderedRange(r.firstFinished, r.afterLastFinished, r.afterLastMap, r.f.mapsPerEpoch); err != nil {
|
if err := newRange.addRenderedRange(r.finished.First(), r.finished.AfterLast(), r.renderBefore, r.f.mapsPerEpoch); err != nil {
|
||||||
return filterMapsRange{}, fmt.Errorf("failed to update rendered range: %v", err)
|
return filterMapsRange{}, fmt.Errorf("failed to update rendered range: %v", err)
|
||||||
}
|
}
|
||||||
if newRange.firstRenderedMap != r.f.indexedRange.firstRenderedMap {
|
if newRange.maps.First() != r.f.indexedRange.maps.First() {
|
||||||
// first rendered map changed; update first indexed block
|
// first rendered map changed; update first indexed block
|
||||||
if newRange.firstRenderedMap > 0 {
|
if newRange.maps.First() > 0 {
|
||||||
lastBlock, _, err := r.f.getLastBlockOfMap(newRange.firstRenderedMap - 1)
|
firstBlock, _, err := r.f.getLastBlockOfMap(newRange.maps.First() - 1)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return filterMapsRange{}, fmt.Errorf("failed to retrieve last block of map %d before rendered range: %v", newRange.firstRenderedMap-1, err)
|
return filterMapsRange{}, fmt.Errorf("failed to retrieve last block of map %d before rendered range: %v", newRange.maps.First()-1, err)
|
||||||
}
|
}
|
||||||
newRange.firstIndexedBlock = lastBlock + 1
|
newRange.blocks.SetFirst(firstBlock + 1) // firstBlock is probably partially rendered
|
||||||
} else {
|
} else {
|
||||||
newRange.firstIndexedBlock = 0
|
newRange.blocks.SetFirst(0)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if newRange.afterLastRenderedMap == r.afterLastFinished {
|
if newRange.maps.AfterLast() == r.finished.AfterLast() {
|
||||||
// last rendered map changed; update last indexed block and head pointers
|
// last rendered map changed; update last indexed block and head pointers
|
||||||
lm := r.finishedMaps[r.afterLastFinished-1]
|
lm := r.finishedMaps[r.finished.Last()]
|
||||||
newRange.headBlockIndexed = lm.finished
|
newRange.headIndexed = lm.finished
|
||||||
if lm.finished {
|
if lm.finished {
|
||||||
newRange.afterLastIndexedBlock = r.f.targetView.headNumber + 1
|
newRange.blocks.SetLast(r.f.targetView.headNumber)
|
||||||
if lm.lastBlock != r.f.targetView.headNumber {
|
if lm.lastBlock != r.f.targetView.headNumber {
|
||||||
panic("map rendering finished but last block != head block")
|
panic("map rendering finished but last block != head block")
|
||||||
}
|
}
|
||||||
newRange.headBlockDelimiter = lm.headDelimiter
|
newRange.headDelimiter = lm.headDelimiter
|
||||||
} else {
|
} else {
|
||||||
newRange.afterLastIndexedBlock = lm.lastBlock
|
newRange.blocks.SetAfterLast(lm.lastBlock) // lastBlock is probably partially rendered
|
||||||
newRange.headBlockDelimiter = 0
|
newRange.headDelimiter = 0
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
// last rendered map not replaced; ensure that target chain view matches
|
// last rendered map not replaced; ensure that target chain view matches
|
||||||
// indexed chain view on the rendered section
|
// indexed chain view on the rendered section
|
||||||
if lastBlock := r.finishedMaps[r.afterLastFinished-1].lastBlock; !matchViews(r.f.indexedView, r.f.targetView, lastBlock) {
|
if lastBlock := r.finishedMaps[r.finished.Last()].lastBlock; !matchViews(r.f.indexedView, r.f.targetView, lastBlock) {
|
||||||
return filterMapsRange{}, errChainUpdate
|
return filterMapsRange{}, errChainUpdate
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -572,9 +570,9 @@ func (fmr *filterMapsRange) addRenderedRange(firstRendered, afterLastRendered, a
|
||||||
m uint32
|
m uint32
|
||||||
d int
|
d int
|
||||||
}
|
}
|
||||||
endpoints := []endpoint{{fmr.firstRenderedMap, 1}, {fmr.afterLastRenderedMap, -1}, {firstRendered, 1}, {afterLastRendered, -101}, {afterLastRemoved, 100}}
|
endpoints := []endpoint{{fmr.maps.First(), 1}, {fmr.maps.AfterLast(), -1}, {firstRendered, 1}, {afterLastRendered, -101}, {afterLastRemoved, 100}}
|
||||||
if fmr.tailPartialEpoch > 0 {
|
if fmr.tailPartialEpoch > 0 {
|
||||||
endpoints = append(endpoints, []endpoint{{fmr.firstRenderedMap - mapsPerEpoch, 1}, {fmr.firstRenderedMap - mapsPerEpoch + fmr.tailPartialEpoch, -1}}...)
|
endpoints = append(endpoints, []endpoint{{fmr.maps.First() - mapsPerEpoch, 1}, {fmr.maps.First() - mapsPerEpoch + fmr.tailPartialEpoch, -1}}...)
|
||||||
}
|
}
|
||||||
sort.Slice(endpoints, func(i, j int) bool { return endpoints[i].m < endpoints[j].m })
|
sort.Slice(endpoints, func(i, j int) bool { return endpoints[i].m < endpoints[j].m })
|
||||||
var (
|
var (
|
||||||
|
|
@ -597,14 +595,12 @@ func (fmr *filterMapsRange) addRenderedRange(firstRendered, afterLastRendered, a
|
||||||
case 0:
|
case 0:
|
||||||
// Initialized database, but no finished maps yet.
|
// Initialized database, but no finished maps yet.
|
||||||
fmr.tailPartialEpoch = 0
|
fmr.tailPartialEpoch = 0
|
||||||
fmr.firstRenderedMap = firstRendered
|
fmr.maps = common.NewRange(firstRendered, 0)
|
||||||
fmr.afterLastRenderedMap = firstRendered
|
|
||||||
|
|
||||||
case 2:
|
case 2:
|
||||||
// One rendered section (no partial tail epoch).
|
// One rendered section (no partial tail epoch).
|
||||||
fmr.tailPartialEpoch = 0
|
fmr.tailPartialEpoch = 0
|
||||||
fmr.firstRenderedMap = merged[0]
|
fmr.maps = common.NewRange(merged[0], merged[1]-merged[0])
|
||||||
fmr.afterLastRenderedMap = merged[1]
|
|
||||||
|
|
||||||
case 4:
|
case 4:
|
||||||
// Two rendered sections (with a gap).
|
// Two rendered sections (with a gap).
|
||||||
|
|
@ -614,8 +610,7 @@ func (fmr *filterMapsRange) addRenderedRange(firstRendered, afterLastRendered, a
|
||||||
return fmt.Errorf("invalid tail partial epoch: %v", merged)
|
return fmt.Errorf("invalid tail partial epoch: %v", merged)
|
||||||
}
|
}
|
||||||
fmr.tailPartialEpoch = merged[1] - merged[0]
|
fmr.tailPartialEpoch = merged[1] - merged[0]
|
||||||
fmr.firstRenderedMap = merged[2]
|
fmr.maps = common.NewRange(merged[2], merged[3]-merged[2])
|
||||||
fmr.afterLastRenderedMap = merged[3]
|
|
||||||
|
|
||||||
default:
|
default:
|
||||||
return fmt.Errorf("invalid number of rendered sections: %v", merged)
|
return fmt.Errorf("invalid number of rendered sections: %v", merged)
|
||||||
|
|
@ -643,12 +638,12 @@ func (f *FilterMaps) newLogIteratorFromBlockDelimiter(blockNumber uint64) (*logI
|
||||||
if blockNumber > f.targetView.headNumber {
|
if blockNumber > f.targetView.headNumber {
|
||||||
return nil, fmt.Errorf("iterator entry point %d after target chain head block %d", blockNumber, f.targetView.headNumber)
|
return nil, fmt.Errorf("iterator entry point %d after target chain head block %d", blockNumber, f.targetView.headNumber)
|
||||||
}
|
}
|
||||||
if blockNumber < f.indexedRange.firstIndexedBlock || blockNumber >= f.indexedRange.afterLastIndexedBlock {
|
if !f.indexedRange.blocks.Includes(blockNumber) {
|
||||||
return nil, errUnindexedRange
|
return nil, errUnindexedRange
|
||||||
}
|
}
|
||||||
var lvIndex uint64
|
var lvIndex uint64
|
||||||
if f.indexedRange.headBlockIndexed && blockNumber+1 == f.indexedRange.afterLastIndexedBlock {
|
if f.indexedRange.headIndexed && blockNumber+1 == f.indexedRange.blocks.AfterLast() {
|
||||||
lvIndex = f.indexedRange.headBlockDelimiter
|
lvIndex = f.indexedRange.headDelimiter
|
||||||
} else {
|
} else {
|
||||||
var err error
|
var err error
|
||||||
lvIndex, err = f.getBlockLvPointer(blockNumber + 1)
|
lvIndex, err = f.getBlockLvPointer(blockNumber + 1)
|
||||||
|
|
|
||||||
|
|
@ -61,11 +61,9 @@ type SyncRange struct {
|
||||||
// block range where the index has not changed since the last matcher sync
|
// block range where the index has not changed since the last matcher sync
|
||||||
// and therefore the set of matches found in this region is guaranteed to
|
// and therefore the set of matches found in this region is guaranteed to
|
||||||
// be valid and complete.
|
// be valid and complete.
|
||||||
Valid bool
|
ValidBlocks common.Range[uint64]
|
||||||
FirstValid, LastValid uint64
|
|
||||||
// block range indexed according to the given chain head.
|
// block range indexed according to the given chain head.
|
||||||
Indexed bool
|
IndexedBlocks common.Range[uint64]
|
||||||
FirstIndexed, LastIndexed uint64
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// GetPotentialMatches returns a list of logs that are potential matches for the
|
// GetPotentialMatches returns a list of logs that are potential matches for the
|
||||||
|
|
@ -143,11 +141,11 @@ func GetPotentialMatches(ctx context.Context, backend MatcherBackend, firstBlock
|
||||||
}
|
}
|
||||||
|
|
||||||
type matcherEnv struct {
|
type matcherEnv struct {
|
||||||
|
getLogStats runtimeStats // 64 bit aligned
|
||||||
ctx context.Context
|
ctx context.Context
|
||||||
backend MatcherBackend
|
backend MatcherBackend
|
||||||
params *Params
|
params *Params
|
||||||
matcher matcher
|
matcher matcher
|
||||||
getLogStats runtimeStats
|
|
||||||
firstIndex, lastIndex uint64
|
firstIndex, lastIndex uint64
|
||||||
firstMap, lastMap uint32
|
firstMap, lastMap uint32
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -19,6 +19,7 @@ package filtermaps
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -27,9 +28,8 @@ type FilterMapsMatcherBackend struct {
|
||||||
f *FilterMaps
|
f *FilterMaps
|
||||||
|
|
||||||
// these fields should be accessed under f.matchersLock mutex.
|
// these fields should be accessed under f.matchersLock mutex.
|
||||||
valid bool
|
validBlocks common.Range[uint64]
|
||||||
firstValid, lastValid uint64
|
syncCh chan SyncRange
|
||||||
syncCh chan SyncRange
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewMatcherBackend returns a FilterMapsMatcherBackend after registering it in
|
// NewMatcherBackend returns a FilterMapsMatcherBackend after registering it in
|
||||||
|
|
@ -43,11 +43,9 @@ func (f *FilterMaps) NewMatcherBackend() *FilterMapsMatcherBackend {
|
||||||
f.indexLock.RUnlock()
|
f.indexLock.RUnlock()
|
||||||
}()
|
}()
|
||||||
|
|
||||||
fm := &FilterMapsMatcherBackend{
|
fm := &FilterMapsMatcherBackend{f: f}
|
||||||
f: f,
|
if f.indexedRange.initialized {
|
||||||
valid: f.indexedRange.initialized && f.indexedRange.afterLastIndexedBlock > f.indexedRange.firstIndexedBlock,
|
fm.validBlocks = f.indexedRange.blocks
|
||||||
firstValid: f.indexedRange.firstIndexedBlock,
|
|
||||||
lastValid: f.indexedRange.afterLastIndexedBlock - 1,
|
|
||||||
}
|
}
|
||||||
f.matchers[fm] = struct{}{}
|
f.matchers[fm] = struct{}{}
|
||||||
return fm
|
return fm
|
||||||
|
|
@ -122,28 +120,16 @@ func (fm *FilterMapsMatcherBackend) synced() {
|
||||||
fm.f.indexLock.RUnlock()
|
fm.f.indexLock.RUnlock()
|
||||||
}()
|
}()
|
||||||
|
|
||||||
var (
|
indexedBlocks := fm.f.indexedRange.blocks
|
||||||
indexed bool
|
if !fm.f.indexedRange.headIndexed && !indexedBlocks.IsEmpty() {
|
||||||
lastIndexed, subLastIndexed uint64
|
indexedBlocks.SetAfterLast(indexedBlocks.Last()) // remove partially indexed last block
|
||||||
)
|
|
||||||
if !fm.f.indexedRange.headBlockIndexed {
|
|
||||||
subLastIndexed = 1
|
|
||||||
}
|
|
||||||
if fm.f.indexedRange.afterLastIndexedBlock-subLastIndexed > fm.f.indexedRange.firstIndexedBlock {
|
|
||||||
indexed, lastIndexed = true, fm.f.indexedRange.afterLastIndexedBlock-subLastIndexed-1
|
|
||||||
}
|
}
|
||||||
fm.syncCh <- SyncRange{
|
fm.syncCh <- SyncRange{
|
||||||
HeadNumber: fm.f.indexedView.headNumber,
|
HeadNumber: fm.f.indexedView.headNumber,
|
||||||
Valid: fm.valid,
|
ValidBlocks: fm.validBlocks,
|
||||||
FirstValid: fm.firstValid,
|
IndexedBlocks: indexedBlocks,
|
||||||
LastValid: fm.lastValid,
|
|
||||||
Indexed: indexed,
|
|
||||||
FirstIndexed: fm.f.indexedRange.firstIndexedBlock,
|
|
||||||
LastIndexed: lastIndexed,
|
|
||||||
}
|
}
|
||||||
fm.valid = indexed
|
fm.validBlocks = indexedBlocks
|
||||||
fm.firstValid = fm.f.indexedRange.firstIndexedBlock
|
|
||||||
fm.lastValid = lastIndexed
|
|
||||||
fm.syncCh = nil
|
fm.syncCh = nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -187,20 +173,10 @@ func (f *FilterMaps) updateMatchersValidRange() {
|
||||||
defer f.matchersLock.Unlock()
|
defer f.matchersLock.Unlock()
|
||||||
|
|
||||||
for fm := range f.matchers {
|
for fm := range f.matchers {
|
||||||
if !f.indexedRange.hasIndexedBlocks() {
|
if !f.indexedRange.initialized {
|
||||||
fm.valid = false
|
fm.validBlocks = common.Range[uint64]{}
|
||||||
}
|
|
||||||
if !fm.valid {
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if fm.firstValid < f.indexedRange.firstIndexedBlock {
|
fm.validBlocks = fm.validBlocks.Intersection(f.indexedRange.blocks)
|
||||||
fm.firstValid = f.indexedRange.firstIndexedBlock
|
|
||||||
}
|
|
||||||
if fm.lastValid >= f.indexedRange.afterLastIndexedBlock {
|
|
||||||
fm.lastValid = f.indexedRange.afterLastIndexedBlock - 1
|
|
||||||
}
|
|
||||||
if fm.firstValid > fm.lastValid {
|
|
||||||
fm.valid = false
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -356,8 +356,8 @@ func WriteFilterMapBaseRows(db ethdb.KeyValueWriter, mapRowIndex uint64, rows []
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func DeleteFilterMapRows(db ethdb.KeyValueRangeDeleter, firstMapRowIndex, afterLastMapRowIndex uint64) {
|
func DeleteFilterMapRows(db ethdb.KeyValueRangeDeleter, mapRows common.Range[uint64]) {
|
||||||
if err := db.DeleteRange(filterMapRowKey(firstMapRowIndex, false), filterMapRowKey(afterLastMapRowIndex, false)); err != nil {
|
if err := db.DeleteRange(filterMapRowKey(mapRows.First(), false), filterMapRowKey(mapRows.AfterLast(), false)); err != nil {
|
||||||
log.Crit("Failed to delete range of filter map rows", "err", err)
|
log.Crit("Failed to delete range of filter map rows", "err", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -396,8 +396,8 @@ func DeleteFilterMapLastBlock(db ethdb.KeyValueWriter, mapIndex uint32) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func DeleteFilterMapLastBlocks(db ethdb.KeyValueRangeDeleter, firstMapIndex, afterLastMapIndex uint32) {
|
func DeleteFilterMapLastBlocks(db ethdb.KeyValueRangeDeleter, maps common.Range[uint32]) {
|
||||||
if err := db.DeleteRange(filterMapLastBlockKey(firstMapIndex), filterMapLastBlockKey(afterLastMapIndex)); err != nil {
|
if err := db.DeleteRange(filterMapLastBlockKey(maps.First()), filterMapLastBlockKey(maps.AfterLast())); err != nil {
|
||||||
log.Crit("Failed to delete range of filter map last block pointers", "err", err)
|
log.Crit("Failed to delete range of filter map last block pointers", "err", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -433,8 +433,8 @@ func DeleteBlockLvPointer(db ethdb.KeyValueWriter, blockNumber uint64) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func DeleteBlockLvPointers(db ethdb.KeyValueRangeDeleter, firstBlockNumber, afterLastBlockNumber uint64) {
|
func DeleteBlockLvPointers(db ethdb.KeyValueRangeDeleter, blocks common.Range[uint64]) {
|
||||||
if err := db.DeleteRange(filterMapBlockLVKey(firstBlockNumber), filterMapBlockLVKey(afterLastBlockNumber)); err != nil {
|
if err := db.DeleteRange(filterMapBlockLVKey(blocks.First()), filterMapBlockLVKey(blocks.AfterLast())); err != nil {
|
||||||
log.Crit("Failed to delete range of block log value pointers", "err", err)
|
log.Crit("Failed to delete range of block log value pointers", "err", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -442,10 +442,11 @@ func DeleteBlockLvPointers(db ethdb.KeyValueRangeDeleter, firstBlockNumber, afte
|
||||||
// FilterMapsRange is a storage representation of the block range covered by the
|
// FilterMapsRange is a storage representation of the block range covered by the
|
||||||
// filter maps structure and the corresponting log value index range.
|
// filter maps structure and the corresponting log value index range.
|
||||||
type FilterMapsRange struct {
|
type FilterMapsRange struct {
|
||||||
HeadBlockIndexed bool
|
HeadIndexed bool
|
||||||
HeadBlockDelimiter uint64
|
HeadDelimiter uint64
|
||||||
FirstIndexedBlock, AfterLastIndexedBlock uint64
|
Blocks common.Range[uint64]
|
||||||
FirstRenderedMap, AfterLastRenderedMap, TailPartialEpoch uint32
|
Maps common.Range[uint32]
|
||||||
|
TailPartialEpoch uint32
|
||||||
}
|
}
|
||||||
|
|
||||||
// ReadFilterMapsRange retrieves the filter maps range data. Note that if the
|
// ReadFilterMapsRange retrieves the filter maps range data. Note that if the
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue