les: wg data race fixed

This commit is contained in:
Evgeny Danienko 2017-10-24 14:39:36 +03:00
parent 6d6a5a9337
commit 853d5d3fd0
No known key found for this signature in database
GPG key ID: BC8C34D8B45BECBF
2 changed files with 19 additions and 4 deletions

View file

@ -20,6 +20,7 @@ package les
import ( import (
"math/big" "math/big"
"sync" "sync"
"sync/atomic"
"time" "time"
"github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common"
@ -123,11 +124,15 @@ func newLightFetcher(pm *ProtocolManager) *lightFetcher {
// syncLoop is the main event loop of the light fetcher // syncLoop is the main event loop of the light fetcher
func (f *lightFetcher) syncLoop() { func (f *lightFetcher) syncLoop() {
f.pm.wg.Add(1) once := true
defer f.pm.wg.Done()
requesting := false requesting := false
for { for {
if once && atomic.LoadInt32(f.pm.isClosed) != closed {
once = false
f.pm.wg.Add(1)
defer f.pm.wg.Done()
}
select { select {
case <-f.pm.quitSync: case <-f.pm.quitSync:
return return

View file

@ -24,6 +24,7 @@ import (
"math/big" "math/big"
"net" "net"
"sync" "sync"
"sync/atomic"
"time" "time"
"github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common"
@ -115,6 +116,7 @@ type ProtocolManager struct {
// channels for fetcher, syncer, txsyncLoop // channels for fetcher, syncer, txsyncLoop
newPeerCh chan *peer newPeerCh chan *peer
isClosed *int32
quitSync chan struct{} quitSync chan struct{}
noMorePeers chan struct{} noMorePeers chan struct{}
@ -123,10 +125,16 @@ type ProtocolManager struct {
wg *sync.WaitGroup wg *sync.WaitGroup
} }
const (
started int32 = iota
closed
)
// NewProtocolManager returns a new ethereum sub protocol manager. The Ethereum sub protocol manages peers capable // NewProtocolManager returns a new ethereum sub protocol manager. The Ethereum sub protocol manages peers capable
// with the ethereum network. // with the ethereum network.
func NewProtocolManager(chainConfig *params.ChainConfig, lightSync bool, networkId uint64, mux *event.TypeMux, engine consensus.Engine, peers *peerSet, blockchain BlockChain, txpool txPool, chainDb ethdb.Database, odr *LesOdr, txrelay *LesTxRelay, quitSync chan struct{}, wg *sync.WaitGroup) (*ProtocolManager, error) { func NewProtocolManager(chainConfig *params.ChainConfig, lightSync bool, networkId uint64, mux *event.TypeMux, engine consensus.Engine, peers *peerSet, blockchain BlockChain, txpool txPool, chainDb ethdb.Database, odr *LesOdr, txrelay *LesTxRelay, quitSync chan struct{}, wg *sync.WaitGroup) (*ProtocolManager, error) {
// Create the protocol manager with the base fields // Create the protocol manager with the base fields
isClosed := started
manager := &ProtocolManager{ manager := &ProtocolManager{
lightSync: lightSync, lightSync: lightSync,
eventMux: mux, eventMux: mux,
@ -139,6 +147,7 @@ func NewProtocolManager(chainConfig *params.ChainConfig, lightSync bool, network
txrelay: txrelay, txrelay: txrelay,
peers: peers, peers: peers,
newPeerCh: make(chan *peer), newPeerCh: make(chan *peer),
isClosed: &isClosed,
quitSync: quitSync, quitSync: quitSync,
wg: wg, wg: wg,
noMorePeers: make(chan struct{}), noMorePeers: make(chan struct{}),
@ -242,7 +251,8 @@ func (pm *ProtocolManager) Stop() {
// will exit when they try to register. // will exit when they try to register.
pm.peers.Close() pm.peers.Close()
// Wait for any process action // Wait for any process action. Wait should be executed after all pm.wg.Add()
atomic.StoreInt32(pm.isClosed, closed)
pm.wg.Wait() pm.wg.Wait()
log.Info("Light Ethereum protocol stopped") log.Info("Light Ethereum protocol stopped")