swarm/network/stream: fix a goroutine leak in Registry

This commit is contained in:
Janos Guljas 2019-02-19 11:35:19 +01:00
parent c283d9b5e8
commit fa785341bf

View file

@ -95,6 +95,7 @@ type Registry struct {
spec *protocols.Spec //this protocol's spec spec *protocols.Spec //this protocol's spec
balance protocols.Balance //implements protocols.Balance, for accounting balance protocols.Balance //implements protocols.Balance, for accounting
prices protocols.Prices //implements protocols.Prices, provides prices to accounting prices protocols.Prices //implements protocols.Prices, provides prices to accounting
close chan struct{} // terminates registry goroutines
} }
// RegistryOptions holds optional values for NewRegistry constructor. // RegistryOptions holds optional values for NewRegistry constructor.
@ -117,6 +118,8 @@ func NewRegistry(localID enode.ID, delivery *Delivery, syncChunkStore storage.Sy
// check if retrieval has been disabled // check if retrieval has been disabled
retrieval := options.Retrieval != RetrievalDisabled retrieval := options.Retrieval != RetrievalDisabled
closeChan := make(chan struct{})
streamer := &Registry{ streamer := &Registry{
addr: localID, addr: localID,
skipCheck: options.SkipCheck, skipCheck: options.SkipCheck,
@ -128,6 +131,7 @@ func NewRegistry(localID enode.ID, delivery *Delivery, syncChunkStore storage.Sy
autoRetrieval: retrieval, autoRetrieval: retrieval,
maxPeerServers: options.MaxPeerServers, maxPeerServers: options.MaxPeerServers,
balance: balance, balance: balance,
close: closeChan,
} }
streamer.setupSpec() streamer.setupSpec()
@ -172,12 +176,17 @@ func NewRegistry(localID enode.ID, delivery *Delivery, syncChunkStore storage.Sy
go func() { go func() {
defer close(out) defer close(out)
for i := range in { for {
select { select {
case <-out: case i := <-in:
default: select {
case <-out:
default:
}
out <- i
case <-closeChan:
return
} }
out <- i
} }
}() }()
@ -229,6 +238,8 @@ func NewRegistry(localID enode.ID, delivery *Delivery, syncChunkStore storage.Sy
<-timer.C <-timer.C
} }
timer.Reset(options.SyncUpdateDelay) timer.Reset(options.SyncUpdateDelay)
case <-closeChan:
break loop
} }
} }
timer.Stop() timer.Stop()
@ -398,6 +409,7 @@ func (r *Registry) Quit(peerId enode.ID, s Stream) error {
} }
func (r *Registry) Close() error { func (r *Registry) Close() error {
close(r.close)
return r.intervalsStore.Close() return r.intervalsStore.Close()
} }