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 + } + } } } }