|
|
@@ -6,27 +6,37 @@ import (
|
|
|
|
|
|
"github.com/ethereum/eth-go/ethchain"
|
|
|
"github.com/ethereum/eth-go/ethlog"
|
|
|
- "github.com/ethereum/eth-go/ethreact"
|
|
|
"github.com/ethereum/eth-go/ethwire"
|
|
|
+ "github.com/ethereum/eth-go/event"
|
|
|
)
|
|
|
|
|
|
var logger = ethlog.NewLogger("MINER")
|
|
|
|
|
|
type Miner struct {
|
|
|
- pow ethchain.PoW
|
|
|
- ethereum ethchain.EthManager
|
|
|
- coinbase []byte
|
|
|
- reactChan chan ethreact.Event
|
|
|
- txs ethchain.Transactions
|
|
|
- uncles []*ethchain.Block
|
|
|
- block *ethchain.Block
|
|
|
- powChan chan []byte
|
|
|
- powQuitChan chan ethreact.Event
|
|
|
- quitChan chan chan error
|
|
|
+ pow ethchain.PoW
|
|
|
+ ethereum ethchain.EthManager
|
|
|
+ coinbase []byte
|
|
|
+ txs ethchain.Transactions
|
|
|
+ uncles []*ethchain.Block
|
|
|
+ block *ethchain.Block
|
|
|
+
|
|
|
+ events event.Subscription
|
|
|
+ powQuitChan chan struct{}
|
|
|
+ powDone chan struct{}
|
|
|
|
|
|
turbo bool
|
|
|
}
|
|
|
|
|
|
+const (
|
|
|
+ Started = iota
|
|
|
+ Stopped
|
|
|
+)
|
|
|
+
|
|
|
+type Event struct {
|
|
|
+ Type int // Started || Stopped
|
|
|
+ Miner *Miner
|
|
|
+}
|
|
|
+
|
|
|
func (self *Miner) GetPow() ethchain.PoW {
|
|
|
return self.pow
|
|
|
}
|
|
|
@@ -48,46 +58,42 @@ func (self *Miner) ToggleTurbo() {
|
|
|
}
|
|
|
|
|
|
func (miner *Miner) Start() {
|
|
|
- miner.reactChan = make(chan ethreact.Event, 1) // This is the channel that receives 'updates' when ever a new transaction or block comes in
|
|
|
- miner.powChan = make(chan []byte, 1) // This is the channel that receives valid sha hashes for a given block
|
|
|
- miner.powQuitChan = make(chan ethreact.Event, 1) // This is the channel that can exit the miner thread
|
|
|
- miner.quitChan = make(chan chan error, 1)
|
|
|
|
|
|
// Insert initial TXs in our little miner 'pool'
|
|
|
miner.txs = miner.ethereum.TxPool().Flush()
|
|
|
miner.block = miner.ethereum.BlockChain().NewBlock(miner.coinbase)
|
|
|
|
|
|
+ mux := miner.ethereum.EventMux()
|
|
|
+ miner.events = mux.Subscribe(ethchain.NewBlockEvent{}, ethchain.TxEvent{})
|
|
|
+
|
|
|
// Prepare inital block
|
|
|
//miner.ethereum.StateManager().Prepare(miner.block.State(), miner.block.State())
|
|
|
go miner.listener()
|
|
|
|
|
|
- reactor := miner.ethereum.Reactor()
|
|
|
- reactor.Subscribe("newBlock", miner.reactChan)
|
|
|
- reactor.Subscribe("newTx:pre", miner.reactChan)
|
|
|
-
|
|
|
- // We need the quit chan to be a Reactor event.
|
|
|
- // The POW search method is actually blocking and if we don't
|
|
|
- // listen to the reactor events inside of the pow itself
|
|
|
- // The miner overseer will never get the reactor events themselves
|
|
|
- // Only after the miner will find the sha
|
|
|
- reactor.Subscribe("newBlock", miner.powQuitChan)
|
|
|
- reactor.Subscribe("newTx:pre", miner.powQuitChan)
|
|
|
-
|
|
|
logger.Infoln("Started")
|
|
|
+ mux.Post(Event{Started, miner})
|
|
|
+}
|
|
|
|
|
|
- reactor.Post("miner:start", miner)
|
|
|
+func (miner *Miner) Stop() {
|
|
|
+ logger.Infoln("Stopping...")
|
|
|
+ miner.events.Unsubscribe()
|
|
|
+ miner.ethereum.EventMux().Post(Event{Stopped, miner})
|
|
|
}
|
|
|
|
|
|
func (miner *Miner) listener() {
|
|
|
for {
|
|
|
+ miner.startMining()
|
|
|
+
|
|
|
select {
|
|
|
- case status := <-miner.quitChan:
|
|
|
- logger.Infoln("Stopped")
|
|
|
- status <- nil
|
|
|
- return
|
|
|
- case chanMessage := <-miner.reactChan:
|
|
|
+ case event, isopen := <-miner.events.Chan():
|
|
|
+ miner.stopMining()
|
|
|
+ if !isopen {
|
|
|
+ return
|
|
|
+ }
|
|
|
|
|
|
- if block, ok := chanMessage.Resource.(*ethchain.Block); ok {
|
|
|
+ switch event := event.(type) {
|
|
|
+ case ethchain.NewBlockEvent:
|
|
|
+ block := event.Block
|
|
|
//logger.Infoln("Got new block via Reactor")
|
|
|
if bytes.Compare(miner.ethereum.BlockChain().CurrentBlock.Hash(), block.Hash()) == 0 {
|
|
|
// TODO: Perhaps continue mining to get some uncle rewards
|
|
|
@@ -117,49 +123,44 @@ func (miner *Miner) listener() {
|
|
|
miner.uncles = append(miner.uncles, block)
|
|
|
}
|
|
|
}
|
|
|
- }
|
|
|
|
|
|
- if tx, ok := chanMessage.Resource.(*ethchain.Transaction); ok {
|
|
|
- found := false
|
|
|
- for _, ctx := range miner.txs {
|
|
|
- if found = bytes.Compare(ctx.Hash(), tx.Hash()) == 0; found {
|
|
|
- break
|
|
|
+ case ethchain.TxEvent:
|
|
|
+ if event.Type == ethchain.TxPre {
|
|
|
+ found := false
|
|
|
+ for _, ctx := range miner.txs {
|
|
|
+ if found = bytes.Compare(ctx.Hash(), event.Tx.Hash()) == 0; found {
|
|
|
+ break
|
|
|
+ }
|
|
|
+ }
|
|
|
+ if found == false {
|
|
|
+ // Undo all previous commits
|
|
|
+ miner.block.Undo()
|
|
|
+ // Apply new transactions
|
|
|
+ miner.txs = append(miner.txs, event.Tx)
|
|
|
}
|
|
|
-
|
|
|
- }
|
|
|
- if found == false {
|
|
|
- // Undo all previous commits
|
|
|
- miner.block.Undo()
|
|
|
- // Apply new transactions
|
|
|
- miner.txs = append(miner.txs, tx)
|
|
|
}
|
|
|
}
|
|
|
- default:
|
|
|
- miner.mineNewBlock()
|
|
|
+
|
|
|
+ case <-miner.powDone:
|
|
|
+ // next iteration will start mining again
|
|
|
}
|
|
|
}
|
|
|
}
|
|
|
|
|
|
-func (miner *Miner) Stop() {
|
|
|
- logger.Infoln("Stopping...")
|
|
|
-
|
|
|
- miner.powQuitChan <- ethreact.Event{}
|
|
|
-
|
|
|
- status := make(chan error)
|
|
|
- miner.quitChan <- status
|
|
|
- <-status
|
|
|
-
|
|
|
- reactor := miner.ethereum.Reactor()
|
|
|
- reactor.Unsubscribe("newBlock", miner.powQuitChan)
|
|
|
- reactor.Unsubscribe("newTx:pre", miner.powQuitChan)
|
|
|
- reactor.Unsubscribe("newBlock", miner.reactChan)
|
|
|
- reactor.Unsubscribe("newTx:pre", miner.reactChan)
|
|
|
+func (miner *Miner) startMining() {
|
|
|
+ if miner.powDone == nil {
|
|
|
+ miner.powDone = make(chan struct{})
|
|
|
+ }
|
|
|
+ miner.powQuitChan = make(chan struct{})
|
|
|
+ go miner.mineNewBlock()
|
|
|
+}
|
|
|
|
|
|
- reactor.Post("miner:stop", miner)
|
|
|
+func (miner *Miner) stopMining() {
|
|
|
+ close(miner.powQuitChan)
|
|
|
+ <-miner.powDone
|
|
|
}
|
|
|
|
|
|
func (self *Miner) mineNewBlock() {
|
|
|
-
|
|
|
stateManager := self.ethereum.StateManager()
|
|
|
|
|
|
self.block = self.ethereum.BlockChain().NewBlock(self.coinbase)
|
|
|
@@ -195,8 +196,9 @@ func (self *Miner) mineNewBlock() {
|
|
|
logger.Infof("Mining on block. Includes %v transactions", len(self.txs))
|
|
|
|
|
|
// Find a valid nonce
|
|
|
- self.block.Nonce = self.pow.Search(self.block, self.powQuitChan)
|
|
|
- if self.block.Nonce != nil {
|
|
|
+ nonce := self.pow.Search(self.block, self.powQuitChan)
|
|
|
+ if nonce != nil {
|
|
|
+ self.block.Nonce = nonce
|
|
|
err := self.ethereum.StateManager().Process(self.block, false)
|
|
|
if err != nil {
|
|
|
logger.Infoln(err)
|
|
|
@@ -208,4 +210,5 @@ func (self *Miner) mineNewBlock() {
|
|
|
self.txs = self.ethereum.TxPool().CurrentTransactions()
|
|
|
}
|
|
|
}
|
|
|
+ self.powDone <- struct{}{}
|
|
|
}
|