From 5a771ed6f05d14f8389ace343aa627f528fa445b Mon Sep 17 00:00:00 2001 From: Tuna Date: Thu, 12 Jul 2018 16:41:25 +0700 Subject: [PATCH 1/2] node should also count to his turn in case some nodes are down --- consensus/clique/clique.go | 14 +++++++----- core/blockchain.go | 3 ++- miner/worker.go | 45 +++++++++++++++++++++++++++++++++++--- 3 files changed, 53 insertions(+), 9 deletions(-) diff --git a/consensus/clique/clique.go b/consensus/clique/clique.go index 4e57afb35c..8546b8cd0f 100644 --- a/consensus/clique/clique.go +++ b/consensus/clique/clique.go @@ -48,7 +48,6 @@ const ( inmemorySnapshots = 128 // Number of recent vote snapshots to keep in memory inmemorySignatures = 4096 // Number of recent block signatures to keep in memory wiggleTime = 500 * time.Millisecond // Random delay (per signer) to allow concurrent signers - ) type Masternode struct { @@ -405,7 +404,9 @@ func (c *Clique) GetMasternodes(chain consensus.ChainReader, header *types.Heade return masternodes } -func YourTurn(masternodes []common.Address, snap *Snapshot, header *types.Header, cur common.Address) (bool, error) { +func (c *Clique) GetPeriod() uint64 { return c.config.Period } + +func YourTurn(masternodes []common.Address, snap *Snapshot, header *types.Header, cur common.Address) (int, int, bool, error) { pre := common.Address{} // masternode[0] has chance to create block 1 var err error @@ -413,16 +414,19 @@ func YourTurn(masternodes []common.Address, snap *Snapshot, header *types.Header if header.Number.Uint64() != 0 { pre, err = ecrecover(header, snap.sigcache) if err != nil { - return false, err + return 0, 0, false, err } preIndex = position(masternodes, pre) } curIndex := position(masternodes, cur) - log.Info("Debugging info", "number of masternodes", len(masternodes), "previous", pre, "position", preIndex, "current", cur, "position", curIndex) + log.Info("Masternodes cycle info", "number of masternodes", len(masternodes), "previous", pre, "position", preIndex, "current", cur, "position", curIndex) for i, s := range masternodes { fmt.Printf("%d - %s\n", i, s.String()) } - return (preIndex+1)%len(masternodes) == curIndex, nil + if (preIndex+1)%len(masternodes) == curIndex { + return preIndex, curIndex, true, nil + } + return preIndex, curIndex, false, nil } // snapshot retrieves the authorization snapshot at a given point in time. diff --git a/core/blockchain.go b/core/blockchain.go index 96039a0896..1656ac3c8d 100644 --- a/core/blockchain.go +++ b/core/blockchain.go @@ -49,6 +49,7 @@ var ( blockInsertTimer = metrics.NewRegisteredTimer("chain/inserts", nil) CheckpointCh = make(chan int) M1Ch = make(chan int) + NewBlockCh = make(chan int) ErrNoGenesis = errors.New("Genesis not found in chain") ) @@ -1243,7 +1244,7 @@ func (st *insertStats) report(chain []*types.Block, index int, cache common.Stor context = append(context, []interface{}{"ignored", st.ignored}...) } log.Info("Imported new chain segment", context...) - + NewBlockCh <- 1 *st = insertStats{startTime: now, lastIndex: index + 1} } } diff --git a/miner/worker.go b/miner/worker.go index ae0aeba029..ed78d5f192 100644 --- a/miner/worker.go +++ b/miner/worker.go @@ -51,6 +51,8 @@ const ( chainHeadChanSize = 10 // chainSideChanSize is the size of channel listening to ChainSideEvent. chainSideChanSize = 10 + // Timeout waiting for M1 + m1Timeout = 1000 ) // Agent can register themself with the worker @@ -415,6 +417,24 @@ func (self *worker) makeCurrent(parent *types.Block, header *types.Header) error return nil } +func abs(x int64) int64 { + if x < 0 { + return -x + } + return x +} + +func hop(len, pre, cur int) int { + switch { + case pre < cur: + return cur - (pre + 1) + case pre > cur: + return (len - pre) + (cur - 1) + default: + return len - 1 + } +} + func (self *worker) commitNewWork() { self.mu.Lock() defer self.mu.Unlock() @@ -439,14 +459,33 @@ func (self *worker) commitNewWork() { log.Error("Failed when trying to commit new work", "err", err) return } - ok, err := clique.YourTurn(masternodes, snap, parent.Header(), self.coinbase) + preIndex, curIndex, ok, err := clique.YourTurn(masternodes, snap, parent.Header(), self.coinbase) if err != nil { log.Error("Failed when trying to commit new work", "err", err) return } if !ok { - log.Info("Not our turn to commit block. Wait for next time") - return + log.Info("Not my turn to commit block. Waiting...") + // in case some nodes are down + if preIndex == -1 { + // first block + return + } + h := hop(len(masternodes), preIndex, curIndex) + gap := int64(c.GetPeriod()) * int64(h) + log.Info("Distance from the parent block", "seconds", gap, "hops", h) + L: + for { + select { + case <-core.NewBlockCh: + log.Info("New block has came already. Skip this turn") + return + case <-time.After(time.Duration(gap+m1Timeout) * time.Second): + // wait enough. It's my turn + log.Info("Wait enough. It's my turn", "waited seconds", gap+m1Timeout) + break L + } + } } } } From f65014512d5efcf2a0c2a3a5f0b5022eb95e8a91 Mon Sep 17 00:00:00 2001 From: Tuna Date: Thu, 19 Jul 2018 17:15:21 +0700 Subject: [PATCH 2/2] node waits to his turn until there is a new block comes in --- core/blockchain.go | 2 -- eth/backend.go | 4 ++-- miner/worker.go | 19 ++++++++++--------- 3 files changed, 12 insertions(+), 13 deletions(-) diff --git a/core/blockchain.go b/core/blockchain.go index 1656ac3c8d..9766ad9021 100644 --- a/core/blockchain.go +++ b/core/blockchain.go @@ -49,7 +49,6 @@ var ( blockInsertTimer = metrics.NewRegisteredTimer("chain/inserts", nil) CheckpointCh = make(chan int) M1Ch = make(chan int) - NewBlockCh = make(chan int) ErrNoGenesis = errors.New("Genesis not found in chain") ) @@ -1244,7 +1243,6 @@ func (st *insertStats) report(chain []*types.Block, index int, cache common.Stor context = append(context, []interface{}{"ignored", st.ignored}...) } log.Info("Imported new chain segment", context...) - NewBlockCh <- 1 *st = insertStats{startTime: now, lastIndex: index + 1} } } diff --git a/eth/backend.go b/eth/backend.go index c8247bdab3..7186462ca9 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -471,8 +471,8 @@ func (s *Ethereum) StartStaking(local bool) error { return nil } -func (s *Ethereum) StopStaking() { s.miner.Stop() } -func (s *Ethereum) IsStaking() bool { return s.miner.Mining() } +func (s *Ethereum) StopStaking() { s.miner.Stop() } +func (s *Ethereum) IsStaking() bool { return s.miner.Mining() } func (s *Ethereum) Miner() *miner.Miner { return s.miner } func (s *Ethereum) AccountManager() *accounts.Manager { return s.accountManager } diff --git a/miner/worker.go b/miner/worker.go index ed78d5f192..5fef89d736 100644 --- a/miner/worker.go +++ b/miner/worker.go @@ -52,7 +52,7 @@ const ( // chainSideChanSize is the size of channel listening to ChainSideEvent. chainSideChanSize = 10 // Timeout waiting for M1 - m1Timeout = 1000 + m1Timeout = 1 ) // Agent can register themself with the worker @@ -475,16 +475,17 @@ func (self *worker) commitNewWork() { gap := int64(c.GetPeriod()) * int64(h) log.Info("Distance from the parent block", "seconds", gap, "hops", h) L: - for { - select { - case <-core.NewBlockCh: - log.Info("New block has came already. Skip this turn") + select { + case newBlock := <-self.chainHeadCh: + if newBlock.Block.NumberU64() > parent.NumberU64() { + log.Info("New block has came already. Skip this turn", "new block", newBlock.Block.NumberU64(), "current block", parent.NumberU64()) + self.chainHeadCh <- newBlock return - case <-time.After(time.Duration(gap+m1Timeout) * time.Second): - // wait enough. It's my turn - log.Info("Wait enough. It's my turn", "waited seconds", gap+m1Timeout) - break L } + case <-time.After(time.Duration(gap+m1Timeout) * time.Second): + // wait enough. It's my turn + log.Info("Wait enough. It's my turn", "waited seconds", gap+m1Timeout) + break L } } }