diff --git a/core/bloombits/matcher.go b/core/bloombits/matcher.go index 2423146e0d..f9d9aefc5a 100644 --- a/core/bloombits/matcher.go +++ b/core/bloombits/matcher.go @@ -150,7 +150,7 @@ func (f *fetcher) deliver(sectionIdxList []uint64, data [][]byte) { } // Matcher is a pipelined structure of fetchers and logic matchers which perform -// binary and/or operations on the bitstreams, finally creating a stream of potential matches. +// binary AND/OR operations on the bitstreams, finally creating a stream of potential matches. type Matcher struct { addresses []types.BloomIndexList topics [][]types.BloomIndexList @@ -209,8 +209,8 @@ loop: // match creates a daisy-chain of sub-matchers, one for the address set and one for each topic set, each // sub-matcher receiving a section only if the previous ones have all found a potential match in one of -// the blocks of the section, then binary and-ing its own matches and forwaring the result to the next one -func (m *Matcher) match(sectionCh chan uint64, stop chan struct{}) (chan uint64, chan []byte) { +// the blocks of the section, then binary AND-ing its own matches and forwaring the result to the next one +func (m *Matcher) match(processCh chan partialMatches, stop chan struct{}) chan partialMatches { subIdx := m.topics if len(m.addresses) > 0 { subIdx = append([][]types.BloomIndexList{m.addresses}, subIdx...) @@ -218,12 +218,18 @@ func (m *Matcher) match(sectionCh chan uint64, stop chan struct{}) (chan uint64, m.getNextReqCh = make(chan chan nextRequests) // should be a blocking channel m.distributeRequests(stop) - s := sectionCh - var bv chan []byte for _, idx := range subIdx { - s, bv = m.subMatch(s, bv, idx, stop) + processCh = m.subMatch(processCh, idx, stop) } - return s, bv + return processCh +} + +// partialMatches with a non-nil vector represents a section in which some sub-matchers have already +// found potential matches. Subsequent sub-matchers will binary AND their matches with this vector. +// If vector is nil, it represents a section to be processed by the first sub-matcher. +type partialMatches struct { + sectionIdx uint64 + vector []byte } // newFetcher adds a fetcher for the given bit index if it has not existed before @@ -238,11 +244,11 @@ func (m *Matcher) newFetcher(idx uint) { m.fetchers[idx] = f } -// subMatch creates a sub-matcher that filters for a set of addresses or topics, binary or-s those matches, then -// binary and-s the result to the daisy-chain input (sectionCh/andVectorCh) and forwards it to the daisy-chain output. +// subMatch creates a sub-matcher that filters for a set of addresses or topics, binary OR-s those matches, then +// binary AND-s the result to the daisy-chain input (processCh) and forwards it to the daisy-chain output. // The matches of each address/topic are calculated by fetching the given sections of the three bloom bit indexes belonging to -// that address/topic, and binary and-ing those vectors together. -func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs []types.BloomIndexList, stop chan struct{}) (chan uint64, chan []byte) { +// that address/topic, and binary AND-ing those vectors together. +func (m *Matcher) subMatch(processCh chan partialMatches, idxs []types.BloomIndexList, stop chan struct{}) chan partialMatches { // set up fetchers fetchIdx := make([][3]chan uint64, len(idxs)) fetchData := make([][3]chan []byte, len(idxs)) @@ -253,9 +259,8 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs } } - processCh := make(chan uint64, channelCap) - resIdxCh := make(chan uint64, channelCap) - resDataCh := make(chan []byte, channelCap) + fetchedCh := make(chan partialMatches, channelCap) // entries from processCh are forwarded here after fetches have been initiated + resultsCh := make(chan partialMatches, channelCap) m.wg.Add(2) // goroutine for starting retrievals @@ -266,9 +271,9 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs select { case <-stop: return - case s, ok := <-sectionCh: + case s, ok := <-processCh: if !ok { - close(processCh) + close(fetchedCh) for _, ff := range fetchIdx { for _, f := range ff { close(f) @@ -277,20 +282,20 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs return } - select { - case <-stop: - return - case processCh <- s: - } for _, ff := range fetchIdx { for _, f := range ff { select { case <-stop: return - case f <- s: + case f <- s.sectionIdx: } } } + select { + case <-stop: + return + case fetchedCh <- s: + } } } }() @@ -303,10 +308,9 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs select { case <-stop: return - case s, ok := <-processCh: + case s, ok := <-fetchedCh: if !ok { - close(resIdxCh) - close(resDataCh) + close(resultsCh) return } @@ -337,31 +341,21 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs if orVector == nil { orVector = make([]byte, int(m.sectionSize/8)) } - if andVectorCh != nil { - select { - case <-stop: - return - case andVector := <-andVectorCh: - bitutil.ANDBytes(orVector, orVector, andVector) - } + if s.vector != nil { + bitutil.ANDBytes(orVector, orVector, s.vector) } if bitutil.TestBytes(orVector) { select { case <-stop: return - case resIdxCh <- s: - } - select { - case <-stop: - return - case resDataCh <- orVector: + case resultsCh <- partialMatches{s.sectionIdx, orVector}: } } } } }() - return resIdxCh, resDataCh + return resultsCh } // GetMatches returns a stream of bloom matches in a given range of blocks. @@ -372,10 +366,10 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs func (m *Matcher) GetMatches(start, end uint64, stop chan struct{}) chan uint64 { m.distWg.Wait() - sectionCh := make(chan uint64, channelCap) + processCh := make(chan partialMatches, channelCap) resultsCh := make(chan uint64, channelCap) - s, bv := m.match(sectionCh, stop) + res := m.match(processCh, stop) startSection := start / m.sectionSize endSection := end / m.sectionSize @@ -383,11 +377,11 @@ func (m *Matcher) GetMatches(start, end uint64, stop chan struct{}) chan uint64 m.wg.Add(2) go func() { defer m.wg.Done() - defer close(sectionCh) + defer close(processCh) for i := startSection; i <= endSection; i++ { select { - case sectionCh <- i: + case processCh <- partialMatches{i, nil}: case <-stop: return } @@ -400,17 +394,11 @@ func (m *Matcher) GetMatches(start, end uint64, stop chan struct{}) chan uint64 for { select { - case idx, ok := <-s: + case r, ok := <-res: if !ok { return } - var match []byte - select { - case <-stop: - return - case match = <-bv: - } - sectionStart := idx * m.sectionSize + sectionStart := r.sectionIdx * m.sectionSize s := sectionStart if start > s { s = start @@ -420,7 +408,7 @@ func (m *Matcher) GetMatches(start, end uint64, stop chan struct{}) chan uint64 e = end } for i := s; i <= e; i++ { - b := match[(i-sectionStart)/8] + b := r.vector[(i-sectionStart)/8] bit := 7 - i%8 if b != 0 { if b&(1<