diff --git a/miner/worker.go b/miner/worker.go index 25b91baa2303..95d18d50536f 100644 --- a/miner/worker.go +++ b/miner/worker.go @@ -144,7 +144,6 @@ type worker struct { chainHeadSub event.Subscription chainSideCh chan core.ChainSideEvent chainSideSub event.Subscription - resetCh chan time.Duration // Channel to request timer resets wg sync.WaitGroup @@ -194,7 +193,6 @@ func newWorker(config *Config, chainConfig *params.ChainConfig, engine consensus txsCh: make(chan core.NewTxsEvent, txChanSize), chainHeadCh: make(chan core.ChainHeadEvent, chainHeadChanSize), chainSideCh: make(chan core.ChainSideEvent, chainSideChanSize), - resetCh: make(chan time.Duration, 1), chainDb: eth.ChainDb(), recv: make(chan *Result, resultQueueSize), chain: eth.BlockChain(), @@ -387,58 +385,46 @@ func (w *worker) update() { timeout := time.NewTimer(time.Duration(minePeriod) * time.Second) defer timeout.Stop() - c := make(chan struct{}, 1) - defer close(c) - finish := make(chan struct{}) - defer close(finish) - - go func() { - for { - // A real event arrived, process interesting content + + // resetTimer rearms the mining timer. It must only ever be called from the + // loop below. Driving the timer from a dedicated goroutine would require a + // two-way channel handshake, which deadlocks as soon as both directions are + // full: the loop blocks handing over the new duration while the helper + // blocks handing over the expiry notification. + resetTimer := func(d time.Duration) { + if !timeout.Stop() { + // Drain the timer channel if it had already expired. select { - case d := <-w.resetCh: - // Reset the timer to the new duration. - if !timeout.Stop() { - // Drain the timer channel if it had already expired. - select { - case <-timeout.C: - default: - } - } - timeout.Reset(d) case <-timeout.C: - c <- struct{}{} - case <-finish: - return + default: } } - }() + timeout.Reset(d) + } + for { // A real event arrived, process interesting content select { case v := <-minePeriodCh: log.Info("[worker] update wait period", "period", v) minePeriod = v - w.resetCh <- time.Duration(minePeriod) * time.Second + resetTimer(time.Duration(minePeriod) * time.Second) - case <-c: + case <-timeout.C: if atomic.LoadInt32(&w.mining) == 1 { w.commitNewWork() } - resetTime := getResetTime(w.chain, minePeriod) - w.resetCh <- resetTime + resetTimer(getResetTime(w.chain, minePeriod)) // Handle ChainHeadEvent case <-w.chainHeadCh: w.commitNewWork() - resetTime := getResetTime(w.chain, minePeriod) - w.resetCh <- resetTime + resetTimer(getResetTime(w.chain, minePeriod)) // Handle new round case <-newRoundCh: w.commitNewWork() - resetTime := getResetTime(w.chain, minePeriod) - w.resetCh <- resetTime + resetTimer(getResetTime(w.chain, minePeriod)) // Handle ChainSideEvent case <-w.chainSideCh: diff --git a/miner/worker_test.go b/miner/worker_test.go index 82c92cf9844d..6a137deafaf0 100644 --- a/miner/worker_test.go +++ b/miner/worker_test.go @@ -50,7 +50,6 @@ func TestWorkerUpdateNonXDPoSStaysRunning(t *testing.T) { engine: ethash.NewFaker(), chainHeadSub: newBlockingSubscription(), chainSideSub: newBlockingSubscription(), - resetCh: make(chan time.Duration, 1), } done := make(chan struct{}) @@ -83,6 +82,79 @@ func TestWorkerUpdateNonXDPoSStaysRunning(t *testing.T) { } } +// TestWorkerUpdateKeepsDrainingChainHead ensures the mining timer never stops +// the update loop from servicing chain events. +// +// The timer used to be owned by a dedicated goroutine that exchanged +// notifications with the update loop over two buffered channels: the loop sent +// the next duration on resetCh and the goroutine reported expiries on c. Once +// both buffers were full the two goroutines blocked on each other forever. The +// worker then stopped draining chainHeadCh, which in turn blocked every +// producer of chain events (block insertion posts them synchronously) and +// wedged the whole node. +func TestWorkerUpdateKeepsDrainingChainHead(t *testing.T) { + chainConfig := ¶ms.ChainConfig{ + ChainID: big.NewInt(1338), + HomesteadBlock: new(big.Int), + Ethash: new(params.EthashConfig), + } + genesis := &core.Genesis{ + Config: chainConfig, + // Anchor the genesis at the current time so getResetTime returns a + // steadily shrinking duration and the timer fires repeatedly while the + // loop is busy handling chain head events. + Timestamp: uint64(time.Now().Unix()), + Difficulty: big.NewInt(1), + GasLimit: params.XDCGenesisGasLimit, + } + db := rawdb.NewMemoryDatabase() + engine := ethash.NewFaker() + chain, err := core.NewBlockChain(db, nil, genesis, engine, vm.Config{}) + if err != nil { + t.Fatalf("failed to create blockchain: %v", err) + } + defer chain.Stop() + + head := chain.GetBlockByNumber(0) + if head == nil { + t.Fatal("expected genesis block") + } + + // announceTxs is false and mining is 0, so commitNewWork bails out early in + // checkPreCommit and the loop stays cheap. + worker := &worker{ + chainConfig: chainConfig, + engine: engine, + chain: chain, + chainHeadCh: make(chan core.ChainHeadEvent, chainHeadChanSize), + chainSideCh: make(chan core.ChainSideEvent, chainSideChanSize), + chainHeadSub: newBlockingSubscription(), + chainSideSub: newBlockingSubscription(), + } + + done := make(chan struct{}) + go func() { + worker.update() + close(done) + }() + + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + select { + case worker.chainHeadCh <- core.ChainHeadEvent{Block: head}: + case <-time.After(time.Second): + t.Fatal("worker.update stopped draining chainHeadCh") + } + } + + worker.chainHeadSub.Unsubscribe() + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("worker.update did not return after unsubscribe") + } +} + // TestWorkerUpdateNewTxsWithoutTRC21Issuer tests worker update new txs without trc 21 issuer. func TestWorkerUpdateNewTxsWithoutTRC21Issuer(t *testing.T) { key, err := crypto.HexToECDSA("b71c71a67e1177ad4e901695e1b4b9ee17ae16c6668d313eac2f96dbcda3f291")