diff --git a/core/filtermaps/filtermaps.go b/core/filtermaps/filtermaps.go index b14a4e6114..9e6136ca27 100644 --- a/core/filtermaps/filtermaps.go +++ b/core/filtermaps/filtermaps.go @@ -95,11 +95,11 @@ type FilterMaps struct { targetView *ChainView matcherSyncRequest *FilterMapsMatcherBackend + historyCutoff uint64 finalBlock, lastFinal uint64 lastFinalEpoch uint32 stop bool - targetViewCh chan *ChainView - finalBlockCh chan uint64 + targetCh chan targetUpdate blockProcessingCh chan bool blockProcessing bool matcherSyncCh chan *FilterMapsMatcherBackend @@ -183,7 +183,7 @@ type Config struct { } // NewFilterMaps creates a new FilterMaps and starts the indexer. -func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, params Params, config Config) *FilterMaps { +func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, historyCutoff, finalBlock uint64, params Params, config Config) *FilterMaps { rs, initialized, err := rawdb.ReadFilterMapsRange(db) if err != nil { log.Error("Error reading log index range", "error", err) @@ -193,8 +193,7 @@ func NewFilterMaps(db ethdb.KeyValueStore, initView *ChainView, params Params, c db: db, closeCh: make(chan struct{}), waitIdleCh: make(chan chan bool), - targetViewCh: make(chan *ChainView, 1), - finalBlockCh: make(chan uint64, 1), + targetCh: make(chan targetUpdate, 1), blockProcessingCh: make(chan bool, 1), history: config.History, disabled: config.Disabled, diff --git a/core/filtermaps/indexer.go b/core/filtermaps/indexer.go index 54a3687648..32416a8c47 100644 --- a/core/filtermaps/indexer.go +++ b/core/filtermaps/indexer.go @@ -66,28 +66,25 @@ func (f *FilterMaps) indexerLoop() { } } +type targetUpdate struct { + targetView *ChainView + historyCutoff, finalBlock uint64 +} + // SetTargetView sets a new target chain view for the indexer to render. // Note that SetTargetView never blocks. -func (f *FilterMaps) SetTargetView(targetView *ChainView) { +func (f *FilterMaps) SetTarget(targetView *ChainView, historyCutoff, finalBlock uint64) { if targetView == nil { panic("nil targetView") } for { select { - case <-f.targetViewCh: - case f.targetViewCh <- targetView: - return - } - } -} - -// SetFinalBlock sets the finalized block number used for exporting checkpoints. -// Note that SetFinalBlock never blocks. -func (f *FilterMaps) SetFinalBlock(finalBlock uint64) { - for { - select { - case <-f.finalBlockCh: - case f.finalBlockCh <- finalBlock: + case <-f.targetCh: + case f.targetCh <- targetUpdate{ + targetView: targetView, + historyCutoff: historyCutoff, + finalBlock: finalBlock, + }: return } } @@ -145,26 +142,24 @@ func (f *FilterMaps) processSingleEvent(blocking bool) bool { } if blocking { select { - case targetView := <-f.targetViewCh: - f.setTargetView(targetView) - case f.finalBlock = <-f.finalBlockCh: + case target := <-f.targetCh: + f.setTarget(target) case f.matcherSyncRequest = <-f.matcherSyncCh: case f.blockProcessing = <-f.blockProcessingCh: case <-f.closeCh: f.stop = true case ch := <-f.waitIdleCh: select { - case targetView := <-f.targetViewCh: - f.setTargetView(targetView) + case target := <-f.targetCh: + f.setTarget(target) default: } ch <- !f.blockProcessing && f.targetHeadIndexed() } } else { select { - case targetView := <-f.targetViewCh: - f.setTargetView(targetView) - case f.finalBlock = <-f.finalBlockCh: + case target := <-f.targetCh: + f.setTarget(target) case f.matcherSyncRequest = <-f.matcherSyncCh: case f.blockProcessing = <-f.blockProcessingCh: case <-f.closeCh: @@ -177,8 +172,10 @@ func (f *FilterMaps) processSingleEvent(blocking bool) bool { } // setTargetView updates the target chain view of the iterator. -func (f *FilterMaps) setTargetView(targetView *ChainView) { - f.targetView = targetView +func (f *FilterMaps) setTarget(target targetUpdate) { + f.targetView = target.targetView + f.historyCutoff = target.historyCutoff + f.finalBlock = target.finalBlock } // tryIndexHead tries to render head maps according to the current targetView @@ -234,7 +231,11 @@ func (f *FilterMaps) tryIndexHead() bool { // rendered according to targetView and is suspended as soon as the targetView // is changed. func (f *FilterMaps) tryIndexTail() bool { - for firstEpoch := f.indexedRange.firstRenderedMap >> f.logMapsPerEpoch; firstEpoch > 0 && f.needTailEpoch(firstEpoch-1); { + for { + firstEpoch := f.indexedRange.firstRenderedMap >> f.logMapsPerEpoch + if firstEpoch == 0 || !f.needTailEpoch(firstEpoch-1) { + return true + } f.processEvents() if f.stop || !f.targetHeadIndexed() { return false @@ -283,7 +284,8 @@ func (f *FilterMaps) tryIndexTail() bool { f.lastLogTailIndex = time.Now() } }) - if err != nil { + if err != nil && f.needTailEpoch(firstEpoch-1) { + // stop silently if cutoff point has move beyond epoch boundary while rendering log.Error("Log index tail rendering failed", "error", err) } if !done { @@ -347,16 +349,22 @@ func (f *FilterMaps) needTailEpoch(epoch uint32) bool { if epoch+1 < firstEpoch { return false } - tailTarget := f.tailTargetBlock() - if tailTarget < f.indexedRange.firstIndexedBlock { - return true + if epoch > 0 { + lastBlockOfPrevEpoch, _, err := f.getLastBlockOfMap(epoch<= firstEpoch + } + if f.historyCutoff > lastBlockOfPrevEpoch { + return false + } } - tailLvIndex, err := f.getBlockLvPointer(tailTarget) + lastBlockOfEpoch, _, err := f.getLastBlockOfMap((epoch+1)<= firstEpoch } - return uint64(epoch+1)<<(f.logValuesPerMap+f.logMapsPerEpoch) >= tailLvIndex + return f.tailTargetBlock() <= lastBlockOfEpoch } // tailTargetBlock returns the target value for the tail block number according diff --git a/core/filtermaps/indexer_test.go b/core/filtermaps/indexer_test.go index 17ad765e04..355a25528c 100644 --- a/core/filtermaps/indexer_test.go +++ b/core/filtermaps/indexer_test.go @@ -235,7 +235,7 @@ func (ts *testSetup) setHistory(history uint64, noHistory bool) { History: history, Disabled: noHistory, } - ts.fm = NewFilterMaps(ts.db, view, ts.params, config) + ts.fm = NewFilterMaps(ts.db, view, 0, 0, ts.params, config) ts.fm.testDisableSnapshots = ts.testDisableSnapshots ts.fm.Start() } @@ -420,7 +420,7 @@ func (tc *testChain) setTargetHead() { if tc.ts.fm != nil { if !tc.ts.fm.disabled { //tc.ts.fm.targetViewCh <- NewChainView(tc, head.Number.Uint64(), head.Hash()) - tc.ts.fm.SetTargetView(NewChainView(tc, head.Number.Uint64(), head.Hash())) + tc.ts.fm.SetTarget(NewChainView(tc, head.Number.Uint64(), head.Hash()), 0, 0) } } }