Post

EthHash PoW Consensus - Block building and the filling up of transactions

The complete process on how a block gets built on geth.

EthHash PoW Consensus - Block building and the filling up of transactions

In my previous post https://pvnotpv.github.io/posts/ethhash/, it was regarding the algorithmic part of EthHash, but the whole process of creating a new block, adding transactions to it, and how the entire process works is pretty much a league of its own, and we’d be going deep into that.

I’m kind of really excited to write this post because this is exactly where we’d be going deep into the actual blockchain. I mean, we’d be going to see the actual incrementing of the block number from the parent and how new transactions are added to the chain!

image

So here’s a mapping I made of the entire mining process on geth:

image

To get a perfect version of Geth that doesn’t include many of the beacon chain methods, make sure to go with v1.10.14.

When going through such a huge codebase like Geth, with multiple packages like networking, tries, txpool, and state dbs, you’d wonder how the whole thing is connected, even more importantly, which is that one core thing that drives all of this. In simple words, how does the blockchain keep on running and keep producing blocks? This is exactly where the miner package comes along, and the whole process is really interesting.

Now the above diagram may seems like a lot but we’d be going into it section by section ;)

image

Let’s directly jump into the worker object which is the main struct which implements the whole mining process:

Worker

image

The worker is what implements the multiple loops, which’d keep on listening for events, transactions, mining, and creating new blocks. And the different channels on how these loops communicate with each other—I was really having a hard time figuring out the whole thing due to the multiple channels and the reason why I started making diagrams to not confuse with the whole thing.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
// worker is the main object which takes care of submitting new work to consensus engine
// and gathering the sealing result.
type worker struct {
	config      *Config
	chainConfig *params.ChainConfig
	engine      consensus.Engine
	eth         Backend
	chain       *core.BlockChain
	merger      *consensus.Merger

	// Feeds
	pendingLogsFeed event.Feed

	// Subscriptions
	mux          *event.TypeMux
	txsCh        chan core.NewTxsEvent
	txsSub       event.Subscription
	chainHeadCh  chan core.ChainHeadEvent
	chainHeadSub event.Subscription
	chainSideCh  chan core.ChainSideEvent
	chainSideSub event.Subscription

	// Channels
	newWorkCh          chan *newWorkReq
	taskCh             chan *task
	resultCh           chan *types.Block
	startCh            chan struct{}
	exitCh             chan struct{}
	resubmitIntervalCh chan time.Duration
	resubmitAdjustCh   chan *intervalAdjust

worker has a method named newWorker() which starts all the loops for the worker, which is of course a goroutine of its own and keeps on running doing certain tasks.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
	// Subscribe NewTxsEvent for tx pool
	worker.txsSub = eth.TxPool().SubscribeNewTxsEvent(worker.txsCh)
	// Subscribe events for blockchain
	worker.chainHeadSub = eth.BlockChain().SubscribeChainHeadEvent(worker.chainHeadCh)
	worker.chainSideSub = eth.BlockChain().SubscribeChainSideEvent(worker.chainSideCh)

	// Sanitize recommit interval if the user-specified one is too short.
	recommit := worker.config.Recommit
	if recommit < minRecommitInterval {
		log.Warn("Sanitizing miner recommit interval", "provided", recommit, "updated", minRecommitInterval)
		recommit = minRecommitInterval
	}

	worker.wg.Add(4)
	go worker.mainLoop()
	go worker.newWorkLoop(recommit)
	go worker.resultLoop()
	go worker.taskLoop()

	// Submit first work to initialize pending state.
	if init {
		worker.startCh <- struct{}{}
	}
	return worker

It subscribes to new transactions, head block events, and the side chain events are sent during reorgs.

image

So there are 4 loops doing their own thing, and the best loop to start with is the workLoop()

workLoop

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
// newWorkLoop is a standalone goroutine to submit new mining work upon received events.
func (w *worker) newWorkLoop(recommit time.Duration) {
	defer w.wg.Done()
	var (
...

	// commit aborts in-flight transaction execution with given signal and resubmits a new one.
	commit := func(noempty bool, s int32) {
		if interrupt != nil {
			atomic.StoreInt32(interrupt, s)
		}
		interrupt = new(int32)
		select {
		case w.newWorkCh <- &newWorkReq{interrupt: interrupt, noempty: noempty, timestamp: timestamp}:
		case <-w.exitCh:
			return
		}
		timer.Reset(recommit)
		atomic.StoreInt32(&w.newTxs, 0)
	}

...


	for {
		select {
		case <-w.startCh:
			clearPending(w.chain.CurrentBlock().NumberU64())
			timestamp = time.Now().Unix()
			commit(false, commitInterruptNewHead)

		case head := <-w.chainHeadCh:
			clearPending(head.Block.NumberU64())
			timestamp = time.Now().Unix()
			commit(false, commitInterruptNewHead)

		case <-timer.C:
			// If mining is running resubmit a new work cycle periodically to pull in
			// higher priced transactions. Disable this overhead for pending blocks.
			if w.isRunning() && (w.chainConfig.Clique == nil || w.chainConfig.Clique.Period > 0) {
				// Short circuit if no new transaction arrives.
				if atomic.LoadInt32(&w.newTxs) == 0 {
					timer.Reset(recommit)
					continue
				}
				commit(true, commitInterruptResubmit)
			}


(I’ve removed certain stuffs to focus on the main parts.)

image

From the code we can see that workLoop() is what does action based on events being received; when a new head block is received, it creates a new work request, or when something is sent to the start channel, a new work request is being sent.

At the end of the newWorker() method, we can see these lines:

1
2
3
4
	// Submit first work to initialize pending state.
	if init {
		worker.startCh <- struct{}{}
	}

So when a new worker is created , if init is set to true , the whole process starts right away.

image

1
2
3
4
5
6
// newWorkReq represents a request for new sealing work submitting with relative interrupt notifier.
type newWorkReq struct {
	interrupt *int32
	noempty   bool
	timestamp int64
}

With the interrupts being:

1
2
3
4
5
const (
	commitInterruptNone int32 = iota
	commitInterruptNewHead
	commitInterruptResubmit
)

The struct would start making a lot of sense once we get to the sealing part.

Currently the workLoop has received a request in the start channel and it sends a new work to the “newWork” channel

1
2
		case w.newWorkCh <- &newWorkReq{interrupt: interrupt, noempty: noempty, timestamp: timestamp}:

mainLoop

The mainLoop is what listens to the “newWork channel”:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// mainLoop is a standalone goroutine to regenerate the sealing task based on the received event.
func (w *worker) mainLoop() {
	defer w.wg.Done()
	defer w.txsSub.Unsubscribe()
	defer w.chainHeadSub.Unsubscribe()
	defer w.chainSideSub.Unsubscribe()
	defer func() {
		if w.current != nil && w.current.state != nil {
			w.current.state.StopPrefetcher()
		}
	}()

	for {
		select {
		case req := <-w.newWorkCh:
			w.commitNewWork(req.interrupt, req.noempty, req.timestamp)

image

Header creation and transactions

The commitNewWork function is called with the request.

This is pretty much the most important function of the whole process.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// commitNewWork generates several new sealing tasks based on the parent block.
func (w *worker) commitNewWork(interrupt *int32, noempty bool, timestamp int64) {
	w.mu.RLock()
	defer w.mu.RUnlock()

	tstart := time.Now()
	parent := w.chain.CurrentBlock()

	if parent.Time() >= uint64(timestamp) {
		timestamp = int64(parent.Time() + 1)
	}
	num := parent.Number()
	header := &types.Header{
		ParentHash: parent.Hash(),
		Number:     num.Add(num, common.Big1),
		GasLimit:   core.CalcGasLimit(parent.GasLimit(), w.config.GasCeil),
		Extra:      w.extra,
		Time:       uint64(timestamp),
	}

Here we can see a header being created with the number being incremented from the parent! Idk how insane that sounds, but damn.

The Prepare method of the engine being called here, which sets the difficulty of the block for ethash:

1
	if err := w.engine.Prepare(w.chain, header); err != nil {

Then we can see the pending transactions being taken in the function:

1
2
	// Fill the block with all available pending transactions.
	pending := w.eth.TxPool().Pending(true)
1
2
3
4
5
6
7
8
9
	// Split the pending transactions into locals and remotes
	localTxs, remoteTxs := make(map[common.Address]types.Transactions), pending
	for _, account := range w.eth.TxPool().Locals() {
		if txs := remoteTxs[account]; len(txs) > 0 {
			delete(remoteTxs, account)
			localTxs[account] = txs
		}
	}

Local transactions being the txs sent from the running node.

1
2
3
4
5
6
	if len(remoteTxs) > 0 {
		txs := types.NewTransactionsByPriceAndNonce(w.current.signer, remoteTxs, header.BaseFee)
		if w.commitTransactions(txs, w.coinbase, interrupt) {
			return
		}
	}

image

commitTransactions function which iteratively calls commitTransaction for each function wherein the actual execution takes place.

Transaction execution

1
2
3
4
5
6
7
func (w *worker) commitTransactions(txs *types.TransactionsByPriceAndNonce, coinbase common.Address, interrupt *int32) bool {
	// Short circuit if current is nil
	if w.current == nil {
		return true
	}

1
2
3
4
5
6
7
8
func (w *worker) commitTransactions(txs *types.TransactionsByPriceAndNonce, coinbase common.Address, interrupt *int32) bool {
	// Short circuit if current is nil
	if w.current == nil {
		return true
		
		...
		
	logs, err := w.commitTransaction(tx, coinbase)

image

1
2
3
4
5
6
7
8
9
10
11
12
13
14
func (w *worker) commitTransaction(tx *types.Transaction, coinbase common.Address) ([]*types.Log, error) {
	snap := w.current.state.Snapshot()

	receipt, err := core.ApplyTransaction(w.chainConfig, w.chain, &coinbase, w.current.gasPool, w.current.state, w.current.header, tx, &w.current.header.GasUsed, *w.chain.GetVMConfig())
	if err != nil {
		w.current.state.RevertToSnapshot(snap)
		return nil, err
	}
	w.current.txs = append(w.current.txs, tx)
	w.current.receipts = append(w.current.receipts, receipt)

	return receipt.Logs, nil
}

The ApplyTransaction function is what creates the evm context and runs the transaction.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
// ApplyTransaction attempts to apply a transaction to the given state database
// and uses the input parameters for its environment. It returns the receipt
// for the transaction, gas used and an error if the transaction failed,
// indicating the block was invalid.
func ApplyTransaction(config *params.ChainConfig, bc ChainContext, author *common.Address, gp *GasPool, statedb *state.StateDB, header *types.Header, tx *types.Transaction, usedGas *uint64, cfg vm.Config) (*types.Receipt, error) {
	msg, err := tx.AsMessage(types.MakeSigner(config, header.Number), header.BaseFee)
	if err != nil {
		return nil, err
	}
	// Create a new context to be used in the EVM environment
	blockContext := NewEVMBlockContext(header, bc, author)
	vmenv := vm.NewEVM(blockContext, vm.TxContext{}, statedb, config, cfg)
	return applyTransaction(msg, config, bc, author, gp, statedb, header.Number, header.Hash(), tx, usedGas, vmenv)
}

Then in the commitNewWork() function we can see commit() being called:

Creation of a block

image

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// commit runs any post-transaction state modifications, assembles the final block
// and commits new work if consensus engine is running.
func (w *worker) commit(uncles []*types.Header, interval func(), update bool, start time.Time) error {
	// Deep copy receipts here to avoid interaction between different tasks.
	receipts := copyReceipts(w.current.receipts)
	s := w.current.state.Copy()
	block, err := w.engine.FinalizeAndAssemble(w.chain, w.current.header, s, w.current.txs, uncles, receipts)
	if err != nil {
		return err
	}
	
	...
	
	}
		select {
		case w.taskCh <- &task{receipts: receipts, state: s, block: block, createdAt: time.Now()}:
			w.unconfirmed.Shift(block.NumberU64() - 1)
			log.Info("Commit new mining work", "number", block.Number(), "sealhash", w.engine.SealHash(block.Header()),
				"uncles", len(uncles), "txs", w.current.tcount,
				"gas", block.GasUsed(), "fees", totalFees(block, receipts),
				"elapsed", common.PrettyDuration(time.Since(start)))

		case <-w.exitCh:
			log.Info("Worker has exited")
		}
	}


This is where the actual block creation takes place, for ethhash:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// Finalize implements consensus.Engine, accumulating the block and uncle rewards,
// setting the final state on the header
func (ethash *Ethash) Finalize(chain consensus.ChainHeaderReader, header *types.Header, state *state.StateDB, txs []*types.Transaction, uncles []*types.Header) {
	// Accumulate any block and uncle rewards and commit the final state root
	accumulateRewards(chain.Config(), state, header, uncles)
	header.Root = state.IntermediateRoot(chain.Config().IsEIP158(header.Number))
}

// FinalizeAndAssemble implements consensus.Engine, accumulating the block and
// uncle rewards, setting the final state and assembling the block.
func (ethash *Ethash) FinalizeAndAssemble(chain consensus.ChainHeaderReader, header *types.Header, state *state.StateDB, txs []*types.Transaction, uncles []*types.Header, receipts []*types.Receipt) (*types.Block, error) {
	// Finalize block
	ethash.Finalize(chain, header, state, txs, uncles)

	// Header seems complete, assemble into a block and return
	return types.NewBlock(header, txs, uncles, receipts, trie.NewStackTrie(nil)), nil
}

The state root is being set and a new block is created with the header and the transactions!

So basically we’ve created a block now!

Now in commit() we can see this line:

1
		case w.taskCh <- &task{receipts: receipts, state: s, block: block, createdAt: time.Now()}:

Where a new task is created with the block and sent to the task channel.

taskLoop

Up next we have the task loop which listens to the task channel, where the actual mining process happens!

image

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
// taskLoop is a standalone goroutine to fetch sealing task from the generator and
// push them to consensus engine.
func (w *worker) taskLoop() {
	defer w.wg.Done()
	var (
		stopCh chan struct{}
		prev   common.Hash
	)
...

	for {
		select {
		case task := <-w.taskCh:
			if w.newTaskHook != nil {
				w.newTaskHook(task)
			}
			// Reject duplicate sealing work due to resubmitting.
			sealHash := w.engine.SealHash(task.block.Header())
			if sealHash == prev {
				continue
			}
			// Interrupt previous sealing operation
			interrupt()
			stopCh, prev = make(chan struct{}), sealHash

			if w.skipSealHook != nil && w.skipSealHook(task) {
				continue
			}
			w.pendingMu.Lock()
			w.pendingTasks[sealHash] = task
			w.pendingMu.Unlock()

			if err := w.engine.Seal(w.chain, task.block, w.resultCh, stopCh); err != nil {
				log.Warn("Block sealing failed", "err", err)
				w.pendingMu.Lock()
				delete(w.pendingTasks, sealHash)
				w.pendingMu.Unlock()
			}
		case <-w.exitCh:
			interrupt()
			return
		}
1
			if err := w.engine.Seal(w.chain, task.block, w.resultCh, stopCh); err != nil {

The sealing method is where the whole block mining process happens, which I’ve explained in my previous post.

https://pvnotpv.github.io/posts/ethhash/

Mining in action

Now let’s see the actual mining process in action! Here I’ve created a private network:

1
pv@arch ~/t/go-ethereum ((v1.10.14))> build/bin/geth --datadir .private-net/pow-node-clean --networkid 1337 --nodiscover --http --http.addr 127.0.0.1 --http.port 8545 --http.api eth,net,web3,miner --mine --miner.threads 1 --miner.etherbase 0x3D27412aC6D1bB84E9621bcE0D95207f3dC2214C

First we can see the generation of the DAG:

image

Then the whole mining process beings:

image

1
2
3
4
5
6
pv@arch ~> curl -s -X POST http://127.0.0.1:8545 \
                 -H 'Content-Type: application/json' \
                 --data '{"jsonrpc":"2.0","method":"eth_blockNumber","params":[],"id":1}'
{"jsonrpc":"2.0","id":1,"result":"0x19"}
pv@arch ~>

This post is licensed under CC BY 4.0 by the author.