From 6259641448800a3e88c07808b0fa01cc5039d73d Mon Sep 17 00:00:00 2001 From: Tuna Date: Tue, 9 Oct 2018 14:22:23 +0700 Subject: [PATCH 1/4] init double validation layer --- consensus/posv/posv.go | 13 +++++++++++- contracts/utils.go | 25 +++++++++------------- core/blockchain.go | 8 +++---- core/tx_pool.go | 8 +++++++ eth/backend.go | 48 +++++++++++++++++++++++++++++++++++++++++- eth/fetcher/fetcher.go | 2 +- 6 files changed, 82 insertions(+), 22 deletions(-) diff --git a/consensus/posv/posv.go b/consensus/posv/posv.go index 40e1381845..2cd0b4cdda 100644 --- a/consensus/posv/posv.go +++ b/consensus/posv/posv.go @@ -414,6 +414,17 @@ func (c *Posv) GetMasternodes(chain consensus.ChainReader, header *types.Header) func (c *Posv) GetPeriod() uint64 { return c.config.Period } +func WhoIsCreator(snap *Snapshot, header *types.Header) (common.Address, error) { + if header.Number.Uint64() == 0 { + return common.Address{}, errors.New("Don't take block 0") + } + m, err := ecrecover(header, snap.sigcache) + if err != nil { + return common.Address{}, err + } + return m, nil +} + func YourTurn(masternodes []common.Address, snap *Snapshot, header *types.Header, cur common.Address) (int, int, bool, error) { if len(masternodes) == 0 { return -1, -1, true, nil @@ -423,7 +434,7 @@ func YourTurn(masternodes []common.Address, snap *Snapshot, header *types.Header var err error preIndex := -1 if header.Number.Uint64() != 0 { - pre, err = ecrecover(header, snap.sigcache) + pre, err = WhoIsCreator(snap, header) if err != nil { return 0, 0, false, err } diff --git a/contracts/utils.go b/contracts/utils.go index 435f6d668d..5f6bc589c5 100644 --- a/contracts/utils.go +++ b/contracts/utils.go @@ -272,6 +272,7 @@ func ExtractValidatorsFromBytes(byteValidators []byte) []int64 { intNumber, err := strconv.Atoi(string(trimByte)) if err != nil { log.Error("Can not convert string to integer", "error", err) + return []int64{} } validators = append(validators, int64(intNumber)) } @@ -568,26 +569,20 @@ func GetMasternodesFromCheckpointHeader(checkpointHeader *types.Header) []common } // Get m2 list from checkpoint block. -func GetM2FromCheckpointBlock(checkpointBlock types.Block) ([]common.Address, error) { +func GetM1M2FromCheckpointBlock(checkpointBlock *types.Block) (map[common.Address]common.Address, error) { if checkpointBlock.Number().Int64()%common.EpocBlockRandomize != 0 { return nil, errors.New("This block is not checkpoint block epoc.") } - - // Get singers from this block. + m1m2 := map[common.Address]common.Address{} + // Get signers from this block. masternodes := GetMasternodesFromCheckpointHeader(checkpointBlock.Header()) validators := ExtractValidatorsFromBytes(checkpointBlock.Header().Validators) - var m2List []common.Address - lenMasternodes := len(masternodes) - var valAddr common.Address - for validatorIndex := range validators { - if validatorIndex < lenMasternodes { - valAddr = masternodes[validatorIndex] - } else { - valAddr = masternodes[validatorIndex-lenMasternodes] - } - m2List = append(m2List, valAddr) + if len(validators) < len(masternodes) { + return nil, errors.New("Len(m2) is less than len(m1)") } - - return m2List, nil + for i, m1 := range masternodes { + m1m2[m1] = masternodes[validators[i]] + } + return m1m2, nil } diff --git a/core/blockchain.go b/core/blockchain.go index efd6c0758e..cf90d52177 100644 --- a/core/blockchain.go +++ b/core/blockchain.go @@ -1635,16 +1635,16 @@ func (bc *BlockChain) UpdateM1() error { ms = append(ms, posv.Masternode{Address: candidate, Stake: v.Uint64()}) } } - log.Info("Ordered list of masternode candidates") - for _, m := range ms { - log.Info("", "address", m.Address.String(), "stake", m.Stake) - } if len(ms) == 0 { log.Info("No masternode candidates found. Keep the current masternodes set for the next epoch") } else { sort.Slice(ms, func(i, j int) bool { return ms[i].Stake >= ms[j].Stake }) + log.Info("Ordered list of masternode candidates") + for _, m := range ms { + log.Info("", "address", m.Address.String(), "stake", m.Stake) + } // update masternodes log.Info("Updating new set of masternodes") if len(ms) > common.MaxMasternodes { diff --git a/core/tx_pool.go b/core/tx_pool.go index 09685405e3..571f700839 100644 --- a/core/tx_pool.go +++ b/core/tx_pool.go @@ -553,6 +553,14 @@ func (pool *TxPool) local() map[common.Address]types.Transactions { return txs } +func (pool *TxPool) GetSender(tx *types.Transaction) (common.Address, error) { + from, err := types.Sender(pool.signer, tx) + if err != nil { + return common.Address{}, ErrInvalidSender + } + return from, nil +} + // validateTx checks whether a transaction is valid according to the consensus // rules and adheres to some heuristic limits of the local node (price and size). func (pool *TxPool) validateTx(tx *types.Transaction, local bool) error { diff --git a/eth/backend.go b/eth/backend.go index 9499b7dea8..ff2cbf467f 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -51,6 +51,7 @@ import ( "github.com/ethereum/go-ethereum/params" "github.com/ethereum/go-ethereum/rlp" "github.com/ethereum/go-ethereum/rpc" + "time" ) const NumOfMasternodes = 99 @@ -201,10 +202,36 @@ func New(ctx *node.ServiceContext, config *Config) (*Ethereum, error) { return } if _, authorized := snap.Signers[eth.etherbase]; authorized { - if err := contracts.CreateTransactionSign(chainConfig, eth.txPool, eth.accountManager, block, chainDb); err != nil { + // double validation + m2, err := getM2(snap, eth, block) + if err != nil { + log.Error("Fail to validate M2 condition for imported block", "error", err) + return + } + if eth.etherbase != m2 { + //wait until signTx from m2 comes into txPool + txCh := make(chan core.TxPreEvent, txChanSize) + eth.txPool.SubscribeTxPreEvent(txCh) + G: + select { + case event := <-txCh: + from, err := eth.txPool.GetSender(event.Tx) + if (err == nil) && (event.Tx.To().String() == common.BlockSigners) && (from == m2) { + if err := contracts.CreateTransactionSign(chainConfig, eth.txPool, eth.accountManager, block, chainDb); err != nil { + log.Error("Fail to create tx sign for imported block", "error", err) + return + } + } + //timeout 10s + case <-time.After(time.Duration(10) * time.Second): + break G + } + close(txCh) + } else if err := contracts.CreateTransactionSign(chainConfig, eth.txPool, eth.accountManager, block, chainDb); err != nil { log.Error("Fail to create tx sign for imported block", "error", err) return } + // end of double validation } } eth.protocolManager.fetcher.SetImportedHook(importedHook) @@ -295,6 +322,25 @@ func New(ctx *node.ServiceContext, config *Config) (*Ethereum, error) { return eth, nil } +func getM2(snap *posv.Snapshot, eth *Ethereum, block *types.Block) (common.Address, error) { + epoch := eth.chainConfig.Posv.Epoch + no := block.NumberU64() + cpNo := no + if no%epoch != 0 { + cpNo = no - (no % epoch) + } + cpBlk := eth.blockchain.GetBlockByNumber(cpNo) + m, err := contracts.GetM1M2FromCheckpointBlock(cpBlk) + if err != nil { + return common.Address{}, err + } + m1, err := posv.WhoIsCreator(snap, block.Header()) + if err != nil { + return common.Address{}, err + } + return m[m1], nil +} + func makeExtraData(extra []byte) []byte { if len(extra) == 0 { // create default extradata diff --git a/eth/fetcher/fetcher.go b/eth/fetcher/fetcher.go index 4f0c916f7b..fefaf91c03 100644 --- a/eth/fetcher/fetcher.go +++ b/eth/fetcher/fetcher.go @@ -674,7 +674,7 @@ func (f *Fetcher) insert(peer string, block *types.Block) { propAnnounceOutTimer.UpdateSince(block.ReceivedAt) go f.broadcastBlock(block, false) - // Invoke the testing hook if needed + // Invoke the imported hook if needed if f.importedHook != nil { f.importedHook(block) } From 80cefd3e5573ba8bcd3cd233a9d972d5f355669d Mon Sep 17 00:00:00 2001 From: Tuna Date: Tue, 9 Oct 2018 15:33:58 +0700 Subject: [PATCH 2/4] lookup txPool before listening to imcoming tx --- eth/backend.go | 26 +++++++++++++++++++++++--- 1 file changed, 23 insertions(+), 3 deletions(-) diff --git a/eth/backend.go b/eth/backend.go index ff2cbf467f..fe7d386e5a 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -209,9 +209,28 @@ func New(ctx *node.ServiceContext, config *Config) (*Ethereum, error) { return } if eth.etherbase != m2 { - //wait until signTx from m2 comes into txPool + // firstly, look into txPool + pendingMap, err := eth.txPool.Pending() + if err != nil { + log.Error("Fail to get txPool pending", "err", err) + //reset pendingMap + pendingMap = map[common.Address]types.Transactions{} + } + txsSentFromM2 := pendingMap[m2] + if len(txsSentFromM2) > 0 { + for _, tx := range txsSentFromM2 { + if tx.To().String() == common.BlockSigners { + if err := contracts.CreateTransactionSign(chainConfig, eth.txPool, eth.accountManager, block, chainDb); err != nil { + log.Error("Fail to create tx sign for imported block", "error", err) + return + } + return + } + } + } + //then wait until signTx from m2 comes into txPool txCh := make(chan core.TxPreEvent, txChanSize) - eth.txPool.SubscribeTxPreEvent(txCh) + subEvent := eth.txPool.SubscribeTxPreEvent(txCh) G: select { case event := <-txCh: @@ -221,12 +240,13 @@ func New(ctx *node.ServiceContext, config *Config) (*Ethereum, error) { log.Error("Fail to create tx sign for imported block", "error", err) return } + return } //timeout 10s case <-time.After(time.Duration(10) * time.Second): break G } - close(txCh) + subEvent.Unsubscribe() } else if err := contracts.CreateTransactionSign(chainConfig, eth.txPool, eth.accountManager, block, chainDb); err != nil { log.Error("Fail to create tx sign for imported block", "error", err) return From 6d7cfe14d89e2e4df86cc18dc90785655f43b8dd Mon Sep 17 00:00:00 2001 From: Tuna Date: Tue, 9 Oct 2018 15:36:41 +0700 Subject: [PATCH 3/4] trim m2 set to fit m1 set --- contracts/utils.go | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/contracts/utils.go b/contracts/utils.go index 5f6bc589c5..6046baed85 100644 --- a/contracts/utils.go +++ b/contracts/utils.go @@ -579,10 +579,12 @@ func GetM1M2FromCheckpointBlock(checkpointBlock *types.Block) (map[common.Addres validators := ExtractValidatorsFromBytes(checkpointBlock.Header().Validators) if len(validators) < len(masternodes) { - return nil, errors.New("Len(m2) is less than len(m1)") + return nil, errors.New("len(m2) is less than len(m1)") } - for i, m1 := range masternodes { - m1m2[m1] = masternodes[validators[i]] + if len(masternodes) > 0 { + for i, m1 := range masternodes { + m1m2[m1] = masternodes[validators[i]%int64(len(masternodes))] + } } return m1m2, nil } From bd57df67ae6794ffb7fe2416c6eae32768259ed9 Mon Sep 17 00:00:00 2001 From: Tuna Date: Tue, 9 Oct 2018 18:03:31 +0700 Subject: [PATCH 4/4] tiny fix unitest --- eth/backend.go | 2 +- eth/protocol_test.go | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/eth/backend.go b/eth/backend.go index fe7d386e5a..bc3c7d2871 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -209,7 +209,7 @@ func New(ctx *node.ServiceContext, config *Config) (*Ethereum, error) { return } if eth.etherbase != m2 { - // firstly, look into txPool + // firstly, look into pending txPool pendingMap, err := eth.txPool.Pending() if err != nil { log.Error("Fail to get txPool pending", "err", err) diff --git a/eth/protocol_test.go b/eth/protocol_test.go index ed30f35bf6..d6ac52e0f4 100644 --- a/eth/protocol_test.go +++ b/eth/protocol_test.go @@ -63,8 +63,8 @@ func testStatusMsgErrors(t *testing.T, protocol int) { wantError: errResp(ErrProtocolVersionMismatch, "10 (!= %d)", protocol), }, { - code: StatusMsg, data: statusData{uint32(protocol), 89, td, head.Hash(), genesis.Hash()}, - wantError: errResp(ErrNetworkIdMismatch, "89 (!= 1)"), + code: StatusMsg, data: statusData{uint32(protocol), 999, td, head.Hash(), genesis.Hash()}, + wantError: errResp(ErrNetworkIdMismatch, "999 (!= 89)"), }, { code: StatusMsg, data: statusData{uint32(protocol), DefaultConfig.NetworkId, td, head.Hash(), common.Hash{3}},