mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-25 06:06:44 +00:00
p2p: allow several asyncFilter lookups in parallel
Signed-off-by: Csaba Kiraly <csaba.kiraly@gmail.com>
This commit is contained in:
parent
721b494854
commit
26feecf3d0
2 changed files with 39 additions and 11 deletions
|
|
@ -491,7 +491,7 @@ func (s *Ethereum) setupDiscovery() error {
|
|||
if s.p2pServer.DiscoveryV4() != nil {
|
||||
asyncFilter := s.p2pServer.DiscoveryV4().RequestENR
|
||||
filter := eth.NewNodeFilter(s.blockchain)
|
||||
iter := enode.AsyncFilter(s.p2pServer.DiscoveryV4().RandomNodes(), asyncFilter)
|
||||
iter := enode.AsyncFilter(s.p2pServer.DiscoveryV4().RandomNodes(), asyncFilter, 16)
|
||||
iter = enode.Filter(iter, filter)
|
||||
s.discmix.AddSource(iter)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -154,25 +154,52 @@ func (f *filterIter) Next() bool {
|
|||
|
||||
// AsyncFilter wraps an iterator such that Next only returns nodes for which
|
||||
// the 'check' function returns a (possibly modified) node.
|
||||
func AsyncFilter(it Iterator, check func(*Node) (*Node, error)) Iterator {
|
||||
return &AsyncFilterIter{it, nil, check}
|
||||
func AsyncFilter(it Iterator, check func(*Node) (*Node, error), workers int) Iterator {
|
||||
f := &AsyncFilterIter{it, nil, check, make(chan *Node), sync.WaitGroup{}}
|
||||
|
||||
taskCh := make(chan *Node)
|
||||
|
||||
worker := func() {
|
||||
for task := range taskCh {
|
||||
if task == nil {
|
||||
break
|
||||
}
|
||||
nn, err := f.check(task)
|
||||
if err == nil {
|
||||
f.passed <- nn
|
||||
}
|
||||
}
|
||||
f.wg.Done()
|
||||
}
|
||||
|
||||
for range workers {
|
||||
f.wg.Add(1)
|
||||
go worker()
|
||||
}
|
||||
|
||||
go func() {
|
||||
for f.it.Next() {
|
||||
taskCh <- f.it.Node()
|
||||
}
|
||||
close(taskCh)
|
||||
f.wg.Wait()
|
||||
close(f.passed)
|
||||
}()
|
||||
|
||||
return f
|
||||
}
|
||||
|
||||
type AsyncFilterIter struct {
|
||||
it Iterator
|
||||
buffer *Node
|
||||
check func(*Node) (*Node, error)
|
||||
passed chan *Node
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
func (f *AsyncFilterIter) Next() bool {
|
||||
for f.it.Next() {
|
||||
nn, err := f.check(f.it.Node())
|
||||
if err == nil {
|
||||
f.buffer = nn
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
f.buffer = <-f.passed
|
||||
return f.buffer != nil
|
||||
}
|
||||
|
||||
func (f *AsyncFilterIter) Node() *Node {
|
||||
|
|
@ -181,6 +208,7 @@ func (f *AsyncFilterIter) Node() *Node {
|
|||
|
||||
func (f *AsyncFilterIter) Close() {
|
||||
f.it.Close()
|
||||
f.wg.Wait()
|
||||
}
|
||||
|
||||
// FairMix aggregates multiple node iterators. The mixer itself is an iterator which ends
|
||||
|
|
|
|||
Loading…
Reference in a new issue