core/filtermaps: do not embed filterMapsRange

This commit is contained in:
Zsolt Felfoldi 2025-03-12 16:15:47 +01:00
parent 61b17c3bce
commit 25786522bf
5 changed files with 117 additions and 119 deletions

View file

@ -21,6 +21,7 @@ import (
"errors" "errors"
"fmt" "fmt"
"os" "os"
"slices"
"sync" "sync"
"time" "time"
@ -60,9 +61,9 @@ type FilterMaps struct {
// fields written by the indexer and read by matcher backend. Indexer can // fields written by the indexer and read by matcher backend. Indexer can
// read them without a lock and write them under indexLock write lock. // read them without a lock and write them under indexLock write lock.
// Matcher backend can read them under indexLock read lock. // Matcher backend can read them under indexLock read lock.
indexLock sync.RWMutex indexLock sync.RWMutex
filterMapsRange indexedRange filterMapsRange
indexedView *ChainView // always consistent with the log index indexedView *ChainView // always consistent with the log index
// also accessed by indexer and matcher backend but no locking needed. // also accessed by indexer and matcher backend but no locking needed.
filterMapCache *lru.Cache[uint32, filterMap] filterMapCache *lru.Cache[uint32, filterMap]
@ -132,15 +133,7 @@ type FilterRow []uint32
// Equal returns true if the given filter rows are equivalent. // Equal returns true if the given filter rows are equivalent.
func (a FilterRow) Equal(b FilterRow) bool { func (a FilterRow) Equal(b FilterRow) bool {
if len(a) != len(b) { return slices.Equal(a, b)
return false
}
for i, v := range a {
if b[i] != v {
return false
}
}
return true
} }
// filterMapsRange describes the rendered range of filter maps and the range // filterMapsRange describes the rendered range of filter maps and the range
@ -192,7 +185,7 @@ func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, params Params, h
unindexLimit: unindexLimit, unindexLimit: unindexLimit,
exportFileName: exportFileName, exportFileName: exportFileName,
Params: params, Params: params,
filterMapsRange: filterMapsRange{ indexedRange: filterMapsRange{
initialized: initialized, initialized: initialized,
headBlockIndexed: rs.HeadBlockIndexed, headBlockIndexed: rs.HeadBlockIndexed,
headBlockDelimiter: rs.HeadBlockDelimiter, headBlockDelimiter: rs.HeadBlockDelimiter,
@ -211,23 +204,26 @@ func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, params Params, h
renderSnapshots: lru.NewCache[uint64, *renderedMap](cachedRenderSnapshots), renderSnapshots: lru.NewCache[uint64, *renderedMap](cachedRenderSnapshots),
} }
f.targetView = initView f.targetView = initView
if f.initialized { if f.indexedRange.initialized {
f.indexedView = f.initChainView(f.targetView) f.indexedView = f.initChainView(f.targetView)
f.headBlockIndexed = f.afterLastIndexedBlock == f.indexedView.headNumber+1 f.indexedRange.headBlockIndexed = f.indexedRange.afterLastIndexedBlock == f.indexedView.headNumber+1
if !f.headBlockIndexed { if !f.indexedRange.headBlockIndexed {
f.headBlockDelimiter = 0 f.indexedRange.headBlockDelimiter = 0
} }
} }
if f.hasIndexedBlocks() { if f.indexedRange.hasIndexedBlocks() {
log.Info("Initialized log indexer", "first block", f.firstIndexedBlock, "last block", f.afterLastIndexedBlock-1, "first map", f.firstRenderedMap, "last map", f.afterLastRenderedMap-1, "head indexed", f.headBlockIndexed) log.Info("Initialized log indexer",
"first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
"first map", f.indexedRange.firstRenderedMap, "last map", f.indexedRange.afterLastRenderedMap-1,
"head indexed", f.indexedRange.headBlockIndexed)
} }
return f return f
} }
// Start starts the indexer. // Start starts the indexer.
func (f *FilterMaps) Start() { func (f *FilterMaps) Start() {
if !f.testDisableSnapshots && f.initialized && f.headBlockIndexed && if !f.testDisableSnapshots && f.indexedRange.initialized && f.indexedRange.headBlockIndexed &&
f.firstRenderedMap < f.afterLastRenderedMap { 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)
@ -249,7 +245,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.afterLastRenderedMap mapIndex := f.indexedRange.afterLastRenderedMap
for { for {
var ok bool var ok bool
mapIndex, ok = f.lastMapBoundaryBefore(mapIndex) mapIndex, ok = f.lastMapBoundaryBefore(mapIndex)
@ -272,7 +268,7 @@ func (f *FilterMaps) initChainView(chainView *ChainView) *ChainView {
// the database. The function returns true if everything was successfully removed. // the database. The function returns true if everything was successfully removed.
func (f *FilterMaps) reset() bool { func (f *FilterMaps) reset() bool {
f.indexLock.Lock() f.indexLock.Lock()
f.filterMapsRange = filterMapsRange{} f.indexedRange = filterMapsRange{}
f.indexedView = nil f.indexedView = nil
f.filterMapCache.Purge() f.filterMapCache.Purge()
f.renderSnapshots.Purge() f.renderSnapshots.Purge()
@ -369,7 +365,7 @@ func (f *FilterMaps) removeDbWithPrefix(prefix []byte, action string) bool {
// Note that this function assumes that the index write lock is being held. // Note that this function assumes that the index write lock is being held.
func (f *FilterMaps) setRange(batch ethdb.KeyValueWriter, newView *ChainView, newRange filterMapsRange) { func (f *FilterMaps) setRange(batch ethdb.KeyValueWriter, newView *ChainView, newRange filterMapsRange) {
f.indexedView = newView f.indexedView = newView
f.filterMapsRange = newRange f.indexedRange = newRange
f.updateMatchersValidRange() f.updateMatchersValidRange()
if newRange.initialized { if newRange.initialized {
rs := rawdb.FilterMapsRange{ rs := rawdb.FilterMapsRange{
@ -397,7 +393,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.firstRenderedMap || mapIndex >= f.afterLastRenderedMap { if mapIndex < f.indexedRange.firstRenderedMap || mapIndex >= f.indexedRange.afterLastRenderedMap {
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
@ -412,8 +408,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.firstIndexedBlock { if firstBlockNumber < f.indexedRange.firstIndexedBlock {
firstBlockNumber = f.firstIndexedBlock firstBlockNumber = f.indexedRange.firstIndexedBlock
} }
// 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 {
@ -567,8 +563,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.afterLastIndexedBlock && f.headBlockIndexed { if blockNumber >= f.indexedRange.afterLastIndexedBlock && f.indexedRange.headBlockIndexed {
return f.headBlockDelimiter, nil return f.indexedRange.headBlockDelimiter, nil
} }
if lvPointer, ok := f.lvPointerCache.Get(blockNumber); ok { if lvPointer, ok := f.lvPointerCache.Get(blockNumber); ok {
return lvPointer, nil return lvPointer, nil
@ -644,11 +640,13 @@ func (f *FilterMaps) deleteTailEpoch(epoch uint32) error {
} }
firstBlock++ firstBlock++
} }
fmr := f.filterMapsRange fmr := f.indexedRange
if f.firstRenderedMap == firstMap && f.afterLastRenderedMap > firstMap+f.mapsPerEpoch && f.tailPartialEpoch == 0 { if f.indexedRange.firstRenderedMap == firstMap &&
f.indexedRange.afterLastRenderedMap > firstMap+f.mapsPerEpoch &&
f.indexedRange.tailPartialEpoch == 0 {
fmr.firstRenderedMap = firstMap + f.mapsPerEpoch fmr.firstRenderedMap = firstMap + f.mapsPerEpoch
fmr.firstIndexedBlock = lastBlock + 1 fmr.firstIndexedBlock = lastBlock + 1
} else if f.firstRenderedMap == firstMap+f.mapsPerEpoch { } else if f.indexedRange.firstRenderedMap == 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")

View file

@ -41,7 +41,7 @@ func (f *FilterMaps) indexerLoop() {
log.Info("Started log indexer") log.Info("Started log indexer")
for !f.stop { for !f.stop {
if !f.initialized { if !f.indexedRange.initialized {
if err := f.init(); err != nil { if err := f.init(); err != nil {
log.Error("Error initializing log index", "error", err) log.Error("Error initializing log index", "error", err)
f.waitForEvent() f.waitForEvent()
@ -199,20 +199,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.afterLastIndexedBlock f.ptrHeadIndex = f.indexedRange.afterLastIndexedBlock
} }
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.hasIndexedBlocks() && f.afterLastIndexedBlock >= f.ptrHeadIndex && if f.indexedRange.hasIndexedBlocks() && f.indexedRange.afterLastIndexedBlock >= 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.firstIndexedBlock, "last block", f.afterLastIndexedBlock-1, "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
"processed", f.afterLastIndexedBlock-f.ptrHeadIndex, "processed", f.indexedRange.afterLastIndexedBlock-f.ptrHeadIndex,
"remaining", f.indexedView.headNumber+1-f.afterLastIndexedBlock, "remaining", f.indexedView.headNumber+1-f.indexedRange.afterLastIndexedBlock,
"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()
@ -223,8 +223,8 @@ func (f *FilterMaps) tryIndexHead() bool {
} }
if f.loggedHeadIndex { if f.loggedHeadIndex {
log.Info("Log index head rendering finished", log.Info("Log index head rendering finished",
"first block", f.firstIndexedBlock, "last block", f.afterLastIndexedBlock-1, "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
"processed", f.afterLastIndexedBlock-f.ptrHeadIndex, "processed", f.indexedRange.afterLastIndexedBlock-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
@ -237,7 +237,7 @@ func (f *FilterMaps) tryIndexHead() bool {
// rendered according to targetView and is suspended as soon as the targetView // rendered according to targetView and is suspended as soon as the targetView
// is changed. // is changed.
func (f *FilterMaps) tryIndexTail() bool { func (f *FilterMaps) tryIndexTail() bool {
for firstEpoch := f.firstRenderedMap >> f.logMapsPerEpoch; firstEpoch > 0 && f.needTailEpoch(firstEpoch-1); { for firstEpoch := f.indexedRange.firstRenderedMap >> f.logMapsPerEpoch; firstEpoch > 0 && f.needTailEpoch(firstEpoch-1); {
f.processEvents() f.processEvents()
if f.stop || !f.targetHeadIndexed() { if f.stop || !f.targetHeadIndexed() {
return false return false
@ -245,12 +245,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.firstRenderedMap { if tailRenderer != nil && tailRenderer.afterLastMap != f.indexedRange.firstRenderedMap {
tailRenderer = nil tailRenderer = nil
} }
if tailRenderer == nil { if tailRenderer == nil {
var err error var err error
tailRenderer, err = f.renderMapsBefore(f.firstRenderedMap) tailRenderer, err = f.renderMapsBefore(f.indexedRange.firstRenderedMap)
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
@ -263,7 +263,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.firstIndexedBlock - f.tailPartialBlocks() f.ptrTailIndex = f.indexedRange.firstIndexedBlock - f.tailPartialBlocks()
} }
done, err := tailRenderer.run(func() bool { done, err := tailRenderer.run(func() bool {
f.processEvents() f.processEvents()
@ -271,16 +271,16 @@ 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.firstIndexedBlock > ttb+tpb { if f.indexedRange.firstIndexedBlock > ttb+tpb {
remaining = f.firstIndexedBlock - ttb - tpb remaining = f.indexedRange.firstIndexedBlock - ttb - tpb
} }
if f.hasIndexedBlocks() && f.ptrTailIndex >= f.firstIndexedBlock && if f.indexedRange.hasIndexedBlocks() && f.ptrTailIndex >= f.indexedRange.firstIndexedBlock &&
(!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.firstIndexedBlock, "last block", f.afterLastIndexedBlock-1, "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
"processed", f.ptrTailIndex-f.firstIndexedBlock+tpb, "processed", f.ptrTailIndex-f.indexedRange.firstIndexedBlock+tpb,
"remaining", remaining, "remaining", remaining,
"next tail epoch percentage", f.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)))
f.loggedTailIndex = true f.loggedTailIndex = true
f.lastLogTailIndex = time.Now() f.lastLogTailIndex = time.Now()
@ -296,8 +296,8 @@ func (f *FilterMaps) tryIndexTail() bool {
} }
if f.loggedTailIndex { if f.loggedTailIndex {
log.Info("Log index tail rendering finished", log.Info("Log index tail rendering finished",
"first block", f.firstIndexedBlock, "last block", f.afterLastIndexedBlock-1, "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
"processed", f.ptrTailIndex-f.firstIndexedBlock, "processed", f.ptrTailIndex-f.indexedRange.firstIndexedBlock,
"elapsed", common.PrettyDuration(time.Since(f.startedTailIndexAt))) "elapsed", common.PrettyDuration(time.Since(f.startedTailIndexAt)))
f.loggedTailIndex = false f.loggedTailIndex = false
} }
@ -310,7 +310,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.firstRenderedMap - f.tailPartialEpoch) >> f.logMapsPerEpoch firstEpoch := (f.indexedRange.firstRenderedMap - f.indexedRange.tailPartialEpoch) >> f.logMapsPerEpoch
if f.needTailEpoch(firstEpoch) { if f.needTailEpoch(firstEpoch) {
break break
} }
@ -321,8 +321,8 @@ 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.firstRenderedMap - f.tailPartialEpoch f.ptrTailUnindexMap = f.indexedRange.firstRenderedMap - f.indexedRange.tailPartialEpoch
f.ptrTailUnindexBlock = f.firstIndexedBlock - f.tailPartialBlocks() f.ptrTailUnindexBlock = f.indexedRange.firstIndexedBlock - 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)
@ -331,9 +331,9 @@ func (f *FilterMaps) tryUnindexTail() bool {
} }
if f.startedTailUnindex { if f.startedTailUnindex {
log.Info("Log index tail unindexing finished", log.Info("Log index tail unindexing finished",
"first block", f.firstIndexedBlock, "last block", f.afterLastIndexedBlock-1, "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1,
"removed maps", f.firstRenderedMap-f.ptrTailUnindexMap, "removed maps", f.indexedRange.firstRenderedMap-f.ptrTailUnindexMap,
"removed blocks", f.firstIndexedBlock-f.tailPartialBlocks()-f.ptrTailUnindexBlock, "removed blocks", f.indexedRange.firstIndexedBlock-f.tailPartialBlocks()-f.ptrTailUnindexBlock,
"elapsed", common.PrettyDuration(time.Since(f.startedTailUnindexAt))) "elapsed", common.PrettyDuration(time.Since(f.startedTailUnindexAt)))
f.startedTailUnindex = false f.startedTailUnindex = false
} }
@ -343,7 +343,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.firstRenderedMap >> f.logMapsPerEpoch firstEpoch := f.indexedRange.firstRenderedMap >> f.logMapsPerEpoch
if epoch > firstEpoch { if epoch > firstEpoch {
return true return true
} }
@ -351,7 +351,7 @@ func (f *FilterMaps) needTailEpoch(epoch uint32) bool {
return false return false
} }
tailTarget := f.tailTargetBlock() tailTarget := f.tailTargetBlock()
if tailTarget < f.firstIndexedBlock { if tailTarget < f.indexedRange.firstIndexedBlock {
return true return true
} }
tailLvIndex, err := f.getBlockLvPointer(tailTarget) tailLvIndex, err := f.getBlockLvPointer(tailTarget)
@ -374,18 +374,18 @@ func (f *FilterMaps) tailTargetBlock() uint64 {
// tailPartialBlocks returns the number of rendered blocks in the partially // tailPartialBlocks returns the number of rendered blocks in the partially
// rendered next tail epoch. // rendered next tail epoch.
func (f *FilterMaps) tailPartialBlocks() uint64 { func (f *FilterMaps) tailPartialBlocks() uint64 {
if f.tailPartialEpoch == 0 { if f.indexedRange.tailPartialEpoch == 0 {
return 0 return 0
} }
end, _, err := f.getLastBlockOfMap(f.firstRenderedMap - f.mapsPerEpoch + f.tailPartialEpoch - 1) end, _, err := f.getLastBlockOfMap(f.indexedRange.firstRenderedMap - f.mapsPerEpoch + f.indexedRange.tailPartialEpoch - 1)
if err != nil { if err != nil {
log.Error("Error fetching last block of map", "mapIndex", f.firstRenderedMap-f.mapsPerEpoch+f.tailPartialEpoch-1, "error", err) log.Error("Error fetching last block of map", "mapIndex", f.indexedRange.firstRenderedMap-f.mapsPerEpoch+f.indexedRange.tailPartialEpoch-1, "error", err)
} }
var start uint64 var start uint64
if f.firstRenderedMap-f.mapsPerEpoch > 0 { if f.indexedRange.firstRenderedMap-f.mapsPerEpoch > 0 {
start, _, err = f.getLastBlockOfMap(f.firstRenderedMap - f.mapsPerEpoch - 1) start, _, err = f.getLastBlockOfMap(f.indexedRange.firstRenderedMap - f.mapsPerEpoch - 1)
if err != nil { if err != nil {
log.Error("Error fetching last block of map", "mapIndex", f.firstRenderedMap-f.mapsPerEpoch-1, "error", err) log.Error("Error fetching last block of map", "mapIndex", f.indexedRange.firstRenderedMap-f.mapsPerEpoch-1, "error", err)
} }
} }
return end - start return end - start
@ -394,5 +394,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.headBlockIndexed return equalViews(f.targetView, f.indexedView) && f.indexedRange.headBlockIndexed
} }

View file

@ -101,12 +101,12 @@ func TestIndexerRandomRange(t *testing.T) {
checkSnapshot = false checkSnapshot = false
} }
if noHistory { if noHistory {
if ts.fm.initialized { if ts.fm.indexedRange.initialized {
t.Fatalf("filterMapsRange initialized while indexing is disabled") t.Fatalf("filterMapsRange initialized while indexing is disabled")
} }
continue continue
} }
if !ts.fm.initialized { if !ts.fm.indexedRange.initialized {
t.Fatalf("filterMapsRange not initialized while indexing is enabled") t.Fatalf("filterMapsRange not initialized while indexing is enabled")
} }
var tailBlock uint64 var tailBlock uint64
@ -124,14 +124,14 @@ func TestIndexerRandomRange(t *testing.T) {
// (expTailBlock-1)*lvPerBlock >= tailLvPtr // (expTailBlock-1)*lvPerBlock >= tailLvPtr
expTailBlock = (tailLvPtr + lvPerBlock*2 - 1) / lvPerBlock expTailBlock = (tailLvPtr + lvPerBlock*2 - 1) / lvPerBlock
} }
if ts.fm.afterLastIndexedBlock != uint64(head+1) { if ts.fm.indexedRange.afterLastIndexedBlock != uint64(head+1) {
ts.t.Fatalf("Invalid index head (expected #%d, got #%d)", head, ts.fm.afterLastIndexedBlock-1) ts.t.Fatalf("Invalid index head (expected #%d, got #%d)", head, ts.fm.indexedRange.afterLastIndexedBlock-1)
} }
if ts.fm.headBlockDelimiter != uint64(head)*lvPerBlock { if ts.fm.indexedRange.headBlockDelimiter != uint64(head)*lvPerBlock {
ts.t.Fatalf("Invalid index head delimiter pointer (expected %d, got %d)", uint64(head)*lvPerBlock, ts.fm.headBlockDelimiter) ts.t.Fatalf("Invalid index head delimiter pointer (expected %d, got %d)", uint64(head)*lvPerBlock, ts.fm.indexedRange.headBlockDelimiter)
} }
if ts.fm.firstIndexedBlock != expTailBlock { if ts.fm.indexedRange.firstIndexedBlock != expTailBlock {
ts.t.Fatalf("Invalid index tail block (expected #%d, got #%d)", expTailBlock, ts.fm.firstIndexedBlock) ts.t.Fatalf("Invalid index tail block (expected #%d, got #%d)", expTailBlock, ts.fm.indexedRange.firstIndexedBlock)
} }
} }
} }

View file

@ -139,7 +139,7 @@ func (f *FilterMaps) renderMapsFromMapBoundary(firstMap, afterLastMap uint32, st
func (f *FilterMaps) lastCanonicalSnapshotBefore(afterLastMap uint32) *renderedMap { func (f *FilterMaps) lastCanonicalSnapshotBefore(afterLastMap 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.afterLastIndexedBlock && if cp, _ := f.renderSnapshots.Get(blockNumber); cp != nil && blockNumber < f.indexedRange.afterLastIndexedBlock &&
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 < afterLastMap && (best == nil || blockNumber > best.lastBlock) {
best = cp best = cp
@ -155,7 +155,7 @@ func (f *FilterMaps) lastCanonicalSnapshotBefore(afterLastMap uint32) *renderedM
// 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(afterLastMap uint32) (nextMap uint32, startBlock, startLvPtr uint64, err error) {
if !f.initialized { if !f.indexedRange.initialized {
return 0, 0, 0, nil return 0, 0, 0, nil
} }
mapIndex := afterLastMap mapIndex := afterLastMap
@ -184,18 +184,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.initialized || f.afterLastRenderedMap == 0 { if !f.indexedRange.initialized || f.indexedRange.afterLastRenderedMap == 0 {
return 0, false return 0, false
} }
if mapIndex > f.afterLastRenderedMap { if mapIndex > f.indexedRange.afterLastRenderedMap {
mapIndex = f.afterLastRenderedMap mapIndex = f.indexedRange.afterLastRenderedMap
} }
if mapIndex > f.firstRenderedMap { if mapIndex > f.indexedRange.firstRenderedMap {
return mapIndex - 1, true return mapIndex - 1, true
} }
if mapIndex+f.mapsPerEpoch > f.firstRenderedMap { if mapIndex+f.mapsPerEpoch > f.indexedRange.firstRenderedMap {
if mapIndex > f.firstRenderedMap-f.mapsPerEpoch+f.tailPartialEpoch { if mapIndex > f.indexedRange.firstRenderedMap-f.mapsPerEpoch+f.indexedRange.tailPartialEpoch {
mapIndex = f.firstRenderedMap - f.mapsPerEpoch + f.tailPartialEpoch mapIndex = f.indexedRange.firstRenderedMap - f.mapsPerEpoch + f.indexedRange.tailPartialEpoch
} }
} else { } else {
mapIndex = (mapIndex >> f.logMapsPerEpoch) << f.logMapsPerEpoch mapIndex = (mapIndex >> f.logMapsPerEpoch) << f.logMapsPerEpoch
@ -214,19 +214,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.afterLastRenderedMap - 1) fm, err := f.getFilterMap(f.indexedRange.afterLastRenderedMap - 1)
if err != nil { if err != nil {
return fmt.Errorf("failed to load head snapshot map %d: %v", f.afterLastRenderedMap-1, err) return fmt.Errorf("failed to load head snapshot map %d: %v", f.indexedRange.afterLastRenderedMap-1, err)
} }
lastBlock, _, err := f.getLastBlockOfMap(f.afterLastRenderedMap - 1) lastBlock, _, err := f.getLastBlockOfMap(f.indexedRange.afterLastRenderedMap - 1)
if err != nil { if err != nil {
return fmt.Errorf("failed to retrieve last block of head snapshot map %d: %v", f.afterLastRenderedMap-1, err) return fmt.Errorf("failed to retrieve last block of head snapshot map %d: %v", f.indexedRange.afterLastRenderedMap-1, err)
} }
var firstBlock uint64 var firstBlock uint64
if f.afterLastRenderedMap > 1 { if f.indexedRange.afterLastRenderedMap > 1 {
prevLastBlock, _, err := f.getLastBlockOfMap(f.afterLastRenderedMap - 2) prevLastBlock, _, err := f.getLastBlockOfMap(f.indexedRange.afterLastRenderedMap - 2)
if err != nil { if err != nil {
return fmt.Errorf("failed to retrieve last block of map %d before head snapshot: %v", f.afterLastRenderedMap-2, err) return fmt.Errorf("failed to retrieve last block of map %d before head snapshot: %v", f.indexedRange.afterLastRenderedMap-2, err)
} }
firstBlock = prevLastBlock + 1 firstBlock = prevLastBlock + 1
} }
@ -237,14 +237,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.afterLastIndexedBlock-1, &renderedMap{ f.renderSnapshots.Add(f.indexedRange.afterLastIndexedBlock-1, &renderedMap{
filterMap: fm, filterMap: fm,
mapIndex: f.afterLastRenderedMap - 1, mapIndex: f.indexedRange.afterLastRenderedMap - 1,
lastBlock: f.afterLastIndexedBlock - 1, lastBlock: f.indexedRange.afterLastIndexedBlock - 1,
lastBlockId: f.indexedView.getBlockId(f.afterLastIndexedBlock - 1), lastBlockId: f.indexedView.getBlockId(f.indexedRange.afterLastIndexedBlock - 1),
blockLvPtrs: lvPtrs, blockLvPtrs: lvPtrs,
finished: true, finished: true,
headDelimiter: f.headBlockDelimiter, headDelimiter: f.indexedRange.headBlockDelimiter,
}) })
return nil return nil
} }
@ -342,7 +342,7 @@ func (r *mapRenderer) renderCurrentMap(stopCb func() bool) (bool, error) {
if err := r.iterator.next(); err != nil { if err := r.iterator.next(); err != nil {
return false, fmt.Errorf("failed to advance log iterator at %d while rendering map %d: %v", r.iterator.lvIndex, r.currentMap.mapIndex, err) return false, fmt.Errorf("failed to advance log iterator at %d while rendering map %d: %v", r.iterator.lvIndex, r.currentMap.mapIndex, err)
} }
if !r.f.testDisableSnapshots && r.afterLastMap >= r.f.afterLastRenderedMap && if !r.f.testDisableSnapshots && r.afterLastMap >= r.f.indexedRange.afterLastRenderedMap &&
(r.iterator.delimiter || r.iterator.finished) { (r.iterator.delimiter || r.iterator.finished) {
r.makeSnapshot() r.makeSnapshot()
} }
@ -364,7 +364,7 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error {
r.f.indexLock.Lock() r.f.indexLock.Lock()
defer r.f.indexLock.Unlock() defer r.f.indexLock.Unlock()
oldRange := r.f.filterMapsRange oldRange := r.f.indexedRange
tempRange, err := r.getTempRange() tempRange, err := r.getTempRange()
if err != nil { if err != nil {
return fmt.Errorf("failed to get temporary rendered range: %v", err) return fmt.Errorf("failed to get temporary rendered range: %v", err)
@ -476,11 +476,11 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error {
// performance so instead safety is ensured by first reverting the valid map // performance so instead safety is ensured by first reverting the valid map
// 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.filterMapsRange tempRange := r.f.indexedRange
if err := tempRange.addRenderedRange(r.firstFinished, r.firstFinished, r.afterLastMap, r.f.mapsPerEpoch); err != nil { if err := tempRange.addRenderedRange(r.firstFinished, r.firstFinished, r.afterLastMap, 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.firstRenderedMap { if tempRange.firstRenderedMap != r.f.indexedRange.firstRenderedMap {
// first rendered map changed; update first indexed block // first rendered map changed; update first indexed block
if tempRange.firstRenderedMap > 0 { if tempRange.firstRenderedMap > 0 {
lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.firstRenderedMap - 1) lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.firstRenderedMap - 1)
@ -492,7 +492,7 @@ func (r *mapRenderer) getTempRange() (filterMapsRange, error) {
tempRange.firstIndexedBlock = 0 tempRange.firstIndexedBlock = 0
} }
} }
if tempRange.afterLastRenderedMap != r.f.afterLastRenderedMap { if tempRange.afterLastRenderedMap != r.f.indexedRange.afterLastRenderedMap {
// first rendered map changed; update first indexed block // first rendered map changed; update first indexed block
if tempRange.afterLastRenderedMap > 0 { if tempRange.afterLastRenderedMap > 0 {
lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.afterLastRenderedMap - 1) lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.afterLastRenderedMap - 1)
@ -512,11 +512,11 @@ func (r *mapRenderer) getTempRange() (filterMapsRange, error) {
// rendered maps. // rendered maps.
func (r *mapRenderer) getUpdatedRange() (filterMapsRange, error) { func (r *mapRenderer) getUpdatedRange() (filterMapsRange, error) {
// update filterMapsRange // update filterMapsRange
newRange := r.f.filterMapsRange newRange := r.f.indexedRange
if err := newRange.addRenderedRange(r.firstFinished, r.afterLastFinished, r.afterLastMap, r.f.mapsPerEpoch); err != nil { if err := newRange.addRenderedRange(r.firstFinished, r.afterLastFinished, r.afterLastMap, 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.firstRenderedMap { if newRange.firstRenderedMap != r.f.indexedRange.firstRenderedMap {
// first rendered map changed; update first indexed block // first rendered map changed; update first indexed block
if newRange.firstRenderedMap > 0 { if newRange.firstRenderedMap > 0 {
lastBlock, _, err := r.f.getLastBlockOfMap(newRange.firstRenderedMap - 1) lastBlock, _, err := r.f.getLastBlockOfMap(newRange.firstRenderedMap - 1)
@ -625,12 +625,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.firstIndexedBlock || blockNumber >= f.afterLastIndexedBlock { if blockNumber < f.indexedRange.firstIndexedBlock || blockNumber >= f.indexedRange.afterLastIndexedBlock {
return nil, errUnindexedRange return nil, errUnindexedRange
} }
var lvIndex uint64 var lvIndex uint64
if f.headBlockIndexed && blockNumber+1 == f.afterLastIndexedBlock { if f.indexedRange.headBlockIndexed && blockNumber+1 == f.indexedRange.afterLastIndexedBlock {
lvIndex = f.headBlockDelimiter lvIndex = f.indexedRange.headBlockDelimiter
} else { } else {
var err error var err error
lvIndex, err = f.getBlockLvPointer(blockNumber + 1) lvIndex, err = f.getBlockLvPointer(blockNumber + 1)

View file

@ -45,9 +45,9 @@ func (f *FilterMaps) NewMatcherBackend() *FilterMapsMatcherBackend {
fm := &FilterMapsMatcherBackend{ fm := &FilterMapsMatcherBackend{
f: f, f: f,
valid: f.initialized && f.afterLastIndexedBlock > f.firstIndexedBlock, valid: f.indexedRange.initialized && f.indexedRange.afterLastIndexedBlock > f.indexedRange.firstIndexedBlock,
firstValid: f.firstIndexedBlock, firstValid: f.indexedRange.firstIndexedBlock,
lastValid: f.afterLastIndexedBlock - 1, lastValid: f.indexedRange.afterLastIndexedBlock - 1,
} }
f.matchers[fm] = struct{}{} f.matchers[fm] = struct{}{}
return fm return fm
@ -126,11 +126,11 @@ func (fm *FilterMapsMatcherBackend) synced() {
indexed bool indexed bool
lastIndexed, subLastIndexed uint64 lastIndexed, subLastIndexed uint64
) )
if !fm.f.headBlockIndexed { if !fm.f.indexedRange.headBlockIndexed {
subLastIndexed = 1 subLastIndexed = 1
} }
if fm.f.afterLastIndexedBlock-subLastIndexed > fm.f.firstIndexedBlock { if fm.f.indexedRange.afterLastIndexedBlock-subLastIndexed > fm.f.indexedRange.firstIndexedBlock {
indexed, lastIndexed = true, fm.f.afterLastIndexedBlock-subLastIndexed-1 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,
@ -138,11 +138,11 @@ func (fm *FilterMapsMatcherBackend) synced() {
FirstValid: fm.firstValid, FirstValid: fm.firstValid,
LastValid: fm.lastValid, LastValid: fm.lastValid,
Indexed: indexed, Indexed: indexed,
FirstIndexed: fm.f.firstIndexedBlock, FirstIndexed: fm.f.indexedRange.firstIndexedBlock,
LastIndexed: lastIndexed, LastIndexed: lastIndexed,
} }
fm.valid = indexed fm.valid = indexed
fm.firstValid = fm.f.firstIndexedBlock fm.firstValid = fm.f.indexedRange.firstIndexedBlock
fm.lastValid = lastIndexed fm.lastValid = lastIndexed
fm.syncCh = nil fm.syncCh = nil
} }
@ -189,17 +189,17 @@ func (f *FilterMaps) updateMatchersValidRange() {
defer f.matchersLock.Unlock() defer f.matchersLock.Unlock()
for fm := range f.matchers { for fm := range f.matchers {
if !f.hasIndexedBlocks() { if !f.indexedRange.hasIndexedBlocks() {
fm.valid = false fm.valid = false
} }
if !fm.valid { if !fm.valid {
continue continue
} }
if fm.firstValid < f.firstIndexedBlock { if fm.firstValid < f.indexedRange.firstIndexedBlock {
fm.firstValid = f.firstIndexedBlock fm.firstValid = f.indexedRange.firstIndexedBlock
} }
if fm.lastValid >= f.afterLastIndexedBlock { if fm.lastValid >= f.indexedRange.afterLastIndexedBlock {
fm.lastValid = f.afterLastIndexedBlock - 1 fm.lastValid = f.indexedRange.afterLastIndexedBlock - 1
} }
if fm.firstValid > fm.lastValid { if fm.firstValid > fm.lastValid {
fm.valid = false fm.valid = false