eth : fix quit channel goroutine leak

This commit is contained in:
ucwong 2020-04-03 08:34:18 +00:00
parent f7b29ec942
commit c26a13aada
7 changed files with 11 additions and 11 deletions

View file

@ -148,7 +148,7 @@ func New(ctx *node.ServiceContext, config *Config) (*Ethereum, error) {
eventMux: ctx.EventMux, eventMux: ctx.EventMux,
accountManager: ctx.AccountManager, accountManager: ctx.AccountManager,
engine: CreateConsensusEngine(ctx, chainConfig, &config.Ethash, config.Miner.Notify, config.Miner.Noverify, chainDb), engine: CreateConsensusEngine(ctx, chainConfig, &config.Ethash, config.Miner.Notify, config.Miner.Noverify, chainDb),
closeBloomHandler: make(chan struct{}), closeBloomHandler: make(chan struct{}, 1),
networkID: config.NetworkId, networkID: config.NetworkId,
gasPrice: config.Miner.GasPrice, gasPrice: config.Miner.GasPrice,
etherbase: config.Miner.Etherbase, etherbase: config.Miner.Etherbase,

View file

@ -232,7 +232,7 @@ func New(checkpoint uint64, stateDb ethdb.Database, stateBloom *trie.SyncBloom,
bodyWakeCh: make(chan bool, 1), bodyWakeCh: make(chan bool, 1),
receiptWakeCh: make(chan bool, 1), receiptWakeCh: make(chan bool, 1),
headerProcCh: make(chan []*types.Header, 1), headerProcCh: make(chan []*types.Header, 1),
quitCh: make(chan struct{}), quitCh: make(chan struct{}, 1),
stateCh: make(chan dataPack), stateCh: make(chan dataPack),
stateSyncStart: make(chan *stateSync), stateSyncStart: make(chan *stateSync),
syncStatsState: stateSyncStats{ syncStatsState: stateSyncStats{
@ -394,7 +394,7 @@ func (d *Downloader) synchronise(id string, hash common.Hash, td *big.Int, mode
} }
// Create cancel channel for aborting mid-flight and mark the master peer // Create cancel channel for aborting mid-flight and mark the master peer
d.cancelLock.Lock() d.cancelLock.Lock()
d.cancelCh = make(chan struct{}) d.cancelCh = make(chan struct{}, 1)
d.cancelPeer = id d.cancelPeer = id
d.cancelLock.Unlock() d.cancelLock.Unlock()

View file

@ -244,8 +244,8 @@ func newStateSync(d *Downloader, root common.Hash) *stateSync {
keccak: sha3.NewLegacyKeccak256(), keccak: sha3.NewLegacyKeccak256(),
tasks: make(map[common.Hash]*stateTask), tasks: make(map[common.Hash]*stateTask),
deliver: make(chan *stateReq), deliver: make(chan *stateReq),
cancel: make(chan struct{}), cancel: make(chan struct{}, 1),
done: make(chan struct{}), done: make(chan struct{}, 1),
} }
} }

View file

@ -176,7 +176,7 @@ func NewBlockFetcher(getBlock blockRetrievalFn, verifyHeader headerVerifierFn, b
headerFilter: make(chan chan *headerFilterTask), headerFilter: make(chan chan *headerFilterTask),
bodyFilter: make(chan chan *bodyFilterTask), bodyFilter: make(chan chan *bodyFilterTask),
done: make(chan common.Hash), done: make(chan common.Hash),
quit: make(chan struct{}), quit: make(chan struct{}, 1),
announces: make(map[string]int), announces: make(map[string]int),
announced: make(map[common.Hash][]*blockAnnounce), announced: make(map[common.Hash][]*blockAnnounce),
fetching: make(map[common.Hash]*blockAnnounce), fetching: make(map[common.Hash]*blockAnnounce),

View file

@ -192,7 +192,7 @@ func NewTxFetcherForTests(
notify: make(chan *txAnnounce), notify: make(chan *txAnnounce),
cleanup: make(chan *txDelivery), cleanup: make(chan *txDelivery),
drop: make(chan *txDrop), drop: make(chan *txDrop),
quit: make(chan struct{}), quit: make(chan struct{}, 1),
waitlist: make(map[common.Hash]map[string]struct{}), waitlist: make(map[common.Hash]map[string]struct{}),
waittime: make(map[common.Hash]mclock.AbsTime), waittime: make(map[common.Hash]mclock.AbsTime),
waitslots: make(map[string]map[common.Hash]struct{}), waitslots: make(map[string]map[common.Hash]struct{}),

View file

@ -111,7 +111,7 @@ func NewProtocolManager(config *params.ChainConfig, checkpoint *params.TrustedCh
peers: newPeerSet(), peers: newPeerSet(),
whitelist: whitelist, whitelist: whitelist,
txsyncCh: make(chan *txsync), txsyncCh: make(chan *txsync),
quitSync: make(chan struct{}), quitSync: make(chan struct{}, 1),
} }
if mode == downloader.FullSync { if mode == downloader.FullSync {

View file

@ -122,7 +122,7 @@ func newPeer(version int, p *p2p.Peer, rw p2p.MsgReadWriter, getPooledTx func(ha
txBroadcast: make(chan []common.Hash), txBroadcast: make(chan []common.Hash),
txAnnounce: make(chan []common.Hash), txAnnounce: make(chan []common.Hash),
getPooledTx: getPooledTx, getPooledTx: getPooledTx,
term: make(chan struct{}), term: make(chan struct{}, 1),
} }
} }
@ -179,7 +179,7 @@ func (p *peer) broadcastTransactions() {
// If there's anything available to transfer, fire up an async writer // If there's anything available to transfer, fire up an async writer
if len(txs) > 0 { if len(txs) > 0 {
done = make(chan struct{}) done = make(chan struct{}, 1)
go func() { go func() {
if err := p.sendTransactions(txs); err != nil { if err := p.sendTransactions(txs); err != nil {
fail <- err fail <- err
@ -241,7 +241,7 @@ func (p *peer) announceTransactions() {
// If there's anything available to transfer, fire up an async writer // If there's anything available to transfer, fire up an async writer
if len(pending) > 0 { if len(pending) > 0 {
done = make(chan struct{}) done = make(chan struct{}, 1)
go func() { go func() {
if err := p.sendPooledTransactionHashes(pending); err != nil { if err := p.sendPooledTransactionHashes(pending); err != nil {
fail <- err fail <- err