Refactor dispute manager to new blockchain architecture

This commit is contained in:
ChronosX88 2021-07-11 03:08:03 +03:00
parent 956de703e8
commit 5aadbab600
Signed by: ChronosXYZ
GPG Key ID: 085A69A82C8C511A
4 changed files with 146 additions and 97 deletions

View File

@ -2,8 +2,19 @@ package consensus
import (
"context"
"encoding/hex"
"time"
types2 "github.com/Secured-Finance/dione/blockchain/types"
"github.com/Secured-Finance/dione/types"
"github.com/fxamacker/cbor/v2"
"github.com/sirupsen/logrus"
"golang.org/x/crypto/sha3"
"github.com/Secured-Finance/dione/blockchain"
"github.com/Secured-Finance/dione/contracts/dioneDispute"
"github.com/Secured-Finance/dione/contracts/dioneOracle"
"github.com/Secured-Finance/dione/ethclient"
@ -16,9 +27,10 @@ type DisputeManager struct {
submissionMap map[string]*dioneOracle.DioneOracleSubmittedOracleRequest
disputeMap map[string]*dioneDispute.DioneDisputeNewDispute
voteWindow time.Duration
blockchain *blockchain.BlockChain
}
func NewDisputeManager(ctx context.Context, ethClient *ethclient.EthereumClient, pcm *PBFTConsensusManager, voteWindow int) (*DisputeManager, error) {
func NewDisputeManager(ctx context.Context, ethClient *ethclient.EthereumClient, pcm *PBFTConsensusManager, voteWindow int, bc *blockchain.BlockChain) (*DisputeManager, error) {
newSubmittionsChan, submSubscription, err := ethClient.SubscribeOnNewSubmittions(ctx)
if err != nil {
return nil, err
@ -36,6 +48,7 @@ func NewDisputeManager(ctx context.Context, ethClient *ethclient.EthereumClient,
submissionMap: map[string]*dioneOracle.DioneOracleSubmittedOracleRequest{},
disputeMap: map[string]*dioneDispute.DioneDisputeNewDispute{},
voteWindow: time.Duration(voteWindow) * time.Second,
blockchain: bc,
}
go func() {
@ -62,95 +75,117 @@ func NewDisputeManager(ctx context.Context, ethClient *ethclient.EthereumClient,
return dm, nil
}
func (dm *DisputeManager) onNewSubmission(submittion *dioneOracle.DioneOracleSubmittedOracleRequest) {
//c := dm.pcm.GetConsensusInfo(submittion.ReqID.String())
//if c == nil {
// // todo: warn
// return
//}
//
//dm.submissionMap[submittion.ReqID.String()] = submittion
//
//submHashBytes := sha3.Sum256(submittion.Data)
//localHashBytes := sha3.Sum256(c.Task.Payload)
//submHash := hex.EncodeToString(submHashBytes[:])
//localHash := hex.EncodeToString(localHashBytes[:])
//if submHash != localHash {
// logrus.Debugf("submission of request id %s isn't valid - beginning dispute", c.Task.RequestID)
// addr := common.HexToAddress(c.Task.MinerEth)
// reqID, ok := big.NewInt(0).SetString(c.Task.RequestID, 10)
// if !ok {
// logrus.Errorf("cannot parse request id: %s", c.Task.RequestID)
// return
// }
// err := dm.ethClient.BeginDispute(addr, reqID)
// if err != nil {
// logrus.Errorf(err.Error())
// return
// }
// disputeFinishTimer := time.NewTimer(dm.voteWindow)
// go func() {
// for {
// select {
// case <-dm.ctx.Done():
// return
// case <-disputeFinishTimer.C:
// {
// d, ok := dm.disputeMap[reqID.String()]
// if !ok {
// logrus.Error("cannot finish dispute: it doesn't exist in manager's dispute map!")
// return
// }
// err := dm.ethClient.FinishDispute(d.Dhash)
// if err != nil {
// logrus.Errorf(err.Error())
// return
// }
// disputeFinishTimer.Stop()
// return
// }
// }
// }
// }()
//}
// TODO refactor due to new architecture with blockchain
func (dm *DisputeManager) onNewSubmission(submission *dioneOracle.DioneOracleSubmittedOracleRequest) {
// find a block that contains the dione task with specified request id
task, block, err := dm.findTaskAndBlockWithRequestID(submission.ReqID.String())
if err != nil {
logrus.Error(err)
return
}
dm.submissionMap[submission.ReqID.String()] = submission
submHashBytes := sha3.Sum256(submission.Data)
localHashBytes := sha3.Sum256(task.Payload)
submHash := hex.EncodeToString(submHashBytes[:])
localHash := hex.EncodeToString(localHashBytes[:])
if submHash != localHash {
logrus.Debugf("submission of request id %s isn't valid - beginning dispute", submission.ReqID)
err := dm.ethClient.BeginDispute(block.Header.ProposerEth, submission.ReqID)
if err != nil {
logrus.Errorf(err.Error())
return
}
disputeFinishTimer := time.NewTimer(dm.voteWindow)
go func() {
for {
select {
case <-dm.ctx.Done():
return
case <-disputeFinishTimer.C:
{
d, ok := dm.disputeMap[submission.ReqID.String()]
if !ok {
logrus.Error("cannot finish dispute: it doesn't exist in manager's dispute map!")
return
}
err := dm.ethClient.FinishDispute(d.Dhash)
if err != nil {
logrus.Errorf(err.Error())
return
}
disputeFinishTimer.Stop()
return
}
}
}
}()
}
}
func (dm *DisputeManager) findTaskAndBlockWithRequestID(requestID string) (*types.DioneTask, *types2.Block, error) {
height, err := dm.blockchain.GetLatestBlockHeight()
if err != nil {
return nil, nil, err
}
for {
block, err := dm.blockchain.FetchBlockByHeight(height)
if err != nil {
return nil, nil, err
}
for _, v := range block.Data {
var task types.DioneTask
err := cbor.Unmarshal(v.Data, &task)
if err != nil {
logrus.Error(err)
continue
}
if task.RequestID == requestID {
return &task, block, nil
}
}
height--
}
}
func (dm *DisputeManager) onNewDispute(dispute *dioneDispute.DioneDisputeNewDispute) {
//c := dm.pcm.GetConsensusInfo(dispute.RequestID.String())
//if c == nil {
// // todo: warn
// return
//}
//
//subm, ok := dm.submissionMap[dispute.RequestID.String()]
//if !ok {
// // todo: warn
// return
//}
//
//dm.disputeMap[dispute.RequestID.String()] = dispute
//
//if dispute.DisputeInitiator.Hex() == dm.ethClient.GetEthAddress().Hex() {
// return
//}
//
//submHashBytes := sha3.Sum256(subm.Data)
//localHashBytes := sha3.Sum256(c.Task.Payload)
//submHash := hex.EncodeToString(submHashBytes[:])
//localHash := hex.EncodeToString(localHashBytes[:])
//if submHash == localHash {
// err := dm.ethClient.VoteDispute(dispute.Dhash, false)
// if err != nil {
// logrus.Errorf(err.Error())
// return
// }
//}
//
//err := dm.ethClient.VoteDispute(dispute.Dhash, true)
//if err != nil {
// logrus.Errorf(err.Error())
// return
//}
// TODO refactor due to new architecture with blockchain
task, _, err := dm.findTaskAndBlockWithRequestID(dispute.RequestID.String())
if err != nil {
logrus.Error(err)
return
}
subm, ok := dm.submissionMap[dispute.RequestID.String()]
if !ok {
logrus.Warn("desired submission isn't found in map")
return
}
dm.disputeMap[dispute.RequestID.String()] = dispute
if dispute.DisputeInitiator.Hex() == dm.ethClient.GetEthAddress().Hex() {
return
}
submHashBytes := sha3.Sum256(subm.Data)
localHashBytes := sha3.Sum256(task.Payload)
submHash := hex.EncodeToString(submHashBytes[:])
localHash := hex.EncodeToString(localHashBytes[:])
if submHash == localHash {
err := dm.ethClient.VoteDispute(dispute.Dhash, false)
if err != nil {
logrus.Errorf(err.Error())
return
}
}
err = dm.ethClient.VoteDispute(dispute.Dhash, true)
if err != nil {
logrus.Errorf(err.Error())
return
}
}

View File

@ -9,6 +9,8 @@ import (
"os"
"time"
"github.com/multiformats/go-multiaddr"
"github.com/asaskevich/EventBus"
"github.com/fxamacker/cbor/v2"
@ -160,7 +162,14 @@ func NewNode(config *config.Config, prvKey crypto.PrivKey, pexDiscoveryUpdateTim
r := provideP2PRPCClient(lhost)
// initialize sync manager
sm, err := provideSyncManager(bus, bc, mp, r, baddrs[0], psb) // FIXME here we just pick up first bootstrap in list
var baddr multiaddr.Multiaddr
if len(baddrs) == 0 {
baddr = nil
} else {
baddr = baddrs[0]
}
sm, err := provideSyncManager(bus, bc, mp, r, baddr, psb) // FIXME here we just pick up first bootstrap in list
if err != nil {
logrus.Fatal(err)
}
@ -186,7 +195,7 @@ func NewNode(config *config.Config, prvKey crypto.PrivKey, pexDiscoveryUpdateTim
logrus.Info("Random beacon subsystem has been initialized!")
// initialize dispute subsystem
disputeManager, err := provideDisputeManager(context.TODO(), ethClient, consensusManager, config)
disputeManager, err := provideDisputeManager(context.TODO(), ethClient, consensusManager, config, bc)
if err != nil {
logrus.Fatal(err)
}

View File

@ -56,8 +56,8 @@ func provideCache(config *config.Config) cache.Cache {
return backend
}
func provideDisputeManager(ctx context.Context, ethClient *ethclient.EthereumClient, pcm *consensus.PBFTConsensusManager, cfg *config.Config) (*consensus.DisputeManager, error) {
return consensus.NewDisputeManager(ctx, ethClient, pcm, cfg.Ethereum.DisputeVoteWindow)
func provideDisputeManager(ctx context.Context, ethClient *ethclient.EthereumClient, pcm *consensus.PBFTConsensusManager, cfg *config.Config, bc *blockchain.BlockChain) (*consensus.DisputeManager, error) {
return consensus.NewDisputeManager(ctx, ethClient, pcm, cfg.Ethereum.DisputeVoteWindow, bc)
}
func provideMiner(peerID peer.ID, ethAddress common.Address, ethClient *ethclient.EthereumClient, privateKey crypto.PrivKey, mempool *pool.Mempool) *consensus.Miner {
@ -165,11 +165,16 @@ func provideMemPool() (*pool.Mempool, error) {
}
func provideSyncManager(bus EventBus.Bus, bp *blockchain.BlockChain, mp *pool.Mempool, r *gorpc.Client, bootstrap multiaddr.Multiaddr, psb *pubsub.PubSubRouter) (sync.SyncManager, error) {
bootstrapPeerID := peer.ID("")
if bootstrap != nil {
addr, err := peer.AddrInfoFromP2pAddr(bootstrap)
if err != nil {
return nil, err
}
return sync.NewSyncManager(bus, bp, mp, r, addr.ID, psb), nil
bootstrapPeerID = addr.ID
}
return sync.NewSyncManager(bus, bp, mp, r, bootstrapPeerID, psb), nil
}
func provideP2PRPCClient(h host.Host) *gorpc.Client {