From 6a6deb1bd892d4d6928ec0a9277ffd8c73d3595f Mon Sep 17 00:00:00 2001 From: Zsolt Felfoldi Date: Fri, 14 Mar 2025 16:08:22 +0100 Subject: [PATCH] common, core/filtermaps: use Range type --- common/types.go | 82 ++++++++++++ core/filtermaps/filtermaps.go | 98 +++++++------- core/filtermaps/indexer.go | 68 +++++----- core/filtermaps/indexer_test.go | 12 +- core/filtermaps/map_renderer.go | 207 ++++++++++++++--------------- core/filtermaps/matcher.go | 8 +- core/filtermaps/matcher_backend.go | 56 +++----- core/rawdb/accessors_indexes.go | 21 +-- 8 files changed, 299 insertions(+), 253 deletions(-) diff --git a/common/types.go b/common/types.go index fdb25f1b34..89a559ef19 100644 --- a/common/types.go +++ b/common/types.go @@ -23,6 +23,7 @@ import ( "encoding/json" "errors" "fmt" + "io" "math/big" "math/rand" "reflect" @@ -30,6 +31,7 @@ import ( "strings" "github.com/ethereum/go-ethereum/common/hexutil" + "github.com/ethereum/go-ethereum/rlp" "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)) } + +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 +} diff --git a/core/filtermaps/filtermaps.go b/core/filtermaps/filtermaps.go index 9e6136ca27..528d24c62e 100644 --- a/core/filtermaps/filtermaps.go +++ b/core/filtermaps/filtermaps.go @@ -145,24 +145,23 @@ func (a FilterRow) Equal(b FilterRow) bool { // filterMapsRange describes the rendered range of filter maps and the range // of fully rendered blocks. type filterMapsRange struct { - initialized bool - headBlockIndexed bool - headBlockDelimiter uint64 // zero if afterLastIndexedBlock != targetBlockNumber - // if initialized then all maps are rendered between firstRenderedMap and - // afterLastRenderedMap-1 - firstRenderedMap, afterLastRenderedMap uint32 + initialized bool + headIndexed bool + headDelimiter uint64 // zero if headIndexed is false + // if initialized then all maps are rendered in the maps range + maps common.Range[uint32] // if tailPartialEpoch > 0 then maps between firstRenderedMap-mapsPerEpoch and // firstRenderedMap-mapsPerEpoch+tailPartialEpoch-1 are rendered tailPartialEpoch uint32 - // if initialized then all log values belonging to blocks between - // firstIndexedBlock and afterLastIndexedBlock are fully rendered - // blockLvPointers are available between firstIndexedBlock and afterLastIndexedBlock-1 - firstIndexedBlock, afterLastIndexedBlock uint64 + // if initialized then all log values in the blocks range are fully + // rendered + // blockLvPointers are available in the blocks range + blocks common.Range[uint64] } // hasIndexedBlocks returns true if the range has at least one fully indexed block. 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 @@ -200,14 +199,12 @@ func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, historyCutoff, f exportFileName: config.ExportFileName, Params: params, indexedRange: filterMapsRange{ - initialized: initialized, - headBlockIndexed: rs.HeadBlockIndexed, - headBlockDelimiter: rs.HeadBlockDelimiter, - firstIndexedBlock: rs.FirstIndexedBlock, - afterLastIndexedBlock: rs.AfterLastIndexedBlock, - firstRenderedMap: rs.FirstRenderedMap, - afterLastRenderedMap: rs.AfterLastRenderedMap, - tailPartialEpoch: rs.TailPartialEpoch, + initialized: initialized, + headIndexed: rs.HeadIndexed, + headDelimiter: rs.HeadDelimiter, + blocks: rs.Blocks, + maps: rs.Maps, + tailPartialEpoch: rs.TailPartialEpoch, }, matcherSyncCh: make(chan *FilterMapsMatcherBackend), matchers: make(map[*FilterMapsMatcherBackend]struct{}), @@ -222,24 +219,23 @@ func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, historyCutoff, f f.targetView = initView if f.indexedRange.initialized { f.indexedView = f.initChainView(f.targetView) - f.indexedRange.headBlockIndexed = f.indexedRange.afterLastIndexedBlock == f.indexedView.headNumber+1 - if !f.indexedRange.headBlockIndexed { - f.indexedRange.headBlockDelimiter = 0 + f.indexedRange.headIndexed = f.indexedRange.blocks.AfterLast() == f.indexedView.headNumber+1 + if !f.indexedRange.headIndexed { + f.indexedRange.headDelimiter = 0 } } if f.indexedRange.hasIndexedBlocks() { 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) + "first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(), + "first map", f.indexedRange.maps.First(), "last map", f.indexedRange.maps.Last(), + "head indexed", f.indexedRange.headIndexed) } return f } // Start starts the indexer. func (f *FilterMaps) Start() { - if !f.testDisableSnapshots && f.indexedRange.initialized && f.indexedRange.headBlockIndexed && - f.indexedRange.firstRenderedMap < f.indexedRange.afterLastRenderedMap { + if !f.testDisableSnapshots && f.indexedRange.hasIndexedBlocks() && f.indexedRange.headIndexed { // previous target head rendered; load last map as snapshot if err := f.loadHeadSnapshot(); err != nil { 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 // the latest maps are not consistent with targetView. func (f *FilterMaps) initChainView(chainView *ChainView) *ChainView { - mapIndex := f.indexedRange.afterLastRenderedMap + mapIndex := f.indexedRange.maps.AfterLast() for { var ok bool mapIndex, ok = f.lastMapBoundaryBefore(mapIndex) @@ -331,10 +327,8 @@ func (f *FilterMaps) init() error { } if bestLen > 0 { cp := checkpoints[bestIdx][bestLen-1] - fmr.firstIndexedBlock = cp.BlockNumber + 1 - fmr.afterLastIndexedBlock = cp.BlockNumber + 1 - fmr.firstRenderedMap = uint32(bestLen) << f.logMapsPerEpoch - fmr.afterLastRenderedMap = uint32(bestLen) << f.logMapsPerEpoch + fmr.blocks = common.NewRange[uint64](cp.BlockNumber+1, 0) + fmr.maps = common.NewRange[uint32](uint32(bestLen)<> f.logValuesPerMap) - if mapIndex < f.indexedRange.firstRenderedMap || mapIndex >= f.indexedRange.afterLastRenderedMap { + if !f.indexedRange.maps.Includes(mapIndex) { return nil, nil } // 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) } } - if firstBlockNumber < f.indexedRange.firstIndexedBlock { - firstBlockNumber = f.indexedRange.firstIndexedBlock + if firstBlockNumber < f.indexedRange.blocks.First() { + firstBlockNumber = f.indexedRange.blocks.First() } // find block with binary search based on block to log value index pointers 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 // called from outside the indexerLoop goroutine. func (f *FilterMaps) getBlockLvPointer(blockNumber uint64) (uint64, error) { - if blockNumber >= f.indexedRange.afterLastIndexedBlock && f.indexedRange.headBlockIndexed { - return f.indexedRange.headBlockDelimiter, nil + if blockNumber >= f.indexedRange.blocks.AfterLast() && f.indexedRange.headIndexed { + return f.indexedRange.headDelimiter, nil } if lvPointer, ok := f.lvPointerCache.Get(blockNumber); ok { return lvPointer, nil @@ -662,26 +654,28 @@ func (f *FilterMaps) deleteTailEpoch(epoch uint32) error { firstBlock++ } fmr := f.indexedRange - if f.indexedRange.firstRenderedMap == firstMap && - f.indexedRange.afterLastRenderedMap > firstMap+f.mapsPerEpoch && + if f.indexedRange.maps.First() == firstMap && + f.indexedRange.maps.AfterLast() > firstMap+f.mapsPerEpoch && f.indexedRange.tailPartialEpoch == 0 { - fmr.firstRenderedMap = firstMap + f.mapsPerEpoch - fmr.firstIndexedBlock = lastBlock + 1 - } else if f.indexedRange.firstRenderedMap == firstMap+f.mapsPerEpoch { + fmr.maps.SetFirst(firstMap + f.mapsPerEpoch) + fmr.blocks.SetFirst(lastBlock + 1) + } else if f.indexedRange.maps.First() == firstMap+f.mapsPerEpoch { fmr.tailPartialEpoch = 0 } else { return errors.New("invalid tail epoch number") } 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++ { 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++ { 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++ { f.lvPointerCache.Remove(blockNumber) } diff --git a/core/filtermaps/indexer.go b/core/filtermaps/indexer.go index 82eda19136..addda03df4 100644 --- a/core/filtermaps/indexer.go +++ b/core/filtermaps/indexer.go @@ -193,20 +193,20 @@ func (f *FilterMaps) tryIndexHead() bool { f.lastLogHeadIndex = time.Now() f.startedHeadIndexAt = f.lastLogHeadIndex f.startedHeadIndex = true - f.ptrHeadIndex = f.indexedRange.afterLastIndexedBlock + f.ptrHeadIndex = f.indexedRange.blocks.AfterLast() } if _, err := headRenderer.run(func() bool { f.processEvents() return f.stop }, func() { 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) || time.Since(f.lastLogHeadIndex) > logFrequency) { log.Info("Log index head rendering in progress", - "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1, - "processed", f.indexedRange.afterLastIndexedBlock-f.ptrHeadIndex, - "remaining", f.indexedView.headNumber+1-f.indexedRange.afterLastIndexedBlock, + "first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(), + "processed", f.indexedRange.blocks.AfterLast()-f.ptrHeadIndex, + "remaining", f.indexedView.headNumber-f.indexedRange.blocks.Last(), "elapsed", common.PrettyDuration(time.Since(f.startedHeadIndexAt))) f.loggedHeadIndex = true f.lastLogHeadIndex = time.Now() @@ -215,10 +215,10 @@ func (f *FilterMaps) tryIndexHead() bool { log.Error("Log index head rendering failed", "error", err) return false } - if f.loggedHeadIndex { + if f.loggedHeadIndex && f.indexedRange.hasIndexedBlocks() { log.Info("Log index head rendering finished", - "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1, - "processed", f.indexedRange.afterLastIndexedBlock-f.ptrHeadIndex, + "first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(), + "processed", f.indexedRange.blocks.AfterLast()-f.ptrHeadIndex, "elapsed", common.PrettyDuration(time.Since(f.startedHeadIndexAt))) } f.loggedHeadIndex, f.startedHeadIndex = false, false @@ -232,7 +232,7 @@ func (f *FilterMaps) tryIndexHead() bool { // is changed. func (f *FilterMaps) tryIndexTail() bool { for { - firstEpoch := f.indexedRange.firstRenderedMap >> f.logMapsPerEpoch + firstEpoch := f.indexedRange.maps.First() >> f.logMapsPerEpoch if firstEpoch == 0 || !f.needTailEpoch(firstEpoch-1) { break } @@ -243,12 +243,12 @@ func (f *FilterMaps) tryIndexTail() bool { // resume process if tail rendering was interrupted because of head rendering tailRenderer := f.tailRenderer f.tailRenderer = nil - if tailRenderer != nil && tailRenderer.afterLastMap != f.indexedRange.firstRenderedMap { + if tailRenderer != nil && tailRenderer.renderBefore != f.indexedRange.maps.First() { tailRenderer = nil } if tailRenderer == nil { var err error - tailRenderer, err = f.renderMapsBefore(f.indexedRange.firstRenderedMap) + tailRenderer, err = f.renderMapsBefore(f.indexedRange.maps.First()) if err != nil { log.Error("Error creating log index tail renderer", "error", err) return false @@ -261,7 +261,7 @@ func (f *FilterMaps) tryIndexTail() bool { f.lastLogTailIndex = time.Now() f.startedTailIndexAt = f.lastLogTailIndex f.startedTailIndex = true - f.ptrTailIndex = f.indexedRange.firstIndexedBlock - f.tailPartialBlocks() + f.ptrTailIndex = f.indexedRange.blocks.First() - f.tailPartialBlocks() } done, err := tailRenderer.run(func() bool { f.processEvents() @@ -269,14 +269,14 @@ func (f *FilterMaps) tryIndexTail() bool { }, func() { tpb, ttb := f.tailPartialBlocks(), f.tailTargetBlock() remaining := uint64(1) - if f.indexedRange.firstIndexedBlock > ttb+tpb { - remaining = f.indexedRange.firstIndexedBlock - ttb - tpb + if f.indexedRange.blocks.First() > 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) { log.Info("Log index tail rendering in progress", - "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1, - "processed", f.ptrTailIndex-f.indexedRange.firstIndexedBlock+tpb, + "first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(), + "processed", f.ptrTailIndex-f.indexedRange.blocks.First()+tpb, "remaining", remaining, "next tail epoch percentage", f.indexedRange.tailPartialEpoch*100/f.mapsPerEpoch, "elapsed", common.PrettyDuration(time.Since(f.startedTailIndexAt))) @@ -293,10 +293,10 @@ func (f *FilterMaps) tryIndexTail() bool { return false } } - if f.loggedTailIndex { + if f.loggedTailIndex && f.indexedRange.hasIndexedBlocks() { log.Info("Log index tail rendering finished", - "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1, - "processed", f.ptrTailIndex-f.indexedRange.firstIndexedBlock, + "first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(), + "processed", f.ptrTailIndex-f.indexedRange.blocks.First(), "elapsed", common.PrettyDuration(time.Since(f.startedTailIndexAt))) f.loggedTailIndex = false } @@ -309,7 +309,7 @@ func (f *FilterMaps) tryIndexTail() bool { // data from the database and is also called while running head indexing. func (f *FilterMaps) tryUnindexTail() bool { for { - firstEpoch := (f.indexedRange.firstRenderedMap - f.indexedRange.tailPartialEpoch) >> f.logMapsPerEpoch + firstEpoch := (f.indexedRange.maps.First() - f.indexedRange.tailPartialEpoch) >> f.logMapsPerEpoch if f.needTailEpoch(firstEpoch) { break } @@ -320,19 +320,19 @@ func (f *FilterMaps) tryUnindexTail() bool { if !f.startedTailUnindex { f.startedTailUnindexAt = time.Now() f.startedTailUnindex = true - f.ptrTailUnindexMap = f.indexedRange.firstRenderedMap - f.indexedRange.tailPartialEpoch - f.ptrTailUnindexBlock = f.indexedRange.firstIndexedBlock - f.tailPartialBlocks() + f.ptrTailUnindexMap = f.indexedRange.maps.First() - f.indexedRange.tailPartialEpoch + f.ptrTailUnindexBlock = f.indexedRange.blocks.First() - f.tailPartialBlocks() } if err := f.deleteTailEpoch(firstEpoch); err != nil { log.Error("Log index tail epoch unindexing failed", "error", err) return false } } - if f.startedTailUnindex { + if f.startedTailUnindex && f.indexedRange.hasIndexedBlocks() { log.Info("Log index tail unindexing finished", - "first block", f.indexedRange.firstIndexedBlock, "last block", f.indexedRange.afterLastIndexedBlock-1, - "removed maps", f.indexedRange.firstRenderedMap-f.ptrTailUnindexMap, - "removed blocks", f.indexedRange.firstIndexedBlock-f.tailPartialBlocks()-f.ptrTailUnindexBlock, + "first block", f.indexedRange.blocks.First(), "last block", f.indexedRange.blocks.Last(), + "removed maps", f.indexedRange.maps.First()-f.ptrTailUnindexMap, + "removed blocks", f.indexedRange.blocks.First()-f.tailPartialBlocks()-f.ptrTailUnindexBlock, "elapsed", common.PrettyDuration(time.Since(f.startedTailUnindexAt))) f.startedTailUnindex = false } @@ -342,7 +342,7 @@ func (f *FilterMaps) tryUnindexTail() bool { // needTailEpoch returns true if the given tail epoch needs to be kept // according to the current tail target, false if it can be removed. func (f *FilterMaps) needTailEpoch(epoch uint32) bool { - firstEpoch := f.indexedRange.firstRenderedMap >> f.logMapsPerEpoch + firstEpoch := f.indexedRange.maps.First() >> f.logMapsPerEpoch if epoch > firstEpoch { return true } @@ -382,15 +382,15 @@ func (f *FilterMaps) tailPartialBlocks() uint64 { if f.indexedRange.tailPartialEpoch == 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 { - 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 - if f.indexedRange.firstRenderedMap-f.mapsPerEpoch > 0 { - start, _, err = f.getLastBlockOfMap(f.indexedRange.firstRenderedMap - f.mapsPerEpoch - 1) + if f.indexedRange.maps.First()-f.mapsPerEpoch > 0 { + start, _, err = f.getLastBlockOfMap(f.indexedRange.maps.First() - f.mapsPerEpoch - 1) 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 @@ -399,5 +399,5 @@ func (f *FilterMaps) tailPartialBlocks() uint64 { // targetHeadIndexed returns true if the current log index is consistent with // targetView with its head block fully rendered. func (f *FilterMaps) targetHeadIndexed() bool { - return equalViews(f.targetView, f.indexedView) && f.indexedRange.headBlockIndexed + return equalViews(f.targetView, f.indexedView) && f.indexedRange.headIndexed } diff --git a/core/filtermaps/indexer_test.go b/core/filtermaps/indexer_test.go index 74ab5b0f5d..a02f8d2459 100644 --- a/core/filtermaps/indexer_test.go +++ b/core/filtermaps/indexer_test.go @@ -144,16 +144,16 @@ func TestIndexerRandomRange(t *testing.T) { expTailBlock++ } } - if ts.fm.indexedRange.afterLastIndexedBlock != uint64(head+1) { - ts.t.Fatalf("Invalid index head (expected #%d, got #%d)", head, ts.fm.indexedRange.afterLastIndexedBlock-1) + if ts.fm.indexedRange.blocks.Last() != uint64(head) { + ts.t.Fatalf("Invalid index head (expected #%d, got #%d)", head, ts.fm.indexedRange.blocks.Last()) } expHeadDelimiter := expdpos(uint64(head)) - if ts.fm.indexedRange.headBlockDelimiter != expHeadDelimiter { - ts.t.Fatalf("Invalid index head delimiter pointer (expected %d, got %d)", expHeadDelimiter, ts.fm.indexedRange.headBlockDelimiter) + if ts.fm.indexedRange.headDelimiter != expHeadDelimiter { + ts.t.Fatalf("Invalid index head delimiter pointer (expected %d, got %d)", expHeadDelimiter, ts.fm.indexedRange.headDelimiter) } - if ts.fm.indexedRange.firstIndexedBlock != expTailBlock { - ts.t.Fatalf("Invalid index tail block (expected #%d, got #%d)", expTailBlock, ts.fm.indexedRange.firstIndexedBlock) + if ts.fm.indexedRange.blocks.First() != expTailBlock { + ts.t.Fatalf("Invalid index tail block (expected #%d, got #%d)", expTailBlock, ts.fm.indexedRange.blocks.First()) } } } diff --git a/core/filtermaps/map_renderer.go b/core/filtermaps/map_renderer.go index 9bed3b41c0..29a677cdec 100644 --- a/core/filtermaps/map_renderer.go +++ b/core/filtermaps/map_renderer.go @@ -46,12 +46,12 @@ var ( // mapRenderer represents a process that renders filter maps in a specified // range according to the actual targetView. type mapRenderer struct { - f *FilterMaps - afterLastMap uint32 - currentMap *renderedMap - finishedMaps map[uint32]*renderedMap - firstFinished, afterLastFinished uint32 - iterator *logIterator + f *FilterMaps + renderBefore uint32 + currentMap *renderedMap + finishedMaps map[uint32]*renderedMap + finished common.Range[uint32] + iterator *logIterator } // 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 // point that is consistent with the current targetView. // 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, // otherwise it is checked that the rendered range is consistent with both // views. -func (f *FilterMaps) renderMapsBefore(afterLastMap uint32) (*mapRenderer, error) { - nextMap, startBlock, startLvPtr, err := f.lastCanonicalMapBoundaryBefore(afterLastMap) +func (f *FilterMaps) renderMapsBefore(renderBefore uint32) (*mapRenderer, error) { + nextMap, startBlock, startLvPtr, err := f.lastCanonicalMapBoundaryBefore(renderBefore) if err != nil { 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) } - if nextMap >= afterLastMap { + if nextMap >= renderBefore { 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 @@ -108,17 +108,16 @@ func (f *FilterMaps) renderMapsFromSnapshot(cp *renderedMap) (*mapRenderer, erro lastBlock: cp.lastBlock, blockLvPtrs: cp.blockLvPtrs, }, - finishedMaps: make(map[uint32]*renderedMap), - firstFinished: cp.mapIndex, - afterLastFinished: cp.mapIndex, - afterLastMap: math.MaxUint32, - iterator: iter, + finishedMaps: make(map[uint32]*renderedMap), + finished: common.NewRange(cp.mapIndex, 0), + renderBefore: math.MaxUint32, + iterator: iter, }, nil } // renderMapsFromMapBoundary creates a mapRenderer that starts rendering at a // 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) if err != nil { 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, lastBlock: iter.blockNumber, }, - finishedMaps: make(map[uint32]*renderedMap), - firstFinished: firstMap, - afterLastFinished: firstMap, - afterLastMap: afterLastMap, - iterator: iter, + finishedMaps: make(map[uint32]*renderedMap), + finished: common.NewRange(firstMap, 0), + renderBefore: renderBefore, + iterator: iter, }, nil } // lastCanonicalSnapshotBefore returns the latest cached snapshot that matches // the current targetView. -func (f *FilterMaps) lastCanonicalSnapshotBefore(afterLastMap uint32) *renderedMap { +func (f *FilterMaps) lastCanonicalSnapshotBefore(renderBefore uint32) *renderedMap { var best *renderedMap 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 && - cp.mapIndex < afterLastMap && (best == nil || blockNumber > best.lastBlock) { + cp.mapIndex < renderBefore && (best == nil || blockNumber > best.lastBlock) { best = cp } } @@ -158,11 +156,11 @@ func (f *FilterMaps) lastCanonicalSnapshotBefore(afterLastMap uint32) *renderedM // or the boundary of a currently rendered map. // 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. -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 { return 0, 0, 0, nil } - mapIndex := afterLastMap + mapIndex := renderBefore for { var ok bool 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 // map index. 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 } - if mapIndex > f.indexedRange.afterLastRenderedMap { - mapIndex = f.indexedRange.afterLastRenderedMap + if mapIndex > f.indexedRange.maps.AfterLast() { + mapIndex = f.indexedRange.maps.AfterLast() } - if mapIndex > f.indexedRange.firstRenderedMap { + if mapIndex > f.indexedRange.maps.First() { return mapIndex - 1, true } - if mapIndex+f.mapsPerEpoch > f.indexedRange.firstRenderedMap { - if mapIndex > f.indexedRange.firstRenderedMap-f.mapsPerEpoch+f.indexedRange.tailPartialEpoch { - mapIndex = f.indexedRange.firstRenderedMap - f.mapsPerEpoch + f.indexedRange.tailPartialEpoch + if mapIndex+f.mapsPerEpoch > f.indexedRange.maps.First() { + if mapIndex > f.indexedRange.maps.First()-f.mapsPerEpoch+f.indexedRange.tailPartialEpoch { + mapIndex = f.indexedRange.maps.First() - f.mapsPerEpoch + f.indexedRange.tailPartialEpoch } } else { 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 // a snapshot. func (f *FilterMaps) loadHeadSnapshot() error { - fm, err := f.getFilterMap(f.indexedRange.afterLastRenderedMap - 1) + fm, err := f.getFilterMap(f.indexedRange.maps.Last()) 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 { - 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 - if f.indexedRange.afterLastRenderedMap > 1 { - prevLastBlock, _, err := f.getLastBlockOfMap(f.indexedRange.afterLastRenderedMap - 2) + if f.indexedRange.maps.AfterLast() > 1 { + prevLastBlock, _, err := f.getLastBlockOfMap(f.indexedRange.maps.Last() - 1) 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 } @@ -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) } } - f.renderSnapshots.Add(f.indexedRange.afterLastIndexedBlock-1, &renderedMap{ + f.renderSnapshots.Add(f.indexedRange.blocks.Last(), &renderedMap{ filterMap: fm, - mapIndex: f.indexedRange.afterLastRenderedMap - 1, - lastBlock: f.indexedRange.afterLastIndexedBlock - 1, - lastBlockId: f.indexedView.getBlockId(f.indexedRange.afterLastIndexedBlock - 1), + mapIndex: f.indexedRange.maps.Last(), + lastBlock: f.indexedRange.blocks.Last(), + lastBlockId: f.indexedView.getBlockId(f.indexedRange.blocks.Last()), blockLvPtrs: lvPtrs, finished: true, - headDelimiter: f.indexedRange.headBlockDelimiter, + headDelimiter: f.indexedRange.headDelimiter, }) return nil } @@ -277,14 +275,14 @@ func (r *mapRenderer) run(stopCb func() bool, writeCb func()) (bool, error) { } // map finished r.finishedMaps[r.currentMap.mapIndex] = r.currentMap - r.afterLastFinished++ - if len(r.finishedMaps) >= maxMapsPerBatch || r.afterLastFinished&(r.f.baseRowGroupLength-1) == 0 { + r.finished.SetLast(r.finished.AfterLast()) + if len(r.finishedMaps) >= maxMapsPerBatch || r.finished.AfterLast()&(r.f.baseRowGroupLength-1) == 0 { if err := r.writeFinishedMaps(stopCb); err != nil { return false, err } 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 { return false, err } @@ -293,7 +291,7 @@ func (r *mapRenderer) run(stopCb func() bool, writeCb func()) (bool, error) { } r.currentMap = &renderedMap{ 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 { 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.makeSnapshot() } @@ -403,7 +401,7 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error { mapIndices []uint32 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] if fm, _ := r.f.filterMapCache.Get(mapIndex); fm != nil && row.Equal(fm[rowIndex]) { continue @@ -411,8 +409,8 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error { mapIndices = append(mapIndices, mapIndex) rows = append(rows, row) } - if newRange.afterLastRenderedMap == r.afterLastFinished { // head updated; remove future entries - for mapIndex := r.afterLastFinished; mapIndex < oldRange.afterLastRenderedMap; mapIndex++ { + if newRange.maps.AfterLast() == r.finished.AfterLast() { // head updated; remove future entries + for mapIndex := r.finished.AfterLast(); mapIndex < oldRange.maps.AfterLast(); mapIndex++ { if fm, _ := r.f.filterMapCache.Get(mapIndex); fm != nil && len(fm[rowIndex]) == 0 { continue } @@ -426,24 +424,24 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error { checkWriteCnt() } // 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 - 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) } - for mapIndex := r.afterLastFinished; mapIndex < oldRange.afterLastRenderedMap; mapIndex++ { + for mapIndex := r.finished.AfterLast(); mapIndex < oldRange.maps.AfterLast(); mapIndex++ { r.f.filterMapCache.Remove(mapIndex) } } else { // head not updated; do not cache maps during tail rendering because we // 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) } } // add or update block pointers - blockNumber := r.finishedMaps[r.firstFinished].firstBlock() - for mapIndex := r.firstFinished; mapIndex < r.afterLastFinished; mapIndex++ { + blockNumber := r.finishedMaps[r.finished.First()].firstBlock() + for mapIndex := r.finished.First(); mapIndex < r.finished.AfterLast(); mapIndex++ { renderedMap := r.finishedMaps[mapIndex] r.f.storeLastBlockOfMap(batch, mapIndex, renderedMap.lastBlock, renderedMap.lastBlockId) checkWriteCnt() @@ -456,18 +454,18 @@ func (r *mapRenderer) writeFinishedMaps(pauseCb func() bool) error { blockNumber++ } } - if newRange.afterLastRenderedMap == r.afterLastFinished { // head updated; remove future entries - for mapIndex := r.afterLastFinished; mapIndex < oldRange.afterLastRenderedMap; mapIndex++ { + if newRange.maps.AfterLast() == r.finished.AfterLast() { // head updated; remove future entries + for mapIndex := r.finished.AfterLast(); mapIndex < oldRange.maps.AfterLast(); mapIndex++ { r.f.deleteLastBlockOfMap(batch, mapIndex) checkWriteCnt() } - for ; blockNumber < oldRange.afterLastIndexedBlock; blockNumber++ { + for ; blockNumber < oldRange.blocks.AfterLast(); blockNumber++ { r.f.deleteBlockLvPointer(batch, blockNumber) checkWriteCnt() } } r.finishedMaps = make(map[uint32]*renderedMap) - r.firstFinished = r.afterLastFinished + r.finished.SetFirst(r.finished.AfterLast()) r.f.setRange(batch, renderedView, newRange) if err := batch.Write(); err != nil { 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. func (r *mapRenderer) getTempRange() (filterMapsRange, error) { 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) } - if tempRange.firstRenderedMap != r.f.indexedRange.firstRenderedMap { + if tempRange.maps.First() != r.f.indexedRange.maps.First() { // first rendered map changed; update first indexed block - if tempRange.firstRenderedMap > 0 { - lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.firstRenderedMap - 1) + if tempRange.maps.First() > 0 { + firstBlock, _, err := r.f.getLastBlockOfMap(tempRange.maps.First() - 1) 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 { - tempRange.firstIndexedBlock = 0 + tempRange.blocks.SetFirst(0) } } - if tempRange.afterLastRenderedMap != r.f.indexedRange.afterLastRenderedMap { - // first rendered map changed; update first indexed block - if tempRange.afterLastRenderedMap > 0 { - lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.afterLastRenderedMap - 1) + if tempRange.maps.AfterLast() != r.f.indexedRange.maps.AfterLast() { + // last rendered map changed; update last indexed block + if !tempRange.maps.IsEmpty() { + lastBlock, _, err := r.f.getLastBlockOfMap(tempRange.maps.Last()) 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 { - tempRange.afterLastIndexedBlock = 0 + tempRange.blocks.SetAfterLast(0) } - tempRange.headBlockDelimiter = 0 + tempRange.headDelimiter = 0 } return tempRange, nil } @@ -518,39 +516,39 @@ func (r *mapRenderer) getTempRange() (filterMapsRange, error) { func (r *mapRenderer) getUpdatedRange() (filterMapsRange, error) { // update 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.finished.First(), r.finished.AfterLast(), r.renderBefore, r.f.mapsPerEpoch); err != nil { 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 - if newRange.firstRenderedMap > 0 { - lastBlock, _, err := r.f.getLastBlockOfMap(newRange.firstRenderedMap - 1) + if newRange.maps.First() > 0 { + firstBlock, _, err := r.f.getLastBlockOfMap(newRange.maps.First() - 1) 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 { - 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 - lm := r.finishedMaps[r.afterLastFinished-1] - newRange.headBlockIndexed = lm.finished + lm := r.finishedMaps[r.finished.Last()] + newRange.headIndexed = 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 { panic("map rendering finished but last block != head block") } - newRange.headBlockDelimiter = lm.headDelimiter + newRange.headDelimiter = lm.headDelimiter } else { - newRange.afterLastIndexedBlock = lm.lastBlock - newRange.headBlockDelimiter = 0 + newRange.blocks.SetAfterLast(lm.lastBlock) // lastBlock is probably partially rendered + newRange.headDelimiter = 0 } } else { // last rendered map not replaced; ensure that target chain view matches // 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 } } @@ -572,9 +570,9 @@ func (fmr *filterMapsRange) addRenderedRange(firstRendered, afterLastRendered, a m uint32 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 { - 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 }) var ( @@ -597,14 +595,12 @@ func (fmr *filterMapsRange) addRenderedRange(firstRendered, afterLastRendered, a case 0: // Initialized database, but no finished maps yet. fmr.tailPartialEpoch = 0 - fmr.firstRenderedMap = firstRendered - fmr.afterLastRenderedMap = firstRendered + fmr.maps = common.NewRange(firstRendered, 0) case 2: // One rendered section (no partial tail epoch). fmr.tailPartialEpoch = 0 - fmr.firstRenderedMap = merged[0] - fmr.afterLastRenderedMap = merged[1] + fmr.maps = common.NewRange(merged[0], merged[1]-merged[0]) case 4: // 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) } fmr.tailPartialEpoch = merged[1] - merged[0] - fmr.firstRenderedMap = merged[2] - fmr.afterLastRenderedMap = merged[3] + fmr.maps = common.NewRange(merged[2], merged[3]-merged[2]) default: 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 { 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 } var lvIndex uint64 - if f.indexedRange.headBlockIndexed && blockNumber+1 == f.indexedRange.afterLastIndexedBlock { - lvIndex = f.indexedRange.headBlockDelimiter + if f.indexedRange.headIndexed && blockNumber+1 == f.indexedRange.blocks.AfterLast() { + lvIndex = f.indexedRange.headDelimiter } else { var err error lvIndex, err = f.getBlockLvPointer(blockNumber + 1) diff --git a/core/filtermaps/matcher.go b/core/filtermaps/matcher.go index eb40629381..973b071937 100644 --- a/core/filtermaps/matcher.go +++ b/core/filtermaps/matcher.go @@ -61,11 +61,9 @@ type SyncRange struct { // 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 // be valid and complete. - Valid bool - FirstValid, LastValid uint64 + ValidBlocks common.Range[uint64] // block range indexed according to the given chain head. - Indexed bool - FirstIndexed, LastIndexed uint64 + IndexedBlocks common.Range[uint64] } // 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 { + getLogStats runtimeStats // 64 bit aligned ctx context.Context backend MatcherBackend params *Params matcher matcher - getLogStats runtimeStats firstIndex, lastIndex uint64 firstMap, lastMap uint32 } diff --git a/core/filtermaps/matcher_backend.go b/core/filtermaps/matcher_backend.go index 66f95f2145..4c4668e321 100644 --- a/core/filtermaps/matcher_backend.go +++ b/core/filtermaps/matcher_backend.go @@ -19,6 +19,7 @@ package filtermaps import ( "context" + "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/core/types" ) @@ -27,9 +28,8 @@ type FilterMapsMatcherBackend struct { f *FilterMaps // these fields should be accessed under f.matchersLock mutex. - valid bool - firstValid, lastValid uint64 - syncCh chan SyncRange + validBlocks common.Range[uint64] + syncCh chan SyncRange } // NewMatcherBackend returns a FilterMapsMatcherBackend after registering it in @@ -43,11 +43,9 @@ func (f *FilterMaps) NewMatcherBackend() *FilterMapsMatcherBackend { f.indexLock.RUnlock() }() - fm := &FilterMapsMatcherBackend{ - f: f, - valid: f.indexedRange.initialized && f.indexedRange.afterLastIndexedBlock > f.indexedRange.firstIndexedBlock, - firstValid: f.indexedRange.firstIndexedBlock, - lastValid: f.indexedRange.afterLastIndexedBlock - 1, + fm := &FilterMapsMatcherBackend{f: f} + if f.indexedRange.initialized { + fm.validBlocks = f.indexedRange.blocks } f.matchers[fm] = struct{}{} return fm @@ -122,28 +120,16 @@ func (fm *FilterMapsMatcherBackend) synced() { fm.f.indexLock.RUnlock() }() - var ( - indexed bool - lastIndexed, subLastIndexed uint64 - ) - 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 + indexedBlocks := fm.f.indexedRange.blocks + if !fm.f.indexedRange.headIndexed && !indexedBlocks.IsEmpty() { + indexedBlocks.SetAfterLast(indexedBlocks.Last()) // remove partially indexed last block } fm.syncCh <- SyncRange{ - HeadNumber: fm.f.indexedView.headNumber, - Valid: fm.valid, - FirstValid: fm.firstValid, - LastValid: fm.lastValid, - Indexed: indexed, - FirstIndexed: fm.f.indexedRange.firstIndexedBlock, - LastIndexed: lastIndexed, + HeadNumber: fm.f.indexedView.headNumber, + ValidBlocks: fm.validBlocks, + IndexedBlocks: indexedBlocks, } - fm.valid = indexed - fm.firstValid = fm.f.indexedRange.firstIndexedBlock - fm.lastValid = lastIndexed + fm.validBlocks = indexedBlocks fm.syncCh = nil } @@ -187,20 +173,10 @@ func (f *FilterMaps) updateMatchersValidRange() { defer f.matchersLock.Unlock() for fm := range f.matchers { - if !f.indexedRange.hasIndexedBlocks() { - fm.valid = false - } - if !fm.valid { + if !f.indexedRange.initialized { + fm.validBlocks = common.Range[uint64]{} continue } - if fm.firstValid < f.indexedRange.firstIndexedBlock { - 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 - } + fm.validBlocks = fm.validBlocks.Intersection(f.indexedRange.blocks) } } diff --git a/core/rawdb/accessors_indexes.go b/core/rawdb/accessors_indexes.go index f679c3aeb3..70482e1847 100644 --- a/core/rawdb/accessors_indexes.go +++ b/core/rawdb/accessors_indexes.go @@ -356,8 +356,8 @@ func WriteFilterMapBaseRows(db ethdb.KeyValueWriter, mapRowIndex uint64, rows [] } } -func DeleteFilterMapRows(db ethdb.KeyValueRangeDeleter, firstMapRowIndex, afterLastMapRowIndex uint64) { - if err := db.DeleteRange(filterMapRowKey(firstMapRowIndex, false), filterMapRowKey(afterLastMapRowIndex, false)); err != nil { +func DeleteFilterMapRows(db ethdb.KeyValueRangeDeleter, mapRows common.Range[uint64]) { + 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) } } @@ -396,8 +396,8 @@ func DeleteFilterMapLastBlock(db ethdb.KeyValueWriter, mapIndex uint32) { } } -func DeleteFilterMapLastBlocks(db ethdb.KeyValueRangeDeleter, firstMapIndex, afterLastMapIndex uint32) { - if err := db.DeleteRange(filterMapLastBlockKey(firstMapIndex), filterMapLastBlockKey(afterLastMapIndex)); err != nil { +func DeleteFilterMapLastBlocks(db ethdb.KeyValueRangeDeleter, maps common.Range[uint32]) { + 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) } } @@ -433,8 +433,8 @@ func DeleteBlockLvPointer(db ethdb.KeyValueWriter, blockNumber uint64) { } } -func DeleteBlockLvPointers(db ethdb.KeyValueRangeDeleter, firstBlockNumber, afterLastBlockNumber uint64) { - if err := db.DeleteRange(filterMapBlockLVKey(firstBlockNumber), filterMapBlockLVKey(afterLastBlockNumber)); err != nil { +func DeleteBlockLvPointers(db ethdb.KeyValueRangeDeleter, blocks common.Range[uint64]) { + 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) } } @@ -442,10 +442,11 @@ func DeleteBlockLvPointers(db ethdb.KeyValueRangeDeleter, firstBlockNumber, afte // FilterMapsRange is a storage representation of the block range covered by the // filter maps structure and the corresponting log value index range. type FilterMapsRange struct { - HeadBlockIndexed bool - HeadBlockDelimiter uint64 - FirstIndexedBlock, AfterLastIndexedBlock uint64 - FirstRenderedMap, AfterLastRenderedMap, TailPartialEpoch uint32 + HeadIndexed bool + HeadDelimiter uint64 + Blocks common.Range[uint64] + Maps common.Range[uint32] + TailPartialEpoch uint32 } // ReadFilterMapsRange retrieves the filter maps range data. Note that if the