mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-25 14:16:44 +00:00
p2p/enode: fixup BufferIterator
Signed-off-by: Csaba Kiraly <csaba.kiraly@gmail.com>
This commit is contained in:
parent
414021a068
commit
83e5556eb4
1 changed files with 18 additions and 6 deletions
|
|
@ -238,26 +238,39 @@ type BufferIter struct {
|
||||||
it Iterator
|
it Iterator
|
||||||
buffer chan *Node
|
buffer chan *Node
|
||||||
head *Node
|
head *Node
|
||||||
|
closed chan struct{}
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewBufferIter creates a new pre-fetch buffer.
|
// NewBufferIter creates a new pre-fetch buffer of a given size.
|
||||||
func NewBufferIter(it Iterator, size int) *BufferIter {
|
func NewBufferIter(it Iterator, size int) *BufferIter {
|
||||||
b := BufferIter{
|
b := BufferIter{
|
||||||
it: it,
|
it: it,
|
||||||
buffer: make(chan *Node, size),
|
buffer: make(chan *Node, size),
|
||||||
|
head: nil,
|
||||||
|
closed: make(chan struct{}),
|
||||||
}
|
}
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
|
// if the wrapped iterator ends, the buffer content will still be served.
|
||||||
defer close(b.buffer)
|
defer close(b.buffer)
|
||||||
|
// If instead the bufferIterator is closed, we bail out of the loop.
|
||||||
for b.it.Next() {
|
for b.it.Next() {
|
||||||
b.buffer <- b.it.Node()
|
select {
|
||||||
|
case b.buffer <- b.it.Node():
|
||||||
|
case <-b.closed:
|
||||||
|
return
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
return &b
|
return &b
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *BufferIter) Next() bool {
|
func (b *BufferIter) Next() bool {
|
||||||
b.head = <-b.buffer
|
select {
|
||||||
|
case b.head = <-b.buffer:
|
||||||
|
case <-b.closed:
|
||||||
|
return false
|
||||||
|
}
|
||||||
return b.head != nil
|
return b.head != nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -266,12 +279,11 @@ func (b *BufferIter) Node() *Node {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *BufferIter) Close() {
|
func (b *BufferIter) Close() {
|
||||||
|
// Close the wrapped iterator first.
|
||||||
b.it.Close()
|
b.it.Close()
|
||||||
// Wait for the buffer to be consumed.
|
close(b.closed)
|
||||||
for range b.buffer {
|
for range b.buffer {
|
||||||
}
|
}
|
||||||
// Close the buffer channel.
|
|
||||||
// close(b.buffer)
|
|
||||||
b.buffer = nil
|
b.buffer = nil
|
||||||
b.head = nil
|
b.head = nil
|
||||||
b.it = nil
|
b.it = nil
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue