mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-27 15:16:43 +00:00
core/bloombits: using partialMatches instead of separate idx/bv channels
This commit is contained in:
parent
0957393047
commit
b9b9db2f99
1 changed files with 40 additions and 52 deletions
|
|
@ -150,7 +150,7 @@ func (f *fetcher) deliver(sectionIdxList []uint64, data [][]byte) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Matcher is a pipelined structure of fetchers and logic matchers which perform
|
// 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 {
|
type Matcher struct {
|
||||||
addresses []types.BloomIndexList
|
addresses []types.BloomIndexList
|
||||||
topics [][]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
|
// 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
|
// 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
|
// 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) {
|
func (m *Matcher) match(processCh chan partialMatches, stop chan struct{}) chan partialMatches {
|
||||||
subIdx := m.topics
|
subIdx := m.topics
|
||||||
if len(m.addresses) > 0 {
|
if len(m.addresses) > 0 {
|
||||||
subIdx = append([][]types.BloomIndexList{m.addresses}, subIdx...)
|
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.getNextReqCh = make(chan chan nextRequests) // should be a blocking channel
|
||||||
m.distributeRequests(stop)
|
m.distributeRequests(stop)
|
||||||
|
|
||||||
s := sectionCh
|
|
||||||
var bv chan []byte
|
|
||||||
for _, idx := range subIdx {
|
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
|
// 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
|
m.fetchers[idx] = f
|
||||||
}
|
}
|
||||||
|
|
||||||
// subMatch creates a sub-matcher that filters for a set of addresses or topics, binary or-s those matches, then
|
// 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.
|
// 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
|
// 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.
|
// 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) {
|
func (m *Matcher) subMatch(processCh chan partialMatches, idxs []types.BloomIndexList, stop chan struct{}) chan partialMatches {
|
||||||
// set up fetchers
|
// set up fetchers
|
||||||
fetchIdx := make([][3]chan uint64, len(idxs))
|
fetchIdx := make([][3]chan uint64, len(idxs))
|
||||||
fetchData := make([][3]chan []byte, 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)
|
fetchedCh := make(chan partialMatches, channelCap) // entries from processCh are forwarded here after fetches have been initiated
|
||||||
resIdxCh := make(chan uint64, channelCap)
|
resultsCh := make(chan partialMatches, channelCap)
|
||||||
resDataCh := make(chan []byte, channelCap)
|
|
||||||
|
|
||||||
m.wg.Add(2)
|
m.wg.Add(2)
|
||||||
// goroutine for starting retrievals
|
// goroutine for starting retrievals
|
||||||
|
|
@ -266,9 +271,9 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs
|
||||||
select {
|
select {
|
||||||
case <-stop:
|
case <-stop:
|
||||||
return
|
return
|
||||||
case s, ok := <-sectionCh:
|
case s, ok := <-processCh:
|
||||||
if !ok {
|
if !ok {
|
||||||
close(processCh)
|
close(fetchedCh)
|
||||||
for _, ff := range fetchIdx {
|
for _, ff := range fetchIdx {
|
||||||
for _, f := range ff {
|
for _, f := range ff {
|
||||||
close(f)
|
close(f)
|
||||||
|
|
@ -277,20 +282,20 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
select {
|
|
||||||
case <-stop:
|
|
||||||
return
|
|
||||||
case processCh <- s:
|
|
||||||
}
|
|
||||||
for _, ff := range fetchIdx {
|
for _, ff := range fetchIdx {
|
||||||
for _, f := range ff {
|
for _, f := range ff {
|
||||||
select {
|
select {
|
||||||
case <-stop:
|
case <-stop:
|
||||||
return
|
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 {
|
select {
|
||||||
case <-stop:
|
case <-stop:
|
||||||
return
|
return
|
||||||
case s, ok := <-processCh:
|
case s, ok := <-fetchedCh:
|
||||||
if !ok {
|
if !ok {
|
||||||
close(resIdxCh)
|
close(resultsCh)
|
||||||
close(resDataCh)
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -337,31 +341,21 @@ func (m *Matcher) subMatch(sectionCh chan uint64, andVectorCh chan []byte, idxs
|
||||||
if orVector == nil {
|
if orVector == nil {
|
||||||
orVector = make([]byte, int(m.sectionSize/8))
|
orVector = make([]byte, int(m.sectionSize/8))
|
||||||
}
|
}
|
||||||
if andVectorCh != nil {
|
if s.vector != nil {
|
||||||
select {
|
bitutil.ANDBytes(orVector, orVector, s.vector)
|
||||||
case <-stop:
|
|
||||||
return
|
|
||||||
case andVector := <-andVectorCh:
|
|
||||||
bitutil.ANDBytes(orVector, orVector, andVector)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
if bitutil.TestBytes(orVector) {
|
if bitutil.TestBytes(orVector) {
|
||||||
select {
|
select {
|
||||||
case <-stop:
|
case <-stop:
|
||||||
return
|
return
|
||||||
case resIdxCh <- s:
|
case resultsCh <- partialMatches{s.sectionIdx, orVector}:
|
||||||
}
|
|
||||||
select {
|
|
||||||
case <-stop:
|
|
||||||
return
|
|
||||||
case resDataCh <- orVector:
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
return resIdxCh, resDataCh
|
return resultsCh
|
||||||
}
|
}
|
||||||
|
|
||||||
// GetMatches returns a stream of bloom matches in a given range of blocks.
|
// 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 {
|
func (m *Matcher) GetMatches(start, end uint64, stop chan struct{}) chan uint64 {
|
||||||
m.distWg.Wait()
|
m.distWg.Wait()
|
||||||
|
|
||||||
sectionCh := make(chan uint64, channelCap)
|
processCh := make(chan partialMatches, channelCap)
|
||||||
resultsCh := make(chan uint64, channelCap)
|
resultsCh := make(chan uint64, channelCap)
|
||||||
|
|
||||||
s, bv := m.match(sectionCh, stop)
|
res := m.match(processCh, stop)
|
||||||
|
|
||||||
startSection := start / m.sectionSize
|
startSection := start / m.sectionSize
|
||||||
endSection := end / 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)
|
m.wg.Add(2)
|
||||||
go func() {
|
go func() {
|
||||||
defer m.wg.Done()
|
defer m.wg.Done()
|
||||||
defer close(sectionCh)
|
defer close(processCh)
|
||||||
|
|
||||||
for i := startSection; i <= endSection; i++ {
|
for i := startSection; i <= endSection; i++ {
|
||||||
select {
|
select {
|
||||||
case sectionCh <- i:
|
case processCh <- partialMatches{i, nil}:
|
||||||
case <-stop:
|
case <-stop:
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -400,17 +394,11 @@ func (m *Matcher) GetMatches(start, end uint64, stop chan struct{}) chan uint64
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case idx, ok := <-s:
|
case r, ok := <-res:
|
||||||
if !ok {
|
if !ok {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
var match []byte
|
sectionStart := r.sectionIdx * m.sectionSize
|
||||||
select {
|
|
||||||
case <-stop:
|
|
||||||
return
|
|
||||||
case match = <-bv:
|
|
||||||
}
|
|
||||||
sectionStart := idx * m.sectionSize
|
|
||||||
s := sectionStart
|
s := sectionStart
|
||||||
if start > s {
|
if start > s {
|
||||||
s = start
|
s = start
|
||||||
|
|
@ -420,7 +408,7 @@ func (m *Matcher) GetMatches(start, end uint64, stop chan struct{}) chan uint64
|
||||||
e = end
|
e = end
|
||||||
}
|
}
|
||||||
for i := s; i <= e; i++ {
|
for i := s; i <= e; i++ {
|
||||||
b := match[(i-sectionStart)/8]
|
b := r.vector[(i-sectionStart)/8]
|
||||||
bit := 7 - i%8
|
bit := 7 - i%8
|
||||||
if b != 0 {
|
if b != 0 {
|
||||||
if b&(1<<bit) != 0 {
|
if b&(1<<bit) != 0 {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue