From aa40192bf639bb3dcbbf04f74c01be4b28e0c1f2 Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?Isabel=20Sch=C3=B6ps=20Thiel=20=40IsabelSchoepd?=
<155141998+IST-Github@users.noreply.github.com>
Date: Thu, 4 Jan 2024 04:32:11 +0100
Subject: [PATCH] Delete miner directory
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
Signed-off-by: Isabel Schöps Thiel @IsabelSchoepd <155141998+IST-Github@users.noreply.github.com>
---
miner/miner.go | 246 -------
miner/miner_test.go | 336 ---------
miner/ordering.go | 147 ----
miner/ordering_test.go | 195 -----
miner/payload_building.go | 243 -------
miner/payload_building_test.go | 160 -----
miner/stress/clique/main.go | 223 ------
miner/worker.go | 1214 --------------------------------
miner/worker_test.go | 509 -------------
9 files changed, 3273 deletions(-)
delete mode 100644 miner/miner.go
delete mode 100644 miner/miner_test.go
delete mode 100644 miner/ordering.go
delete mode 100644 miner/ordering_test.go
delete mode 100644 miner/payload_building.go
delete mode 100644 miner/payload_building_test.go
delete mode 100644 miner/stress/clique/main.go
delete mode 100644 miner/worker.go
delete mode 100644 miner/worker_test.go
diff --git a/miner/miner.go b/miner/miner.go
deleted file mode 100644
index b7273948f5..0000000000
--- a/miner/miner.go
+++ /dev/null
@@ -1,246 +0,0 @@
-// Copyright 2014 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see .
-
-// Package miner implements Ethereum block creation and mining.
-package miner
-
-import (
- "fmt"
- "math/big"
- "sync"
- "time"
-
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/common/hexutil"
- "github.com/ethereum/go-ethereum/consensus"
- "github.com/ethereum/go-ethereum/core"
- "github.com/ethereum/go-ethereum/core/state"
- "github.com/ethereum/go-ethereum/core/txpool"
- "github.com/ethereum/go-ethereum/core/types"
- "github.com/ethereum/go-ethereum/eth/downloader"
- "github.com/ethereum/go-ethereum/event"
- "github.com/ethereum/go-ethereum/log"
- "github.com/ethereum/go-ethereum/params"
-)
-
-// Backend wraps all methods required for mining. Only full node is capable
-// to offer all the functions here.
-type Backend interface {
- BlockChain() *core.BlockChain
- TxPool() *txpool.TxPool
-}
-
-// Config is the configuration parameters of mining.
-type Config struct {
- Etherbase common.Address `toml:",omitempty"` // Public address for block mining rewards
- ExtraData hexutil.Bytes `toml:",omitempty"` // Block extra data set by the miner
- GasFloor uint64 // Target gas floor for mined blocks.
- GasCeil uint64 // Target gas ceiling for mined blocks.
- GasPrice *big.Int // Minimum gas price for mining a transaction
- Recommit time.Duration // The time interval for miner to re-create mining work.
-
- NewPayloadTimeout time.Duration // The maximum time allowance for creating a new payload
-}
-
-// DefaultConfig contains default settings for miner.
-var DefaultConfig = Config{
- GasCeil: 30000000,
- GasPrice: big.NewInt(params.GWei),
-
- // The default recommit time is chosen as two seconds since
- // consensus-layer usually will wait a half slot of time(6s)
- // for payload generation. It should be enough for Geth to
- // run 3 rounds.
- Recommit: 2 * time.Second,
- NewPayloadTimeout: 2 * time.Second,
-}
-
-// Miner creates blocks and searches for proof-of-work values.
-type Miner struct {
- mux *event.TypeMux
- eth Backend
- engine consensus.Engine
- exitCh chan struct{}
- startCh chan struct{}
- stopCh chan struct{}
- worker *worker
-
- wg sync.WaitGroup
-}
-
-func New(eth Backend, config *Config, chainConfig *params.ChainConfig, mux *event.TypeMux, engine consensus.Engine, isLocalBlock func(header *types.Header) bool) *Miner {
- miner := &Miner{
- mux: mux,
- eth: eth,
- engine: engine,
- exitCh: make(chan struct{}),
- startCh: make(chan struct{}),
- stopCh: make(chan struct{}),
- worker: newWorker(config, chainConfig, engine, eth, mux, isLocalBlock, true),
- }
- miner.wg.Add(1)
- go miner.update()
- return miner
-}
-
-// update keeps track of the downloader events. Please be aware that this is a one shot type of update loop.
-// It's entered once and as soon as `Done` or `Failed` has been broadcasted the events are unregistered and
-// the loop is exited. This to prevent a major security vuln where external parties can DOS you with blocks
-// and halt your mining operation for as long as the DOS continues.
-func (miner *Miner) update() {
- defer miner.wg.Done()
-
- events := miner.mux.Subscribe(downloader.StartEvent{}, downloader.DoneEvent{}, downloader.FailedEvent{})
- defer func() {
- if !events.Closed() {
- events.Unsubscribe()
- }
- }()
-
- shouldStart := false
- canStart := true
- dlEventCh := events.Chan()
- for {
- select {
- case ev := <-dlEventCh:
- if ev == nil {
- // Unsubscription done, stop listening
- dlEventCh = nil
- continue
- }
- switch ev.Data.(type) {
- case downloader.StartEvent:
- wasMining := miner.Mining()
- miner.worker.stop()
- canStart = false
- if wasMining {
- // Resume mining after sync was finished
- shouldStart = true
- log.Info("Mining aborted due to sync")
- }
- miner.worker.syncing.Store(true)
-
- case downloader.FailedEvent:
- canStart = true
- if shouldStart {
- miner.worker.start()
- }
- miner.worker.syncing.Store(false)
-
- case downloader.DoneEvent:
- canStart = true
- if shouldStart {
- miner.worker.start()
- }
- miner.worker.syncing.Store(false)
-
- // Stop reacting to downloader events
- events.Unsubscribe()
- }
- case <-miner.startCh:
- if canStart {
- miner.worker.start()
- }
- shouldStart = true
- case <-miner.stopCh:
- shouldStart = false
- miner.worker.stop()
- case <-miner.exitCh:
- miner.worker.close()
- return
- }
- }
-}
-
-func (miner *Miner) Start() {
- miner.startCh <- struct{}{}
-}
-
-func (miner *Miner) Stop() {
- miner.stopCh <- struct{}{}
-}
-
-func (miner *Miner) Close() {
- close(miner.exitCh)
- miner.wg.Wait()
-}
-
-func (miner *Miner) Mining() bool {
- return miner.worker.isRunning()
-}
-
-func (miner *Miner) Hashrate() uint64 {
- if pow, ok := miner.engine.(consensus.PoW); ok {
- return uint64(pow.Hashrate())
- }
- return 0
-}
-
-func (miner *Miner) SetExtra(extra []byte) error {
- if uint64(len(extra)) > params.MaximumExtraDataSize {
- return fmt.Errorf("extra exceeds max length. %d > %v", len(extra), params.MaximumExtraDataSize)
- }
- miner.worker.setExtra(extra)
- return nil
-}
-
-// SetRecommitInterval sets the interval for sealing work resubmitting.
-func (miner *Miner) SetRecommitInterval(interval time.Duration) {
- miner.worker.setRecommitInterval(interval)
-}
-
-// Pending returns the currently pending block and associated state. The returned
-// values can be nil in case the pending block is not initialized
-func (miner *Miner) Pending() (*types.Block, *state.StateDB) {
- return miner.worker.pending()
-}
-
-// PendingBlock returns the currently pending block. The returned block can be
-// nil in case the pending block is not initialized.
-//
-// Note, to access both the pending block and the pending state
-// simultaneously, please use Pending(), as the pending state can
-// change between multiple method calls
-func (miner *Miner) PendingBlock() *types.Block {
- return miner.worker.pendingBlock()
-}
-
-// PendingBlockAndReceipts returns the currently pending block and corresponding receipts.
-// The returned values can be nil in case the pending block is not initialized.
-func (miner *Miner) PendingBlockAndReceipts() (*types.Block, types.Receipts) {
- return miner.worker.pendingBlockAndReceipts()
-}
-
-func (miner *Miner) SetEtherbase(addr common.Address) {
- miner.worker.setEtherbase(addr)
-}
-
-// SetGasCeil sets the gaslimit to strive for when mining blocks post 1559.
-// For pre-1559 blocks, it sets the ceiling.
-func (miner *Miner) SetGasCeil(ceil uint64) {
- miner.worker.setGasCeil(ceil)
-}
-
-// SubscribePendingLogs starts delivering logs from pending transactions
-// to the given channel.
-func (miner *Miner) SubscribePendingLogs(ch chan<- []*types.Log) event.Subscription {
- return miner.worker.pendingLogsFeed.Subscribe(ch)
-}
-
-// BuildPayload builds the payload according to the provided parameters.
-func (miner *Miner) BuildPayload(args *BuildPayloadArgs) (*Payload, error) {
- return miner.worker.buildPayload(args)
-}
diff --git a/miner/miner_test.go b/miner/miner_test.go
deleted file mode 100644
index 411d6026ce..0000000000
--- a/miner/miner_test.go
+++ /dev/null
@@ -1,336 +0,0 @@
-// Copyright 2020 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see .
-
-// Package miner implements Ethereum block creation and mining.
-package miner
-
-import (
- "errors"
- "math/big"
- "testing"
- "time"
-
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/consensus/clique"
- "github.com/ethereum/go-ethereum/core"
- "github.com/ethereum/go-ethereum/core/rawdb"
- "github.com/ethereum/go-ethereum/core/state"
- "github.com/ethereum/go-ethereum/core/txpool"
- "github.com/ethereum/go-ethereum/core/txpool/legacypool"
- "github.com/ethereum/go-ethereum/core/types"
- "github.com/ethereum/go-ethereum/core/vm"
- "github.com/ethereum/go-ethereum/crypto"
- "github.com/ethereum/go-ethereum/eth/downloader"
- "github.com/ethereum/go-ethereum/event"
- "github.com/ethereum/go-ethereum/params"
- "github.com/ethereum/go-ethereum/trie"
-)
-
-type mockBackend struct {
- bc *core.BlockChain
- txPool *txpool.TxPool
-}
-
-func NewMockBackend(bc *core.BlockChain, txPool *txpool.TxPool) *mockBackend {
- return &mockBackend{
- bc: bc,
- txPool: txPool,
- }
-}
-
-func (m *mockBackend) BlockChain() *core.BlockChain {
- return m.bc
-}
-
-func (m *mockBackend) TxPool() *txpool.TxPool {
- return m.txPool
-}
-
-func (m *mockBackend) StateAtBlock(block *types.Block, reexec uint64, base *state.StateDB, checkLive bool, preferDisk bool) (statedb *state.StateDB, err error) {
- return nil, errors.New("not supported")
-}
-
-type testBlockChain struct {
- root common.Hash
- config *params.ChainConfig
- statedb *state.StateDB
- gasLimit uint64
- chainHeadFeed *event.Feed
-}
-
-func (bc *testBlockChain) Config() *params.ChainConfig {
- return bc.config
-}
-
-func (bc *testBlockChain) CurrentBlock() *types.Header {
- return &types.Header{
- Number: new(big.Int),
- GasLimit: bc.gasLimit,
- }
-}
-
-func (bc *testBlockChain) GetBlock(hash common.Hash, number uint64) *types.Block {
- return types.NewBlock(bc.CurrentBlock(), nil, nil, nil, trie.NewStackTrie(nil))
-}
-
-func (bc *testBlockChain) StateAt(common.Hash) (*state.StateDB, error) {
- return bc.statedb, nil
-}
-
-func (bc *testBlockChain) HasState(root common.Hash) bool {
- return bc.root == root
-}
-
-func (bc *testBlockChain) SubscribeChainHeadEvent(ch chan<- core.ChainHeadEvent) event.Subscription {
- return bc.chainHeadFeed.Subscribe(ch)
-}
-
-func TestMiner(t *testing.T) {
- t.Parallel()
- miner, mux, cleanup := createMiner(t)
- defer cleanup(false)
-
- miner.Start()
- waitForMiningState(t, miner, true)
- // Start the downloader
- mux.Post(downloader.StartEvent{})
- waitForMiningState(t, miner, false)
- // Stop the downloader and wait for the update loop to run
- mux.Post(downloader.DoneEvent{})
- waitForMiningState(t, miner, true)
-
- // Subsequent downloader events after a successful DoneEvent should not cause the
- // miner to start or stop. This prevents a security vulnerability
- // that would allow entities to present fake high blocks that would
- // stop mining operations by causing a downloader sync
- // until it was discovered they were invalid, whereon mining would resume.
- mux.Post(downloader.StartEvent{})
- waitForMiningState(t, miner, true)
-
- mux.Post(downloader.FailedEvent{})
- waitForMiningState(t, miner, true)
-}
-
-// TestMinerDownloaderFirstFails tests that mining is only
-// permitted to run indefinitely once the downloader sees a DoneEvent (success).
-// An initial FailedEvent should allow mining to stop on a subsequent
-// downloader StartEvent.
-func TestMinerDownloaderFirstFails(t *testing.T) {
- t.Parallel()
- miner, mux, cleanup := createMiner(t)
- defer cleanup(false)
-
- miner.Start()
- waitForMiningState(t, miner, true)
- // Start the downloader
- mux.Post(downloader.StartEvent{})
- waitForMiningState(t, miner, false)
-
- // Stop the downloader and wait for the update loop to run
- mux.Post(downloader.FailedEvent{})
- waitForMiningState(t, miner, true)
-
- // Since the downloader hasn't yet emitted a successful DoneEvent,
- // we expect the miner to stop on next StartEvent.
- mux.Post(downloader.StartEvent{})
- waitForMiningState(t, miner, false)
-
- // Downloader finally succeeds.
- mux.Post(downloader.DoneEvent{})
- waitForMiningState(t, miner, true)
-
- // Downloader starts again.
- // Since it has achieved a DoneEvent once, we expect miner
- // state to be unchanged.
- mux.Post(downloader.StartEvent{})
- waitForMiningState(t, miner, true)
-
- mux.Post(downloader.FailedEvent{})
- waitForMiningState(t, miner, true)
-}
-
-func TestMinerStartStopAfterDownloaderEvents(t *testing.T) {
- t.Parallel()
- miner, mux, cleanup := createMiner(t)
- defer cleanup(false)
-
- miner.Start()
- waitForMiningState(t, miner, true)
- // Start the downloader
- mux.Post(downloader.StartEvent{})
- waitForMiningState(t, miner, false)
-
- // Downloader finally succeeds.
- mux.Post(downloader.DoneEvent{})
- waitForMiningState(t, miner, true)
-
- miner.Stop()
- waitForMiningState(t, miner, false)
-
- miner.Start()
- waitForMiningState(t, miner, true)
-
- miner.Stop()
- waitForMiningState(t, miner, false)
-}
-
-func TestStartWhileDownload(t *testing.T) {
- t.Parallel()
- miner, mux, cleanup := createMiner(t)
- defer cleanup(false)
- waitForMiningState(t, miner, false)
- miner.Start()
- waitForMiningState(t, miner, true)
- // Stop the downloader and wait for the update loop to run
- mux.Post(downloader.StartEvent{})
- waitForMiningState(t, miner, false)
- // Starting the miner after the downloader should not work
- miner.Start()
- waitForMiningState(t, miner, false)
-}
-
-func TestStartStopMiner(t *testing.T) {
- t.Parallel()
- miner, _, cleanup := createMiner(t)
- defer cleanup(false)
- waitForMiningState(t, miner, false)
- miner.Start()
- waitForMiningState(t, miner, true)
- miner.Stop()
- waitForMiningState(t, miner, false)
-}
-
-func TestCloseMiner(t *testing.T) {
- t.Parallel()
- miner, _, cleanup := createMiner(t)
- defer cleanup(true)
- waitForMiningState(t, miner, false)
- miner.Start()
- waitForMiningState(t, miner, true)
- // Terminate the miner and wait for the update loop to run
- miner.Close()
- waitForMiningState(t, miner, false)
-}
-
-// TestMinerSetEtherbase checks that etherbase becomes set even if mining isn't
-// possible at the moment
-func TestMinerSetEtherbase(t *testing.T) {
- t.Parallel()
- miner, mux, cleanup := createMiner(t)
- defer cleanup(false)
- miner.Start()
- waitForMiningState(t, miner, true)
- // Start the downloader
- mux.Post(downloader.StartEvent{})
- waitForMiningState(t, miner, false)
- // Now user tries to configure proper mining address
- miner.Start()
- // Stop the downloader and wait for the update loop to run
- mux.Post(downloader.DoneEvent{})
- waitForMiningState(t, miner, true)
-
- coinbase := common.HexToAddress("0xdeedbeef")
- miner.SetEtherbase(coinbase)
- if addr := miner.worker.etherbase(); addr != coinbase {
- t.Fatalf("Unexpected etherbase want %x got %x", coinbase, addr)
- }
-}
-
-// waitForMiningState waits until either
-// * the desired mining state was reached
-// * a timeout was reached which fails the test
-func waitForMiningState(t *testing.T, m *Miner, mining bool) {
- t.Helper()
-
- var state bool
- for i := 0; i < 100; i++ {
- time.Sleep(10 * time.Millisecond)
- if state = m.Mining(); state == mining {
- return
- }
- }
- t.Fatalf("Mining() == %t, want %t", state, mining)
-}
-
-func minerTestGenesisBlock(period uint64, gasLimit uint64, faucet common.Address) *core.Genesis {
- config := *params.AllCliqueProtocolChanges
- config.Clique = ¶ms.CliqueConfig{
- Period: period,
- Epoch: config.Clique.Epoch,
- }
-
- // Assemble and return the genesis with the precompiles and faucet pre-funded
- return &core.Genesis{
- Config: &config,
- ExtraData: append(append(make([]byte, 32), faucet[:]...), make([]byte, crypto.SignatureLength)...),
- GasLimit: gasLimit,
- BaseFee: big.NewInt(params.InitialBaseFee),
- Difficulty: big.NewInt(1),
- Alloc: map[common.Address]core.GenesisAccount{
- common.BytesToAddress([]byte{1}): {Balance: big.NewInt(1)}, // ECRecover
- common.BytesToAddress([]byte{2}): {Balance: big.NewInt(1)}, // SHA256
- common.BytesToAddress([]byte{3}): {Balance: big.NewInt(1)}, // RIPEMD
- common.BytesToAddress([]byte{4}): {Balance: big.NewInt(1)}, // Identity
- common.BytesToAddress([]byte{5}): {Balance: big.NewInt(1)}, // ModExp
- common.BytesToAddress([]byte{6}): {Balance: big.NewInt(1)}, // ECAdd
- common.BytesToAddress([]byte{7}): {Balance: big.NewInt(1)}, // ECScalarMul
- common.BytesToAddress([]byte{8}): {Balance: big.NewInt(1)}, // ECPairing
- common.BytesToAddress([]byte{9}): {Balance: big.NewInt(1)}, // BLAKE2b
- faucet: {Balance: new(big.Int).Sub(new(big.Int).Lsh(big.NewInt(1), 256), big.NewInt(9))},
- },
- }
-}
-func createMiner(t *testing.T) (*Miner, *event.TypeMux, func(skipMiner bool)) {
- // Create Ethash config
- config := Config{
- Etherbase: common.HexToAddress("123456789"),
- }
- // Create chainConfig
- chainDB := rawdb.NewMemoryDatabase()
- triedb := trie.NewDatabase(chainDB, nil)
- genesis := minerTestGenesisBlock(15, 11_500_000, common.HexToAddress("12345"))
- chainConfig, _, err := core.SetupGenesisBlock(chainDB, triedb, genesis)
- if err != nil {
- t.Fatalf("can't create new chain config: %v", err)
- }
- // Create consensus engine
- engine := clique.New(chainConfig.Clique, chainDB)
- // Create Ethereum backend
- bc, err := core.NewBlockChain(chainDB, nil, genesis, nil, engine, vm.Config{}, nil, nil)
- if err != nil {
- t.Fatalf("can't create new chain %v", err)
- }
- statedb, _ := state.New(bc.Genesis().Root(), bc.StateCache(), nil)
- blockchain := &testBlockChain{bc.Genesis().Root(), chainConfig, statedb, 10000000, new(event.Feed)}
-
- pool := legacypool.New(testTxPoolConfig, blockchain)
- txpool, _ := txpool.New(new(big.Int).SetUint64(testTxPoolConfig.PriceLimit), blockchain, []txpool.SubPool{pool})
-
- backend := NewMockBackend(bc, txpool)
- // Create event Mux
- mux := new(event.TypeMux)
- // Create Miner
- miner := New(backend, &config, chainConfig, mux, engine, nil)
- cleanup := func(skipMiner bool) {
- bc.Stop()
- engine.Close()
- txpool.Close()
- if !skipMiner {
- miner.Close()
- }
- }
- return miner, mux, cleanup
-}
diff --git a/miner/ordering.go b/miner/ordering.go
deleted file mode 100644
index 4c3055f0d3..0000000000
--- a/miner/ordering.go
+++ /dev/null
@@ -1,147 +0,0 @@
-// Copyright 2014 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see .
-
-package miner
-
-import (
- "container/heap"
- "math/big"
-
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/common/math"
- "github.com/ethereum/go-ethereum/core/txpool"
- "github.com/ethereum/go-ethereum/core/types"
-)
-
-// txWithMinerFee wraps a transaction with its gas price or effective miner gasTipCap
-type txWithMinerFee struct {
- tx *txpool.LazyTransaction
- from common.Address
- fees *big.Int
-}
-
-// newTxWithMinerFee creates a wrapped transaction, calculating the effective
-// miner gasTipCap if a base fee is provided.
-// Returns error in case of a negative effective miner gasTipCap.
-func newTxWithMinerFee(tx *txpool.LazyTransaction, from common.Address, baseFee *big.Int) (*txWithMinerFee, error) {
- tip := new(big.Int).Set(tx.GasTipCap)
- if baseFee != nil {
- if tx.GasFeeCap.Cmp(baseFee) < 0 {
- return nil, types.ErrGasFeeCapTooLow
- }
- tip = math.BigMin(tx.GasTipCap, new(big.Int).Sub(tx.GasFeeCap, baseFee))
- }
- return &txWithMinerFee{
- tx: tx,
- from: from,
- fees: tip,
- }, nil
-}
-
-// txByPriceAndTime implements both the sort and the heap interface, making it useful
-// for all at once sorting as well as individually adding and removing elements.
-type txByPriceAndTime []*txWithMinerFee
-
-func (s txByPriceAndTime) Len() int { return len(s) }
-func (s txByPriceAndTime) Less(i, j int) bool {
- // If the prices are equal, use the time the transaction was first seen for
- // deterministic sorting
- cmp := s[i].fees.Cmp(s[j].fees)
- if cmp == 0 {
- return s[i].tx.Time.Before(s[j].tx.Time)
- }
- return cmp > 0
-}
-func (s txByPriceAndTime) Swap(i, j int) { s[i], s[j] = s[j], s[i] }
-
-func (s *txByPriceAndTime) Push(x interface{}) {
- *s = append(*s, x.(*txWithMinerFee))
-}
-
-func (s *txByPriceAndTime) Pop() interface{} {
- old := *s
- n := len(old)
- x := old[n-1]
- old[n-1] = nil
- *s = old[0 : n-1]
- return x
-}
-
-// transactionsByPriceAndNonce represents a set of transactions that can return
-// transactions in a profit-maximizing sorted order, while supporting removing
-// entire batches of transactions for non-executable accounts.
-type transactionsByPriceAndNonce struct {
- txs map[common.Address][]*txpool.LazyTransaction // Per account nonce-sorted list of transactions
- heads txByPriceAndTime // Next transaction for each unique account (price heap)
- signer types.Signer // Signer for the set of transactions
- baseFee *big.Int // Current base fee
-}
-
-// newTransactionsByPriceAndNonce creates a transaction set that can retrieve
-// price sorted transactions in a nonce-honouring way.
-//
-// Note, the input map is reowned so the caller should not interact any more with
-// if after providing it to the constructor.
-func newTransactionsByPriceAndNonce(signer types.Signer, txs map[common.Address][]*txpool.LazyTransaction, baseFee *big.Int) *transactionsByPriceAndNonce {
- // Initialize a price and received time based heap with the head transactions
- heads := make(txByPriceAndTime, 0, len(txs))
- for from, accTxs := range txs {
- wrapped, err := newTxWithMinerFee(accTxs[0], from, baseFee)
- if err != nil {
- delete(txs, from)
- continue
- }
- heads = append(heads, wrapped)
- txs[from] = accTxs[1:]
- }
- heap.Init(&heads)
-
- // Assemble and return the transaction set
- return &transactionsByPriceAndNonce{
- txs: txs,
- heads: heads,
- signer: signer,
- baseFee: baseFee,
- }
-}
-
-// Peek returns the next transaction by price.
-func (t *transactionsByPriceAndNonce) Peek() *txpool.LazyTransaction {
- if len(t.heads) == 0 {
- return nil
- }
- return t.heads[0].tx
-}
-
-// Shift replaces the current best head with the next one from the same account.
-func (t *transactionsByPriceAndNonce) Shift() {
- acc := t.heads[0].from
- if txs, ok := t.txs[acc]; ok && len(txs) > 0 {
- if wrapped, err := newTxWithMinerFee(txs[0], acc, t.baseFee); err == nil {
- t.heads[0], t.txs[acc] = wrapped, txs[1:]
- heap.Fix(&t.heads, 0)
- return
- }
- }
- heap.Pop(&t.heads)
-}
-
-// Pop removes the best transaction, *not* replacing it with the next one from
-// the same account. This should be used when a transaction cannot be executed
-// and hence all subsequent ones should be discarded from the same account.
-func (t *transactionsByPriceAndNonce) Pop() {
- heap.Pop(&t.heads)
-}
diff --git a/miner/ordering_test.go b/miner/ordering_test.go
deleted file mode 100644
index e5868d7a06..0000000000
--- a/miner/ordering_test.go
+++ /dev/null
@@ -1,195 +0,0 @@
-// Copyright 2014 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see .
-
-package miner
-
-import (
- "crypto/ecdsa"
- "math/big"
- "math/rand"
- "testing"
- "time"
-
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/core/txpool"
- "github.com/ethereum/go-ethereum/core/types"
- "github.com/ethereum/go-ethereum/crypto"
-)
-
-func TestTransactionPriceNonceSortLegacy(t *testing.T) {
- t.Parallel()
- testTransactionPriceNonceSort(t, nil)
-}
-
-func TestTransactionPriceNonceSort1559(t *testing.T) {
- t.Parallel()
- testTransactionPriceNonceSort(t, big.NewInt(0))
- testTransactionPriceNonceSort(t, big.NewInt(5))
- testTransactionPriceNonceSort(t, big.NewInt(50))
-}
-
-// Tests that transactions can be correctly sorted according to their price in
-// decreasing order, but at the same time with increasing nonces when issued by
-// the same account.
-func testTransactionPriceNonceSort(t *testing.T, baseFee *big.Int) {
- // Generate a batch of accounts to start with
- keys := make([]*ecdsa.PrivateKey, 25)
- for i := 0; i < len(keys); i++ {
- keys[i], _ = crypto.GenerateKey()
- }
- signer := types.LatestSignerForChainID(common.Big1)
-
- // Generate a batch of transactions with overlapping values, but shifted nonces
- groups := map[common.Address][]*txpool.LazyTransaction{}
- expectedCount := 0
- for start, key := range keys {
- addr := crypto.PubkeyToAddress(key.PublicKey)
- count := 25
- for i := 0; i < 25; i++ {
- var tx *types.Transaction
- gasFeeCap := rand.Intn(50)
- if baseFee == nil {
- tx = types.NewTx(&types.LegacyTx{
- Nonce: uint64(start + i),
- To: &common.Address{},
- Value: big.NewInt(100),
- Gas: 100,
- GasPrice: big.NewInt(int64(gasFeeCap)),
- Data: nil,
- })
- } else {
- tx = types.NewTx(&types.DynamicFeeTx{
- Nonce: uint64(start + i),
- To: &common.Address{},
- Value: big.NewInt(100),
- Gas: 100,
- GasFeeCap: big.NewInt(int64(gasFeeCap)),
- GasTipCap: big.NewInt(int64(rand.Intn(gasFeeCap + 1))),
- Data: nil,
- })
- if count == 25 && int64(gasFeeCap) < baseFee.Int64() {
- count = i
- }
- }
- tx, err := types.SignTx(tx, signer, key)
- if err != nil {
- t.Fatalf("failed to sign tx: %s", err)
- }
- groups[addr] = append(groups[addr], &txpool.LazyTransaction{
- Hash: tx.Hash(),
- Tx: tx,
- Time: tx.Time(),
- GasFeeCap: tx.GasFeeCap(),
- GasTipCap: tx.GasTipCap(),
- Gas: tx.Gas(),
- BlobGas: tx.BlobGas(),
- })
- }
- expectedCount += count
- }
- // Sort the transactions and cross check the nonce ordering
- txset := newTransactionsByPriceAndNonce(signer, groups, baseFee)
-
- txs := types.Transactions{}
- for tx := txset.Peek(); tx != nil; tx = txset.Peek() {
- txs = append(txs, tx.Tx)
- txset.Shift()
- }
- if len(txs) != expectedCount {
- t.Errorf("expected %d transactions, found %d", expectedCount, len(txs))
- }
- for i, txi := range txs {
- fromi, _ := types.Sender(signer, txi)
-
- // Make sure the nonce order is valid
- for j, txj := range txs[i+1:] {
- fromj, _ := types.Sender(signer, txj)
- if fromi == fromj && txi.Nonce() > txj.Nonce() {
- t.Errorf("invalid nonce ordering: tx #%d (A=%x N=%v) < tx #%d (A=%x N=%v)", i, fromi[:4], txi.Nonce(), i+j, fromj[:4], txj.Nonce())
- }
- }
- // If the next tx has different from account, the price must be lower than the current one
- if i+1 < len(txs) {
- next := txs[i+1]
- fromNext, _ := types.Sender(signer, next)
- tip, err := txi.EffectiveGasTip(baseFee)
- nextTip, nextErr := next.EffectiveGasTip(baseFee)
- if err != nil || nextErr != nil {
- t.Errorf("error calculating effective tip: %v, %v", err, nextErr)
- }
- if fromi != fromNext && tip.Cmp(nextTip) < 0 {
- t.Errorf("invalid gasprice ordering: tx #%d (A=%x P=%v) < tx #%d (A=%x P=%v)", i, fromi[:4], txi.GasPrice(), i+1, fromNext[:4], next.GasPrice())
- }
- }
- }
-}
-
-// Tests that if multiple transactions have the same price, the ones seen earlier
-// are prioritized to avoid network spam attacks aiming for a specific ordering.
-func TestTransactionTimeSort(t *testing.T) {
- t.Parallel()
- // Generate a batch of accounts to start with
- keys := make([]*ecdsa.PrivateKey, 5)
- for i := 0; i < len(keys); i++ {
- keys[i], _ = crypto.GenerateKey()
- }
- signer := types.HomesteadSigner{}
-
- // Generate a batch of transactions with overlapping prices, but different creation times
- groups := map[common.Address][]*txpool.LazyTransaction{}
- for start, key := range keys {
- addr := crypto.PubkeyToAddress(key.PublicKey)
-
- tx, _ := types.SignTx(types.NewTransaction(0, common.Address{}, big.NewInt(100), 100, big.NewInt(1), nil), signer, key)
- tx.SetTime(time.Unix(0, int64(len(keys)-start)))
-
- groups[addr] = append(groups[addr], &txpool.LazyTransaction{
- Hash: tx.Hash(),
- Tx: tx,
- Time: tx.Time(),
- GasFeeCap: tx.GasFeeCap(),
- GasTipCap: tx.GasTipCap(),
- Gas: tx.Gas(),
- BlobGas: tx.BlobGas(),
- })
- }
- // Sort the transactions and cross check the nonce ordering
- txset := newTransactionsByPriceAndNonce(signer, groups, nil)
-
- txs := types.Transactions{}
- for tx := txset.Peek(); tx != nil; tx = txset.Peek() {
- txs = append(txs, tx.Tx)
- txset.Shift()
- }
- if len(txs) != len(keys) {
- t.Errorf("expected %d transactions, found %d", len(keys), len(txs))
- }
- for i, txi := range txs {
- fromi, _ := types.Sender(signer, txi)
- if i+1 < len(txs) {
- next := txs[i+1]
- fromNext, _ := types.Sender(signer, next)
-
- if txi.GasPrice().Cmp(next.GasPrice()) < 0 {
- t.Errorf("invalid gasprice ordering: tx #%d (A=%x P=%v) < tx #%d (A=%x P=%v)", i, fromi[:4], txi.GasPrice(), i+1, fromNext[:4], next.GasPrice())
- }
- // Make sure time order is ascending if the txs have the same gas price
- if txi.GasPrice().Cmp(next.GasPrice()) == 0 && txi.Time().After(next.Time()) {
- t.Errorf("invalid received time ordering: tx #%d (A=%x T=%v) > tx #%d (A=%x T=%v)", i, fromi[:4], txi.Time(), i+1, fromNext[:4], next.Time())
- }
- }
- }
-}
diff --git a/miner/payload_building.go b/miner/payload_building.go
deleted file mode 100644
index 69ffab75b5..0000000000
--- a/miner/payload_building.go
+++ /dev/null
@@ -1,243 +0,0 @@
-// Copyright 2022 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see
-
-package miner
-
-import (
- "crypto/sha256"
- "encoding/binary"
- "math/big"
- "sync"
- "time"
-
- "github.com/ethereum/go-ethereum/beacon/engine"
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/core/types"
- "github.com/ethereum/go-ethereum/log"
- "github.com/ethereum/go-ethereum/params"
- "github.com/ethereum/go-ethereum/rlp"
-)
-
-// BuildPayloadArgs contains the provided parameters for building payload.
-// Check engine-api specification for more details.
-// https://github.com/ethereum/execution-apis/blob/main/src/engine/cancun.md#payloadattributesv3
-type BuildPayloadArgs struct {
- Parent common.Hash // The parent block to build payload on top
- Timestamp uint64 // The provided timestamp of generated payload
- FeeRecipient common.Address // The provided recipient address for collecting transaction fee
- Random common.Hash // The provided randomness value
- Withdrawals types.Withdrawals // The provided withdrawals
- BeaconRoot *common.Hash // The provided beaconRoot (Cancun)
-}
-
-// Id computes an 8-byte identifier by hashing the components of the payload arguments.
-func (args *BuildPayloadArgs) Id() engine.PayloadID {
- // Hash
- hasher := sha256.New()
- hasher.Write(args.Parent[:])
- binary.Write(hasher, binary.BigEndian, args.Timestamp)
- hasher.Write(args.Random[:])
- hasher.Write(args.FeeRecipient[:])
- rlp.Encode(hasher, args.Withdrawals)
- if args.BeaconRoot != nil {
- hasher.Write(args.BeaconRoot[:])
- }
- var out engine.PayloadID
- copy(out[:], hasher.Sum(nil)[:8])
- return out
-}
-
-// Payload wraps the built payload(block waiting for sealing). According to the
-// engine-api specification, EL should build the initial version of the payload
-// which has an empty transaction set and then keep update it in order to maximize
-// the revenue. Therefore, the empty-block here is always available and full-block
-// will be set/updated afterwards.
-type Payload struct {
- id engine.PayloadID
- empty *types.Block
- full *types.Block
- sidecars []*types.BlobTxSidecar
- fullFees *big.Int
- stop chan struct{}
- lock sync.Mutex
- cond *sync.Cond
-}
-
-// newPayload initializes the payload object.
-func newPayload(empty *types.Block, id engine.PayloadID) *Payload {
- payload := &Payload{
- id: id,
- empty: empty,
- stop: make(chan struct{}),
- }
- log.Info("Starting work on payload", "id", payload.id)
- payload.cond = sync.NewCond(&payload.lock)
- return payload
-}
-
-// update updates the full-block with latest built version.
-func (payload *Payload) update(r *newPayloadResult, elapsed time.Duration) {
- payload.lock.Lock()
- defer payload.lock.Unlock()
-
- select {
- case <-payload.stop:
- return // reject stale update
- default:
- }
- // Ensure the newly provided full block has a higher transaction fee.
- // In post-merge stage, there is no uncle reward anymore and transaction
- // fee(apart from the mev revenue) is the only indicator for comparison.
- if payload.full == nil || r.fees.Cmp(payload.fullFees) > 0 {
- payload.full = r.block
- payload.fullFees = r.fees
- payload.sidecars = r.sidecars
-
- feesInEther := new(big.Float).Quo(new(big.Float).SetInt(r.fees), big.NewFloat(params.Ether))
- log.Info("Updated payload",
- "id", payload.id,
- "number", r.block.NumberU64(),
- "hash", r.block.Hash(),
- "txs", len(r.block.Transactions()),
- "withdrawals", len(r.block.Withdrawals()),
- "gas", r.block.GasUsed(),
- "fees", feesInEther,
- "root", r.block.Root(),
- "elapsed", common.PrettyDuration(elapsed),
- )
- }
- payload.cond.Broadcast() // fire signal for notifying full block
-}
-
-// Resolve returns the latest built payload and also terminates the background
-// thread for updating payload. It's safe to be called multiple times.
-func (payload *Payload) Resolve() *engine.ExecutionPayloadEnvelope {
- payload.lock.Lock()
- defer payload.lock.Unlock()
-
- select {
- case <-payload.stop:
- default:
- close(payload.stop)
- }
- if payload.full != nil {
- return engine.BlockToExecutableData(payload.full, payload.fullFees, payload.sidecars)
- }
- return engine.BlockToExecutableData(payload.empty, big.NewInt(0), nil)
-}
-
-// ResolveEmpty is basically identical to Resolve, but it expects empty block only.
-// It's only used in tests.
-func (payload *Payload) ResolveEmpty() *engine.ExecutionPayloadEnvelope {
- payload.lock.Lock()
- defer payload.lock.Unlock()
-
- return engine.BlockToExecutableData(payload.empty, big.NewInt(0), nil)
-}
-
-// ResolveFull is basically identical to Resolve, but it expects full block only.
-// Don't call Resolve until ResolveFull returns, otherwise it might block forever.
-func (payload *Payload) ResolveFull() *engine.ExecutionPayloadEnvelope {
- payload.lock.Lock()
- defer payload.lock.Unlock()
-
- if payload.full == nil {
- select {
- case <-payload.stop:
- return nil
- default:
- }
- // Wait the full payload construction. Note it might block
- // forever if Resolve is called in the meantime which
- // terminates the background construction process.
- payload.cond.Wait()
- }
- // Terminate the background payload construction
- select {
- case <-payload.stop:
- default:
- close(payload.stop)
- }
- return engine.BlockToExecutableData(payload.full, payload.fullFees, payload.sidecars)
-}
-
-// buildPayload builds the payload according to the provided parameters.
-func (w *worker) buildPayload(args *BuildPayloadArgs) (*Payload, error) {
- // Build the initial version with no transaction included. It should be fast
- // enough to run. The empty payload can at least make sure there is something
- // to deliver for not missing slot.
- emptyParams := &generateParams{
- timestamp: args.Timestamp,
- forceTime: true,
- parentHash: args.Parent,
- coinbase: args.FeeRecipient,
- random: args.Random,
- withdrawals: args.Withdrawals,
- beaconRoot: args.BeaconRoot,
- noTxs: true,
- }
- empty := w.getSealingBlock(emptyParams)
- if empty.err != nil {
- return nil, empty.err
- }
-
- // Construct a payload object for return.
- payload := newPayload(empty.block, args.Id())
-
- // Spin up a routine for updating the payload in background. This strategy
- // can maximum the revenue for including transactions with highest fee.
- go func() {
- // Setup the timer for re-building the payload. The initial clock is kept
- // for triggering process immediately.
- timer := time.NewTimer(0)
- defer timer.Stop()
-
- // Setup the timer for terminating the process if SECONDS_PER_SLOT (12s in
- // the Mainnet configuration) have passed since the point in time identified
- // by the timestamp parameter.
- endTimer := time.NewTimer(time.Second * 12)
-
- fullParams := &generateParams{
- timestamp: args.Timestamp,
- forceTime: true,
- parentHash: args.Parent,
- coinbase: args.FeeRecipient,
- random: args.Random,
- withdrawals: args.Withdrawals,
- beaconRoot: args.BeaconRoot,
- noTxs: false,
- }
-
- for {
- select {
- case <-timer.C:
- start := time.Now()
- r := w.getSealingBlock(fullParams)
- if r.err == nil {
- payload.update(r, time.Since(start))
- }
- timer.Reset(w.recommit)
- case <-payload.stop:
- log.Info("Stopping work on payload", "id", payload.id, "reason", "delivery")
- return
- case <-endTimer.C:
- log.Info("Stopping work on payload", "id", payload.id, "reason", "timeout")
- return
- }
- }
- }()
- return payload, nil
-}
diff --git a/miner/payload_building_test.go b/miner/payload_building_test.go
deleted file mode 100644
index 9283635224..0000000000
--- a/miner/payload_building_test.go
+++ /dev/null
@@ -1,160 +0,0 @@
-// Copyright 2022 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see
-
-package miner
-
-import (
- "reflect"
- "testing"
- "time"
-
- "github.com/ethereum/go-ethereum/beacon/engine"
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/consensus/ethash"
- "github.com/ethereum/go-ethereum/core/rawdb"
- "github.com/ethereum/go-ethereum/core/types"
- "github.com/ethereum/go-ethereum/params"
-)
-
-func TestBuildPayload(t *testing.T) {
- t.Parallel()
- var (
- db = rawdb.NewMemoryDatabase()
- recipient = common.HexToAddress("0xdeadbeef")
- )
- w, b := newTestWorker(t, params.TestChainConfig, ethash.NewFaker(), db, 0)
- defer w.close()
-
- timestamp := uint64(time.Now().Unix())
- args := &BuildPayloadArgs{
- Parent: b.chain.CurrentBlock().Hash(),
- Timestamp: timestamp,
- Random: common.Hash{},
- FeeRecipient: recipient,
- }
- payload, err := w.buildPayload(args)
- if err != nil {
- t.Fatalf("Failed to build payload %v", err)
- }
- verify := func(outer *engine.ExecutionPayloadEnvelope, txs int) {
- payload := outer.ExecutionPayload
- if payload.ParentHash != b.chain.CurrentBlock().Hash() {
- t.Fatal("Unexpect parent hash")
- }
- if payload.Random != (common.Hash{}) {
- t.Fatal("Unexpect random value")
- }
- if payload.Timestamp != timestamp {
- t.Fatal("Unexpect timestamp")
- }
- if payload.FeeRecipient != recipient {
- t.Fatal("Unexpect fee recipient")
- }
- if len(payload.Transactions) != txs {
- t.Fatal("Unexpect transaction set")
- }
- }
- empty := payload.ResolveEmpty()
- verify(empty, 0)
-
- full := payload.ResolveFull()
- verify(full, len(pendingTxs))
-
- // Ensure resolve can be called multiple times and the
- // result should be unchanged
- dataOne := payload.Resolve()
- dataTwo := payload.Resolve()
- if !reflect.DeepEqual(dataOne, dataTwo) {
- t.Fatal("Unexpected payload data")
- }
-}
-
-func TestPayloadId(t *testing.T) {
- t.Parallel()
- ids := make(map[string]int)
- for i, tt := range []*BuildPayloadArgs{
- {
- Parent: common.Hash{1},
- Timestamp: 1,
- Random: common.Hash{0x1},
- FeeRecipient: common.Address{0x1},
- },
- // Different parent
- {
- Parent: common.Hash{2},
- Timestamp: 1,
- Random: common.Hash{0x1},
- FeeRecipient: common.Address{0x1},
- },
- // Different timestamp
- {
- Parent: common.Hash{2},
- Timestamp: 2,
- Random: common.Hash{0x1},
- FeeRecipient: common.Address{0x1},
- },
- // Different Random
- {
- Parent: common.Hash{2},
- Timestamp: 2,
- Random: common.Hash{0x2},
- FeeRecipient: common.Address{0x1},
- },
- // Different fee-recipient
- {
- Parent: common.Hash{2},
- Timestamp: 2,
- Random: common.Hash{0x2},
- FeeRecipient: common.Address{0x2},
- },
- // Different withdrawals (non-empty)
- {
- Parent: common.Hash{2},
- Timestamp: 2,
- Random: common.Hash{0x2},
- FeeRecipient: common.Address{0x2},
- Withdrawals: []*types.Withdrawal{
- {
- Index: 0,
- Validator: 0,
- Address: common.Address{},
- Amount: 0,
- },
- },
- },
- // Different withdrawals (non-empty)
- {
- Parent: common.Hash{2},
- Timestamp: 2,
- Random: common.Hash{0x2},
- FeeRecipient: common.Address{0x2},
- Withdrawals: []*types.Withdrawal{
- {
- Index: 2,
- Validator: 0,
- Address: common.Address{},
- Amount: 0,
- },
- },
- },
- } {
- id := tt.Id().String()
- if prev, exists := ids[id]; exists {
- t.Errorf("ID collision, case %d and case %d: id %v", prev, i, id)
- }
- ids[id] = i
- }
-}
diff --git a/miner/stress/clique/main.go b/miner/stress/clique/main.go
deleted file mode 100644
index 13336cd83c..0000000000
--- a/miner/stress/clique/main.go
+++ /dev/null
@@ -1,223 +0,0 @@
-// Copyright 2018 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see .
-
-// This file contains a miner stress test based on the Clique consensus engine.
-package main
-
-import (
- "bytes"
- "crypto/ecdsa"
- "math/big"
- "math/rand"
- "os"
- "os/signal"
- "time"
-
- "github.com/ethereum/go-ethereum/accounts/keystore"
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/common/fdlimit"
- "github.com/ethereum/go-ethereum/core"
- "github.com/ethereum/go-ethereum/core/txpool/legacypool"
- "github.com/ethereum/go-ethereum/core/types"
- "github.com/ethereum/go-ethereum/crypto"
- "github.com/ethereum/go-ethereum/eth"
- "github.com/ethereum/go-ethereum/eth/downloader"
- "github.com/ethereum/go-ethereum/eth/ethconfig"
- "github.com/ethereum/go-ethereum/log"
- "github.com/ethereum/go-ethereum/miner"
- "github.com/ethereum/go-ethereum/node"
- "github.com/ethereum/go-ethereum/p2p"
- "github.com/ethereum/go-ethereum/p2p/enode"
- "github.com/ethereum/go-ethereum/params"
-)
-
-func main() {
- log.SetDefault(log.NewLogger(log.NewTerminalHandlerWithLevel(os.Stderr, log.LevelInfo, true)))
- fdlimit.Raise(2048)
-
- // Generate a batch of accounts to seal and fund with
- faucets := make([]*ecdsa.PrivateKey, 128)
- for i := 0; i < len(faucets); i++ {
- faucets[i], _ = crypto.GenerateKey()
- }
- sealers := make([]*ecdsa.PrivateKey, 4)
- for i := 0; i < len(sealers); i++ {
- sealers[i], _ = crypto.GenerateKey()
- }
- // Create a Clique network based off of the Sepolia config
- genesis := makeGenesis(faucets, sealers)
-
- // Handle interrupts.
- interruptCh := make(chan os.Signal, 5)
- signal.Notify(interruptCh, os.Interrupt)
-
- var (
- stacks []*node.Node
- nodes []*eth.Ethereum
- enodes []*enode.Node
- )
- for _, sealer := range sealers {
- // Start the node and wait until it's up
- stack, ethBackend, err := makeSealer(genesis)
- if err != nil {
- panic(err)
- }
- defer stack.Close()
-
- for stack.Server().NodeInfo().Ports.Listener == 0 {
- time.Sleep(250 * time.Millisecond)
- }
- // Connect the node to all the previous ones
- for _, n := range enodes {
- stack.Server().AddPeer(n)
- }
- // Start tracking the node and its enode
- stacks = append(stacks, stack)
- nodes = append(nodes, ethBackend)
- enodes = append(enodes, stack.Server().Self())
-
- // Inject the signer key and start sealing with it
- ks := keystore.NewKeyStore(stack.KeyStoreDir(), keystore.LightScryptN, keystore.LightScryptP)
- signer, err := ks.ImportECDSA(sealer, "")
- if err != nil {
- panic(err)
- }
- if err := ks.Unlock(signer, ""); err != nil {
- panic(err)
- }
- stack.AccountManager().AddBackend(ks)
- }
-
- // Iterate over all the nodes and start signing on them
- time.Sleep(3 * time.Second)
- for _, node := range nodes {
- if err := node.StartMining(); err != nil {
- panic(err)
- }
- }
- time.Sleep(3 * time.Second)
-
- // Start injecting transactions from the faucet like crazy
- nonces := make([]uint64, len(faucets))
- for {
- // Stop when interrupted.
- select {
- case <-interruptCh:
- for _, node := range stacks {
- node.Close()
- }
- return
- default:
- }
-
- // Pick a random signer node
- index := rand.Intn(len(faucets))
- backend := nodes[index%len(nodes)]
-
- // Create a self transaction and inject into the pool
- tx, err := types.SignTx(types.NewTransaction(nonces[index], crypto.PubkeyToAddress(faucets[index].PublicKey), new(big.Int), 21000, big.NewInt(100000000000), nil), types.HomesteadSigner{}, faucets[index])
- if err != nil {
- panic(err)
- }
- if err := backend.TxPool().Add([]*types.Transaction{tx}, true, false); err != nil {
- panic(err)
- }
- nonces[index]++
-
- // Wait if we're too saturated
- if pend, _ := backend.TxPool().Stats(); pend > 2048 {
- time.Sleep(100 * time.Millisecond)
- }
- }
-}
-
-// makeGenesis creates a custom Clique genesis block based on some pre-defined
-// signer and faucet accounts.
-func makeGenesis(faucets []*ecdsa.PrivateKey, sealers []*ecdsa.PrivateKey) *core.Genesis {
- // Create a Clique network based off of the Sepolia config
- genesis := core.DefaultSepoliaGenesisBlock()
- genesis.GasLimit = 25000000
-
- genesis.Config.ChainID = big.NewInt(18)
- genesis.Config.Clique.Period = 1
-
- genesis.Alloc = core.GenesisAlloc{}
- for _, faucet := range faucets {
- genesis.Alloc[crypto.PubkeyToAddress(faucet.PublicKey)] = core.GenesisAccount{
- Balance: new(big.Int).Exp(big.NewInt(2), big.NewInt(128), nil),
- }
- }
- // Sort the signers and embed into the extra-data section
- signers := make([]common.Address, len(sealers))
- for i, sealer := range sealers {
- signers[i] = crypto.PubkeyToAddress(sealer.PublicKey)
- }
- for i := 0; i < len(signers); i++ {
- for j := i + 1; j < len(signers); j++ {
- if bytes.Compare(signers[i][:], signers[j][:]) > 0 {
- signers[i], signers[j] = signers[j], signers[i]
- }
- }
- }
- genesis.ExtraData = make([]byte, 32+len(signers)*common.AddressLength+65)
- for i, signer := range signers {
- copy(genesis.ExtraData[32+i*common.AddressLength:], signer[:])
- }
- // Return the genesis block for initialization
- return genesis
-}
-
-func makeSealer(genesis *core.Genesis) (*node.Node, *eth.Ethereum, error) {
- // Define the basic configurations for the Ethereum node
- datadir, _ := os.MkdirTemp("", "")
-
- config := &node.Config{
- Name: "geth",
- Version: params.Version,
- DataDir: datadir,
- P2P: p2p.Config{
- ListenAddr: "0.0.0.0:0",
- NoDiscovery: true,
- MaxPeers: 25,
- },
- }
- // Start the node and configure a full Ethereum node on it
- stack, err := node.New(config)
- if err != nil {
- return nil, nil, err
- }
- // Create and register the backend
- ethBackend, err := eth.New(stack, ðconfig.Config{
- Genesis: genesis,
- NetworkId: genesis.Config.ChainID.Uint64(),
- SyncMode: downloader.FullSync,
- DatabaseCache: 256,
- DatabaseHandles: 256,
- TxPool: legacypool.DefaultConfig,
- GPO: ethconfig.Defaults.GPO,
- Miner: miner.Config{
- GasCeil: genesis.GasLimit * 11 / 10,
- GasPrice: big.NewInt(1),
- Recommit: time.Second,
- },
- })
- if err != nil {
- return nil, nil, err
- }
-
- err = stack.Start()
- return stack, ethBackend, err
-}
diff --git a/miner/worker.go b/miner/worker.go
deleted file mode 100644
index 2ed91cc187..0000000000
--- a/miner/worker.go
+++ /dev/null
@@ -1,1214 +0,0 @@
-// Copyright 2015 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see .
-
-package miner
-
-import (
- "errors"
- "fmt"
- "math/big"
- "sync"
- "sync/atomic"
- "time"
-
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/consensus"
- "github.com/ethereum/go-ethereum/consensus/misc/eip1559"
- "github.com/ethereum/go-ethereum/consensus/misc/eip4844"
- "github.com/ethereum/go-ethereum/core"
- "github.com/ethereum/go-ethereum/core/state"
- "github.com/ethereum/go-ethereum/core/txpool"
- "github.com/ethereum/go-ethereum/core/types"
- "github.com/ethereum/go-ethereum/core/vm"
- "github.com/ethereum/go-ethereum/event"
- "github.com/ethereum/go-ethereum/log"
- "github.com/ethereum/go-ethereum/params"
- "github.com/ethereum/go-ethereum/trie"
-)
-
-const (
- // resultQueueSize is the size of channel listening to sealing result.
- resultQueueSize = 10
-
- // txChanSize is the size of channel listening to NewTxsEvent.
- // The number is referenced from the size of tx pool.
- txChanSize = 4096
-
- // chainHeadChanSize is the size of channel listening to ChainHeadEvent.
- chainHeadChanSize = 10
-
- // resubmitAdjustChanSize is the size of resubmitting interval adjustment channel.
- resubmitAdjustChanSize = 10
-
- // minRecommitInterval is the minimal time interval to recreate the sealing block with
- // any newly arrived transactions.
- minRecommitInterval = 1 * time.Second
-
- // maxRecommitInterval is the maximum time interval to recreate the sealing block with
- // any newly arrived transactions.
- maxRecommitInterval = 15 * time.Second
-
- // intervalAdjustRatio is the impact a single interval adjustment has on sealing work
- // resubmitting interval.
- intervalAdjustRatio = 0.1
-
- // intervalAdjustBias is applied during the new resubmit interval calculation in favor of
- // increasing upper limit or decreasing lower limit so that the limit can be reachable.
- intervalAdjustBias = 200 * 1000.0 * 1000.0
-
- // staleThreshold is the maximum depth of the acceptable stale block.
- staleThreshold = 7
-)
-
-var (
- errBlockInterruptedByNewHead = errors.New("new head arrived while building block")
- errBlockInterruptedByRecommit = errors.New("recommit interrupt while building block")
- errBlockInterruptedByTimeout = errors.New("timeout while building block")
-)
-
-// environment is the worker's current environment and holds all
-// information of the sealing block generation.
-type environment struct {
- signer types.Signer
- state *state.StateDB // apply state changes here
- tcount int // tx count in cycle
- gasPool *core.GasPool // available gas used to pack transactions
- coinbase common.Address
-
- header *types.Header
- txs []*types.Transaction
- receipts []*types.Receipt
- sidecars []*types.BlobTxSidecar
- blobs int
-}
-
-// copy creates a deep copy of environment.
-func (env *environment) copy() *environment {
- cpy := &environment{
- signer: env.signer,
- state: env.state.Copy(),
- tcount: env.tcount,
- coinbase: env.coinbase,
- header: types.CopyHeader(env.header),
- receipts: copyReceipts(env.receipts),
- }
- if env.gasPool != nil {
- gasPool := *env.gasPool
- cpy.gasPool = &gasPool
- }
- cpy.txs = make([]*types.Transaction, len(env.txs))
- copy(cpy.txs, env.txs)
-
- cpy.sidecars = make([]*types.BlobTxSidecar, len(env.sidecars))
- copy(cpy.sidecars, env.sidecars)
-
- return cpy
-}
-
-// discard terminates the background prefetcher go-routine. It should
-// always be called for all created environment instances otherwise
-// the go-routine leak can happen.
-func (env *environment) discard() {
- if env.state == nil {
- return
- }
- env.state.StopPrefetcher()
-}
-
-// task contains all information for consensus engine sealing and result submitting.
-type task struct {
- receipts []*types.Receipt
- state *state.StateDB
- block *types.Block
- createdAt time.Time
-}
-
-const (
- commitInterruptNone int32 = iota
- commitInterruptNewHead
- commitInterruptResubmit
- commitInterruptTimeout
-)
-
-// newWorkReq represents a request for new sealing work submitting with relative interrupt notifier.
-type newWorkReq struct {
- interrupt *atomic.Int32
- timestamp int64
-}
-
-// newPayloadResult is the result of payload generation.
-type newPayloadResult struct {
- err error
- block *types.Block
- fees *big.Int // total block fees
- sidecars []*types.BlobTxSidecar // collected blobs of blob transactions
-}
-
-// getWorkReq represents a request for getting a new sealing work with provided parameters.
-type getWorkReq struct {
- params *generateParams
- result chan *newPayloadResult // non-blocking channel
-}
-
-// intervalAdjust represents a resubmitting interval adjustment.
-type intervalAdjust struct {
- ratio float64
- inc bool
-}
-
-// 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
-
- // Feeds
- pendingLogsFeed event.Feed
-
- // Subscriptions
- mux *event.TypeMux
- txsCh chan core.NewTxsEvent
- txsSub event.Subscription
- chainHeadCh chan core.ChainHeadEvent
- chainHeadSub event.Subscription
-
- // Channels
- newWorkCh chan *newWorkReq
- getWorkCh chan *getWorkReq
- taskCh chan *task
- resultCh chan *types.Block
- startCh chan struct{}
- exitCh chan struct{}
- resubmitIntervalCh chan time.Duration
- resubmitAdjustCh chan *intervalAdjust
-
- wg sync.WaitGroup
-
- current *environment // An environment for current running cycle.
-
- mu sync.RWMutex // The lock used to protect the coinbase and extra fields
- coinbase common.Address
- extra []byte
-
- pendingMu sync.RWMutex
- pendingTasks map[common.Hash]*task
-
- snapshotMu sync.RWMutex // The lock used to protect the snapshots below
- snapshotBlock *types.Block
- snapshotReceipts types.Receipts
- snapshotState *state.StateDB
-
- // atomic status counters
- running atomic.Bool // The indicator whether the consensus engine is running or not.
- newTxs atomic.Int32 // New arrival transaction count since last sealing work submitting.
- syncing atomic.Bool // The indicator whether the node is still syncing.
-
- // newpayloadTimeout is the maximum timeout allowance for creating payload.
- // The default value is 2 seconds but node operator can set it to arbitrary
- // large value. A large timeout allowance may cause Geth to fail creating
- // a non-empty payload within the specified time and eventually miss the slot
- // in case there are some computation expensive transactions in txpool.
- newpayloadTimeout time.Duration
-
- // recommit is the time interval to re-create sealing work or to re-build
- // payload in proof-of-stake stage.
- recommit time.Duration
-
- // External functions
- isLocalBlock func(header *types.Header) bool // Function used to determine whether the specified block is mined by local miner.
-
- // Test hooks
- newTaskHook func(*task) // Method to call upon receiving a new sealing task.
- skipSealHook func(*task) bool // Method to decide whether skipping the sealing.
- fullTaskHook func() // Method to call before pushing the full sealing task.
- resubmitHook func(time.Duration, time.Duration) // Method to call upon updating resubmitting interval.
-}
-
-func newWorker(config *Config, chainConfig *params.ChainConfig, engine consensus.Engine, eth Backend, mux *event.TypeMux, isLocalBlock func(header *types.Header) bool, init bool) *worker {
- worker := &worker{
- config: config,
- chainConfig: chainConfig,
- engine: engine,
- eth: eth,
- chain: eth.BlockChain(),
- mux: mux,
- isLocalBlock: isLocalBlock,
- coinbase: config.Etherbase,
- extra: config.ExtraData,
- pendingTasks: make(map[common.Hash]*task),
- txsCh: make(chan core.NewTxsEvent, txChanSize),
- chainHeadCh: make(chan core.ChainHeadEvent, chainHeadChanSize),
- newWorkCh: make(chan *newWorkReq),
- getWorkCh: make(chan *getWorkReq),
- taskCh: make(chan *task),
- resultCh: make(chan *types.Block, resultQueueSize),
- startCh: make(chan struct{}, 1),
- exitCh: make(chan struct{}),
- resubmitIntervalCh: make(chan time.Duration),
- resubmitAdjustCh: make(chan *intervalAdjust, resubmitAdjustChanSize),
- }
- // Subscribe for transaction insertion events (whether from network or resurrects)
- worker.txsSub = eth.TxPool().SubscribeTransactions(worker.txsCh, true)
- // Subscribe events for blockchain
- worker.chainHeadSub = eth.BlockChain().SubscribeChainHeadEvent(worker.chainHeadCh)
-
- // 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.recommit = recommit
-
- // Sanitize the timeout config for creating payload.
- newpayloadTimeout := worker.config.NewPayloadTimeout
- if newpayloadTimeout == 0 {
- log.Warn("Sanitizing new payload timeout to default", "provided", newpayloadTimeout, "updated", DefaultConfig.NewPayloadTimeout)
- newpayloadTimeout = DefaultConfig.NewPayloadTimeout
- }
- if newpayloadTimeout < time.Millisecond*100 {
- log.Warn("Low payload timeout may cause high amount of non-full blocks", "provided", newpayloadTimeout, "default", DefaultConfig.NewPayloadTimeout)
- }
- worker.newpayloadTimeout = newpayloadTimeout
-
- 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
-}
-
-// setEtherbase sets the etherbase used to initialize the block coinbase field.
-func (w *worker) setEtherbase(addr common.Address) {
- w.mu.Lock()
- defer w.mu.Unlock()
- w.coinbase = addr
-}
-
-// etherbase retrieves the configured etherbase address.
-func (w *worker) etherbase() common.Address {
- w.mu.RLock()
- defer w.mu.RUnlock()
- return w.coinbase
-}
-
-func (w *worker) setGasCeil(ceil uint64) {
- w.mu.Lock()
- defer w.mu.Unlock()
- w.config.GasCeil = ceil
-}
-
-// setExtra sets the content used to initialize the block extra field.
-func (w *worker) setExtra(extra []byte) {
- w.mu.Lock()
- defer w.mu.Unlock()
- w.extra = extra
-}
-
-// setRecommitInterval updates the interval for miner sealing work recommitting.
-func (w *worker) setRecommitInterval(interval time.Duration) {
- select {
- case w.resubmitIntervalCh <- interval:
- case <-w.exitCh:
- }
-}
-
-// pending returns the pending state and corresponding block. The returned
-// values can be nil in case the pending block is not initialized.
-func (w *worker) pending() (*types.Block, *state.StateDB) {
- w.snapshotMu.RLock()
- defer w.snapshotMu.RUnlock()
- if w.snapshotState == nil {
- return nil, nil
- }
- return w.snapshotBlock, w.snapshotState.Copy()
-}
-
-// pendingBlock returns pending block. The returned block can be nil in case the
-// pending block is not initialized.
-func (w *worker) pendingBlock() *types.Block {
- w.snapshotMu.RLock()
- defer w.snapshotMu.RUnlock()
- return w.snapshotBlock
-}
-
-// pendingBlockAndReceipts returns pending block and corresponding receipts.
-// The returned values can be nil in case the pending block is not initialized.
-func (w *worker) pendingBlockAndReceipts() (*types.Block, types.Receipts) {
- w.snapshotMu.RLock()
- defer w.snapshotMu.RUnlock()
- return w.snapshotBlock, w.snapshotReceipts
-}
-
-// start sets the running status as 1 and triggers new work submitting.
-func (w *worker) start() {
- w.running.Store(true)
- w.startCh <- struct{}{}
-}
-
-// stop sets the running status as 0.
-func (w *worker) stop() {
- w.running.Store(false)
-}
-
-// isRunning returns an indicator whether worker is running or not.
-func (w *worker) isRunning() bool {
- return w.running.Load()
-}
-
-// close terminates all background threads maintained by the worker.
-// Note the worker does not support being closed multiple times.
-func (w *worker) close() {
- w.running.Store(false)
- close(w.exitCh)
- w.wg.Wait()
-}
-
-// recalcRecommit recalculates the resubmitting interval upon feedback.
-func recalcRecommit(minRecommit, prev time.Duration, target float64, inc bool) time.Duration {
- var (
- prevF = float64(prev.Nanoseconds())
- next float64
- )
- if inc {
- next = prevF*(1-intervalAdjustRatio) + intervalAdjustRatio*(target+intervalAdjustBias)
- max := float64(maxRecommitInterval.Nanoseconds())
- if next > max {
- next = max
- }
- } else {
- next = prevF*(1-intervalAdjustRatio) + intervalAdjustRatio*(target-intervalAdjustBias)
- min := float64(minRecommit.Nanoseconds())
- if next < min {
- next = min
- }
- }
- return time.Duration(int64(next))
-}
-
-// newWorkLoop is a standalone goroutine to submit new sealing work upon received events.
-func (w *worker) newWorkLoop(recommit time.Duration) {
- defer w.wg.Done()
- var (
- interrupt *atomic.Int32
- minRecommit = recommit // minimal resubmit interval specified by user.
- timestamp int64 // timestamp for each round of sealing.
- )
-
- timer := time.NewTimer(0)
- defer timer.Stop()
- <-timer.C // discard the initial tick
-
- // commit aborts in-flight transaction execution with given signal and resubmits a new one.
- commit := func(s int32) {
- if interrupt != nil {
- interrupt.Store(s)
- }
- interrupt = new(atomic.Int32)
- select {
- case w.newWorkCh <- &newWorkReq{interrupt: interrupt, timestamp: timestamp}:
- case <-w.exitCh:
- return
- }
- timer.Reset(recommit)
- w.newTxs.Store(0)
- }
- // clearPending cleans the stale pending tasks.
- clearPending := func(number uint64) {
- w.pendingMu.Lock()
- for h, t := range w.pendingTasks {
- if t.block.NumberU64()+staleThreshold <= number {
- delete(w.pendingTasks, h)
- }
- }
- w.pendingMu.Unlock()
- }
-
- for {
- select {
- case <-w.startCh:
- clearPending(w.chain.CurrentBlock().Number.Uint64())
- timestamp = time.Now().Unix()
- commit(commitInterruptNewHead)
-
- case head := <-w.chainHeadCh:
- clearPending(head.Block.NumberU64())
- timestamp = time.Now().Unix()
- commit(commitInterruptNewHead)
-
- case <-timer.C:
- // If sealing 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 w.newTxs.Load() == 0 {
- timer.Reset(recommit)
- continue
- }
- commit(commitInterruptResubmit)
- }
-
- case interval := <-w.resubmitIntervalCh:
- // Adjust resubmit interval explicitly by user.
- if interval < minRecommitInterval {
- log.Warn("Sanitizing miner recommit interval", "provided", interval, "updated", minRecommitInterval)
- interval = minRecommitInterval
- }
- log.Info("Miner recommit interval update", "from", minRecommit, "to", interval)
- minRecommit, recommit = interval, interval
-
- if w.resubmitHook != nil {
- w.resubmitHook(minRecommit, recommit)
- }
-
- case adjust := <-w.resubmitAdjustCh:
- // Adjust resubmit interval by feedback.
- if adjust.inc {
- before := recommit
- target := float64(recommit.Nanoseconds()) / adjust.ratio
- recommit = recalcRecommit(minRecommit, recommit, target, true)
- log.Trace("Increase miner recommit interval", "from", before, "to", recommit)
- } else {
- before := recommit
- recommit = recalcRecommit(minRecommit, recommit, float64(minRecommit.Nanoseconds()), false)
- log.Trace("Decrease miner recommit interval", "from", before, "to", recommit)
- }
-
- if w.resubmitHook != nil {
- w.resubmitHook(minRecommit, recommit)
- }
-
- case <-w.exitCh:
- return
- }
- }
-}
-
-// mainLoop is responsible for generating and submitting sealing work based on
-// the received event. It can support two modes: automatically generate task and
-// submit it or return task according to given parameters for various proposes.
-func (w *worker) mainLoop() {
- defer w.wg.Done()
- defer w.txsSub.Unsubscribe()
- defer w.chainHeadSub.Unsubscribe()
- defer func() {
- if w.current != nil {
- w.current.discard()
- }
- }()
-
- for {
- select {
- case req := <-w.newWorkCh:
- w.commitWork(req.interrupt, req.timestamp)
-
- case req := <-w.getWorkCh:
- req.result <- w.generateWork(req.params)
-
- case ev := <-w.txsCh:
- // Apply transactions to the pending state if we're not sealing
- //
- // Note all transactions received may not be continuous with transactions
- // already included in the current sealing block. These transactions will
- // be automatically eliminated.
- if !w.isRunning() && w.current != nil {
- // If block is already full, abort
- if gp := w.current.gasPool; gp != nil && gp.Gas() < params.TxGas {
- continue
- }
- txs := make(map[common.Address][]*txpool.LazyTransaction, len(ev.Txs))
- for _, tx := range ev.Txs {
- acc, _ := types.Sender(w.current.signer, tx)
- txs[acc] = append(txs[acc], &txpool.LazyTransaction{
- Pool: w.eth.TxPool(), // We don't know where this came from, yolo resolve from everywhere
- Hash: tx.Hash(),
- Tx: nil, // Do *not* set this! We need to resolve it later to pull blobs in
- Time: tx.Time(),
- GasFeeCap: tx.GasFeeCap(),
- GasTipCap: tx.GasTipCap(),
- Gas: tx.Gas(),
- BlobGas: tx.BlobGas(),
- })
- }
- txset := newTransactionsByPriceAndNonce(w.current.signer, txs, w.current.header.BaseFee)
- tcount := w.current.tcount
- w.commitTransactions(w.current, txset, nil)
-
- // Only update the snapshot if any new transactions were added
- // to the pending block
- if tcount != w.current.tcount {
- w.updateSnapshot(w.current)
- }
- } else {
- // Special case, if the consensus engine is 0 period clique(dev mode),
- // submit sealing work here since all empty submission will be rejected
- // by clique. Of course the advance sealing(empty submission) is disabled.
- if w.chainConfig.Clique != nil && w.chainConfig.Clique.Period == 0 {
- w.commitWork(nil, time.Now().Unix())
- }
- }
- w.newTxs.Add(int32(len(ev.Txs)))
-
- // System stopped
- case <-w.exitCh:
- return
- case <-w.txsSub.Err():
- return
- case <-w.chainHeadSub.Err():
- return
- }
- }
-}
-
-// 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
- )
-
- // interrupt aborts the in-flight sealing task.
- interrupt := func() {
- if stopCh != nil {
- close(stopCh)
- stopCh = nil
- }
- }
- 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
- }
- }
-}
-
-// resultLoop is a standalone goroutine to handle sealing result submitting
-// and flush relative data to the database.
-func (w *worker) resultLoop() {
- defer w.wg.Done()
- for {
- select {
- case block := <-w.resultCh:
- // Short circuit when receiving empty result.
- if block == nil {
- continue
- }
- // Short circuit when receiving duplicate result caused by resubmitting.
- if w.chain.HasBlock(block.Hash(), block.NumberU64()) {
- continue
- }
- var (
- sealhash = w.engine.SealHash(block.Header())
- hash = block.Hash()
- )
- w.pendingMu.RLock()
- task, exist := w.pendingTasks[sealhash]
- w.pendingMu.RUnlock()
- if !exist {
- log.Error("Block found but no relative pending task", "number", block.Number(), "sealhash", sealhash, "hash", hash)
- continue
- }
- // Different block could share same sealhash, deep copy here to prevent write-write conflict.
- var (
- receipts = make([]*types.Receipt, len(task.receipts))
- logs []*types.Log
- )
- for i, taskReceipt := range task.receipts {
- receipt := new(types.Receipt)
- receipts[i] = receipt
- *receipt = *taskReceipt
-
- // add block location fields
- receipt.BlockHash = hash
- receipt.BlockNumber = block.Number()
- receipt.TransactionIndex = uint(i)
-
- // Update the block hash in all logs since it is now available and not when the
- // receipt/log of individual transactions were created.
- receipt.Logs = make([]*types.Log, len(taskReceipt.Logs))
- for i, taskLog := range taskReceipt.Logs {
- log := new(types.Log)
- receipt.Logs[i] = log
- *log = *taskLog
- log.BlockHash = hash
- }
- logs = append(logs, receipt.Logs...)
- }
- // Commit block and state to database.
- _, err := w.chain.WriteBlockAndSetHead(block, receipts, logs, task.state, true)
- if err != nil {
- log.Error("Failed writing block to chain", "err", err)
- continue
- }
- log.Info("Successfully sealed new block", "number", block.Number(), "sealhash", sealhash, "hash", hash,
- "elapsed", common.PrettyDuration(time.Since(task.createdAt)))
-
- // Broadcast the block and announce chain insertion event
- w.mux.Post(core.NewMinedBlockEvent{Block: block})
-
- case <-w.exitCh:
- return
- }
- }
-}
-
-// makeEnv creates a new environment for the sealing block.
-func (w *worker) makeEnv(parent *types.Header, header *types.Header, coinbase common.Address) (*environment, error) {
- // Retrieve the parent state to execute on top and start a prefetcher for
- // the miner to speed block sealing up a bit.
- state, err := w.chain.StateAt(parent.Root)
- if err != nil {
- return nil, err
- }
- state.StartPrefetcher("miner")
-
- // Note the passed coinbase may be different with header.Coinbase.
- env := &environment{
- signer: types.MakeSigner(w.chainConfig, header.Number, header.Time),
- state: state,
- coinbase: coinbase,
- header: header,
- }
- // Keep track of transactions which return errors so they can be removed
- env.tcount = 0
- return env, nil
-}
-
-// updateSnapshot updates pending snapshot block, receipts and state.
-func (w *worker) updateSnapshot(env *environment) {
- w.snapshotMu.Lock()
- defer w.snapshotMu.Unlock()
-
- w.snapshotBlock = types.NewBlock(
- env.header,
- env.txs,
- nil,
- env.receipts,
- trie.NewStackTrie(nil),
- )
- w.snapshotReceipts = copyReceipts(env.receipts)
- w.snapshotState = env.state.Copy()
-}
-
-func (w *worker) commitTransaction(env *environment, tx *types.Transaction) ([]*types.Log, error) {
- if tx.Type() == types.BlobTxType {
- return w.commitBlobTransaction(env, tx)
- }
- receipt, err := w.applyTransaction(env, tx)
- if err != nil {
- return nil, err
- }
- env.txs = append(env.txs, tx)
- env.receipts = append(env.receipts, receipt)
- return receipt.Logs, nil
-}
-
-func (w *worker) commitBlobTransaction(env *environment, tx *types.Transaction) ([]*types.Log, error) {
- sc := tx.BlobTxSidecar()
- if sc == nil {
- panic("blob transaction without blobs in miner")
- }
- // Checking against blob gas limit: It's kind of ugly to perform this check here, but there
- // isn't really a better place right now. The blob gas limit is checked at block validation time
- // and not during execution. This means core.ApplyTransaction will not return an error if the
- // tx has too many blobs. So we have to explicitly check it here.
- if (env.blobs+len(sc.Blobs))*params.BlobTxBlobGasPerBlob > params.MaxBlobGasPerBlock {
- return nil, errors.New("max data blobs reached")
- }
- receipt, err := w.applyTransaction(env, tx)
- if err != nil {
- return nil, err
- }
- env.txs = append(env.txs, tx.WithoutBlobTxSidecar())
- env.receipts = append(env.receipts, receipt)
- env.sidecars = append(env.sidecars, sc)
- env.blobs += len(sc.Blobs)
- *env.header.BlobGasUsed += receipt.BlobGasUsed
- return receipt.Logs, nil
-}
-
-// applyTransaction runs the transaction. If execution fails, state and gas pool are reverted.
-func (w *worker) applyTransaction(env *environment, tx *types.Transaction) (*types.Receipt, error) {
- var (
- snap = env.state.Snapshot()
- gp = env.gasPool.Gas()
- )
- receipt, err := core.ApplyTransaction(w.chainConfig, w.chain, &env.coinbase, env.gasPool, env.state, env.header, tx, &env.header.GasUsed, *w.chain.GetVMConfig())
- if err != nil {
- env.state.RevertToSnapshot(snap)
- env.gasPool.SetGas(gp)
- }
- return receipt, err
-}
-
-func (w *worker) commitTransactions(env *environment, txs *transactionsByPriceAndNonce, interrupt *atomic.Int32) error {
- gasLimit := env.header.GasLimit
- if env.gasPool == nil {
- env.gasPool = new(core.GasPool).AddGas(gasLimit)
- }
- var coalescedLogs []*types.Log
-
- for {
- // Check interruption signal and abort building if it's fired.
- if interrupt != nil {
- if signal := interrupt.Load(); signal != commitInterruptNone {
- return signalToErr(signal)
- }
- }
- // If we don't have enough gas for any further transactions then we're done.
- if env.gasPool.Gas() < params.TxGas {
- log.Trace("Not enough gas for further transactions", "have", env.gasPool, "want", params.TxGas)
- break
- }
- // Retrieve the next transaction and abort if all done.
- ltx := txs.Peek()
- if ltx == nil {
- break
- }
- // If we don't have enough space for the next transaction, skip the account.
- if env.gasPool.Gas() < ltx.Gas {
- log.Trace("Not enough gas left for transaction", "hash", ltx.Hash, "left", env.gasPool.Gas(), "needed", ltx.Gas)
- txs.Pop()
- continue
- }
- if left := uint64(params.MaxBlobGasPerBlock - env.blobs*params.BlobTxBlobGasPerBlob); left < ltx.BlobGas {
- log.Trace("Not enough blob gas left for transaction", "hash", ltx.Hash, "left", left, "needed", ltx.BlobGas)
- txs.Pop()
- continue
- }
- // Transaction seems to fit, pull it up from the pool
- tx := ltx.Resolve()
- if tx == nil {
- log.Trace("Ignoring evicted transaction", "hash", ltx.Hash)
- txs.Pop()
- continue
- }
- // Error may be ignored here. The error has already been checked
- // during transaction acceptance is the transaction pool.
- from, _ := types.Sender(env.signer, tx)
-
- // Check whether the tx is replay protected. If we're not in the EIP155 hf
- // phase, start ignoring the sender until we do.
- if tx.Protected() && !w.chainConfig.IsEIP155(env.header.Number) {
- log.Trace("Ignoring replay protected transaction", "hash", ltx.Hash, "eip155", w.chainConfig.EIP155Block)
- txs.Pop()
- continue
- }
- // Start executing the transaction
- env.state.SetTxContext(tx.Hash(), env.tcount)
-
- logs, err := w.commitTransaction(env, tx)
- switch {
- case errors.Is(err, core.ErrNonceTooLow):
- // New head notification data race between the transaction pool and miner, shift
- log.Trace("Skipping transaction with low nonce", "hash", ltx.Hash, "sender", from, "nonce", tx.Nonce())
- txs.Shift()
-
- case errors.Is(err, nil):
- // Everything ok, collect the logs and shift in the next transaction from the same account
- coalescedLogs = append(coalescedLogs, logs...)
- env.tcount++
- txs.Shift()
-
- default:
- // Transaction is regarded as invalid, drop all consecutive transactions from
- // the same sender because of `nonce-too-high` clause.
- log.Debug("Transaction failed, account skipped", "hash", ltx.Hash, "err", err)
- txs.Pop()
- }
- }
- if !w.isRunning() && len(coalescedLogs) > 0 {
- // We don't push the pendingLogsEvent while we are sealing. The reason is that
- // when we are sealing, the worker will regenerate a sealing block every 3 seconds.
- // In order to avoid pushing the repeated pendingLog, we disable the pending log pushing.
-
- // make a copy, the state caches the logs and these logs get "upgraded" from pending to mined
- // logs by filling in the block hash when the block was mined by the local miner. This can
- // cause a race condition if a log was "upgraded" before the PendingLogsEvent is processed.
- cpy := make([]*types.Log, len(coalescedLogs))
- for i, l := range coalescedLogs {
- cpy[i] = new(types.Log)
- *cpy[i] = *l
- }
- w.pendingLogsFeed.Send(cpy)
- }
- return nil
-}
-
-// generateParams wraps various of settings for generating sealing task.
-type generateParams struct {
- timestamp uint64 // The timstamp for sealing task
- forceTime bool // Flag whether the given timestamp is immutable or not
- parentHash common.Hash // Parent block hash, empty means the latest chain head
- coinbase common.Address // The fee recipient address for including transaction
- random common.Hash // The randomness generated by beacon chain, empty before the merge
- withdrawals types.Withdrawals // List of withdrawals to include in block.
- beaconRoot *common.Hash // The beacon root (cancun field).
- noTxs bool // Flag whether an empty block without any transaction is expected
-}
-
-// prepareWork constructs the sealing task according to the given parameters,
-// either based on the last chain head or specified parent. In this function
-// the pending transactions are not filled yet, only the empty task returned.
-func (w *worker) prepareWork(genParams *generateParams) (*environment, error) {
- w.mu.RLock()
- defer w.mu.RUnlock()
-
- // Find the parent block for sealing task
- parent := w.chain.CurrentBlock()
- if genParams.parentHash != (common.Hash{}) {
- block := w.chain.GetBlockByHash(genParams.parentHash)
- if block == nil {
- return nil, fmt.Errorf("missing parent")
- }
- parent = block.Header()
- }
- // Sanity check the timestamp correctness, recap the timestamp
- // to parent+1 if the mutation is allowed.
- timestamp := genParams.timestamp
- if parent.Time >= timestamp {
- if genParams.forceTime {
- return nil, fmt.Errorf("invalid timestamp, parent %d given %d", parent.Time, timestamp)
- }
- timestamp = parent.Time + 1
- }
- // Construct the sealing block header.
- header := &types.Header{
- ParentHash: parent.Hash(),
- Number: new(big.Int).Add(parent.Number, common.Big1),
- GasLimit: core.CalcGasLimit(parent.GasLimit, w.config.GasCeil),
- Time: timestamp,
- Coinbase: genParams.coinbase,
- }
- // Set the extra field.
- if len(w.extra) != 0 {
- header.Extra = w.extra
- }
- // Set the randomness field from the beacon chain if it's available.
- if genParams.random != (common.Hash{}) {
- header.MixDigest = genParams.random
- }
- // Set baseFee and GasLimit if we are on an EIP-1559 chain
- if w.chainConfig.IsLondon(header.Number) {
- header.BaseFee = eip1559.CalcBaseFee(w.chainConfig, parent)
- if !w.chainConfig.IsLondon(parent.Number) {
- parentGasLimit := parent.GasLimit * w.chainConfig.ElasticityMultiplier()
- header.GasLimit = core.CalcGasLimit(parentGasLimit, w.config.GasCeil)
- }
- }
- // Apply EIP-4844, EIP-4788.
- if w.chainConfig.IsCancun(header.Number, header.Time) {
- var excessBlobGas uint64
- if w.chainConfig.IsCancun(parent.Number, parent.Time) {
- excessBlobGas = eip4844.CalcExcessBlobGas(*parent.ExcessBlobGas, *parent.BlobGasUsed)
- } else {
- // For the first post-fork block, both parent.data_gas_used and parent.excess_data_gas are evaluated as 0
- excessBlobGas = eip4844.CalcExcessBlobGas(0, 0)
- }
- header.BlobGasUsed = new(uint64)
- header.ExcessBlobGas = &excessBlobGas
- header.ParentBeaconRoot = genParams.beaconRoot
- }
- // Run the consensus preparation with the default or customized consensus engine.
- if err := w.engine.Prepare(w.chain, header); err != nil {
- log.Error("Failed to prepare header for sealing", "err", err)
- return nil, err
- }
- // Could potentially happen if starting to mine in an odd state.
- // Note genParams.coinbase can be different with header.Coinbase
- // since clique algorithm can modify the coinbase field in header.
- env, err := w.makeEnv(parent, header, genParams.coinbase)
- if err != nil {
- log.Error("Failed to create sealing context", "err", err)
- return nil, err
- }
- if header.ParentBeaconRoot != nil {
- context := core.NewEVMBlockContext(header, w.chain, nil)
- vmenv := vm.NewEVM(context, vm.TxContext{}, env.state, w.chainConfig, vm.Config{})
- core.ProcessBeaconBlockRoot(*header.ParentBeaconRoot, vmenv, env.state)
- }
- return env, nil
-}
-
-// fillTransactions retrieves the pending transactions from the txpool and fills them
-// into the given sealing block. The transaction selection and ordering strategy can
-// be customized with the plugin in the future.
-func (w *worker) fillTransactions(interrupt *atomic.Int32, env *environment) error {
- pending := w.eth.TxPool().Pending(true)
-
- // Split the pending transactions into locals and remotes.
- localTxs, remoteTxs := make(map[common.Address][]*txpool.LazyTransaction), pending
- for _, account := range w.eth.TxPool().Locals() {
- if txs := remoteTxs[account]; len(txs) > 0 {
- delete(remoteTxs, account)
- localTxs[account] = txs
- }
- }
-
- // Fill the block with all available pending transactions.
- if len(localTxs) > 0 {
- txs := newTransactionsByPriceAndNonce(env.signer, localTxs, env.header.BaseFee)
- if err := w.commitTransactions(env, txs, interrupt); err != nil {
- return err
- }
- }
- if len(remoteTxs) > 0 {
- txs := newTransactionsByPriceAndNonce(env.signer, remoteTxs, env.header.BaseFee)
- if err := w.commitTransactions(env, txs, interrupt); err != nil {
- return err
- }
- }
- return nil
-}
-
-// generateWork generates a sealing block based on the given parameters.
-func (w *worker) generateWork(params *generateParams) *newPayloadResult {
- work, err := w.prepareWork(params)
- if err != nil {
- return &newPayloadResult{err: err}
- }
- defer work.discard()
-
- if !params.noTxs {
- interrupt := new(atomic.Int32)
- timer := time.AfterFunc(w.newpayloadTimeout, func() {
- interrupt.Store(commitInterruptTimeout)
- })
- defer timer.Stop()
-
- err := w.fillTransactions(interrupt, work)
- if errors.Is(err, errBlockInterruptedByTimeout) {
- log.Warn("Block building is interrupted", "allowance", common.PrettyDuration(w.newpayloadTimeout))
- }
- }
- block, err := w.engine.FinalizeAndAssemble(w.chain, work.header, work.state, work.txs, nil, work.receipts, params.withdrawals)
- if err != nil {
- return &newPayloadResult{err: err}
- }
- return &newPayloadResult{
- block: block,
- fees: totalFees(block, work.receipts),
- sidecars: work.sidecars,
- }
-}
-
-// commitWork generates several new sealing tasks based on the parent block
-// and submit them to the sealer.
-func (w *worker) commitWork(interrupt *atomic.Int32, timestamp int64) {
- // Abort committing if node is still syncing
- if w.syncing.Load() {
- return
- }
- start := time.Now()
-
- // Set the coinbase if the worker is running or it's required
- var coinbase common.Address
- if w.isRunning() {
- coinbase = w.etherbase()
- if coinbase == (common.Address{}) {
- log.Error("Refusing to mine without etherbase")
- return
- }
- }
- work, err := w.prepareWork(&generateParams{
- timestamp: uint64(timestamp),
- coinbase: coinbase,
- })
- if err != nil {
- return
- }
- // Fill pending transactions from the txpool into the block.
- err = w.fillTransactions(interrupt, work)
- switch {
- case err == nil:
- // The entire block is filled, decrease resubmit interval in case
- // of current interval is larger than the user-specified one.
- w.adjustResubmitInterval(&intervalAdjust{inc: false})
-
- case errors.Is(err, errBlockInterruptedByRecommit):
- // Notify resubmit loop to increase resubmitting interval if the
- // interruption is due to frequent commits.
- gaslimit := work.header.GasLimit
- ratio := float64(gaslimit-work.gasPool.Gas()) / float64(gaslimit)
- if ratio < 0.1 {
- ratio = 0.1
- }
- w.adjustResubmitInterval(&intervalAdjust{
- ratio: ratio,
- inc: true,
- })
-
- case errors.Is(err, errBlockInterruptedByNewHead):
- // If the block building is interrupted by newhead event, discard it
- // totally. Committing the interrupted block introduces unnecessary
- // delay, and possibly causes miner to mine on the previous head,
- // which could result in higher uncle rate.
- work.discard()
- return
- }
- // Submit the generated block for consensus sealing.
- w.commit(work.copy(), w.fullTaskHook, true, start)
-
- // Swap out the old work with the new one, terminating any leftover
- // prefetcher processes in the mean time and starting a new one.
- if w.current != nil {
- w.current.discard()
- }
- w.current = work
-}
-
-// commit runs any post-transaction state modifications, assembles the final block
-// and commits new work if consensus engine is running.
-// Note the assumption is held that the mutation is allowed to the passed env, do
-// the deep copy first.
-func (w *worker) commit(env *environment, interval func(), update bool, start time.Time) error {
- if w.isRunning() {
- if interval != nil {
- interval()
- }
- // Create a local environment copy, avoid the data race with snapshot state.
- // https://github.com/ethereum/go-ethereum/issues/24299
- env := env.copy()
- // Withdrawals are set to nil here, because this is only called in PoW.
- block, err := w.engine.FinalizeAndAssemble(w.chain, env.header, env.state, env.txs, nil, env.receipts, nil)
- if err != nil {
- return err
- }
- // If we're post merge, just ignore
- if !w.isTTDReached(block.Header()) {
- select {
- case w.taskCh <- &task{receipts: env.receipts, state: env.state, block: block, createdAt: time.Now()}:
- fees := totalFees(block, env.receipts)
- feesInEther := new(big.Float).Quo(new(big.Float).SetInt(fees), big.NewFloat(params.Ether))
- log.Info("Commit new sealing work", "number", block.Number(), "sealhash", w.engine.SealHash(block.Header()),
- "txs", env.tcount, "gas", block.GasUsed(), "fees", feesInEther,
- "elapsed", common.PrettyDuration(time.Since(start)))
-
- case <-w.exitCh:
- log.Info("Worker has exited")
- }
- }
- }
- if update {
- w.updateSnapshot(env)
- }
- return nil
-}
-
-// getSealingBlock generates the sealing block based on the given parameters.
-// The generation result will be passed back via the given channel no matter
-// the generation itself succeeds or not.
-func (w *worker) getSealingBlock(params *generateParams) *newPayloadResult {
- req := &getWorkReq{
- params: params,
- result: make(chan *newPayloadResult, 1),
- }
- select {
- case w.getWorkCh <- req:
- return <-req.result
- case <-w.exitCh:
- return &newPayloadResult{err: errors.New("miner closed")}
- }
-}
-
-// isTTDReached returns the indicator if the given block has reached the total
-// terminal difficulty for The Merge transition.
-func (w *worker) isTTDReached(header *types.Header) bool {
- td, ttd := w.chain.GetTd(header.ParentHash, header.Number.Uint64()-1), w.chain.Config().TerminalTotalDifficulty
- return td != nil && ttd != nil && td.Cmp(ttd) >= 0
-}
-
-// adjustResubmitInterval adjusts the resubmit interval.
-func (w *worker) adjustResubmitInterval(message *intervalAdjust) {
- select {
- case w.resubmitAdjustCh <- message:
- default:
- log.Warn("the resubmitAdjustCh is full, discard the message")
- }
-}
-
-// copyReceipts makes a deep copy of the given receipts.
-func copyReceipts(receipts []*types.Receipt) []*types.Receipt {
- result := make([]*types.Receipt, len(receipts))
- for i, l := range receipts {
- cpy := *l
- result[i] = &cpy
- }
- return result
-}
-
-// totalFees computes total consumed miner fees in Wei. Block transactions and receipts have to have the same order.
-func totalFees(block *types.Block, receipts []*types.Receipt) *big.Int {
- feesWei := new(big.Int)
- for i, tx := range block.Transactions() {
- minerFee, _ := tx.EffectiveGasTip(block.BaseFee())
- feesWei.Add(feesWei, new(big.Int).Mul(new(big.Int).SetUint64(receipts[i].GasUsed), minerFee))
- }
- return feesWei
-}
-
-// signalToErr converts the interruption signal to a concrete error type for return.
-// The given signal must be a valid interruption signal.
-func signalToErr(signal int32) error {
- switch signal {
- case commitInterruptNewHead:
- return errBlockInterruptedByNewHead
- case commitInterruptResubmit:
- return errBlockInterruptedByRecommit
- case commitInterruptTimeout:
- return errBlockInterruptedByTimeout
- default:
- panic(fmt.Errorf("undefined signal %d", signal))
- }
-}
diff --git a/miner/worker_test.go b/miner/worker_test.go
deleted file mode 100644
index 59fbbbcdca..0000000000
--- a/miner/worker_test.go
+++ /dev/null
@@ -1,509 +0,0 @@
-// Copyright 2018 The go-ethereum Authors
-// This file is part of the go-ethereum library.
-//
-// The go-ethereum library is free software: you can redistribute it and/or modify
-// it under the terms of the GNU Lesser General Public License as published by
-// the Free Software Foundation, either version 3 of the License, or
-// (at your option) any later version.
-//
-// The go-ethereum library is distributed in the hope that it will be useful,
-// but WITHOUT ANY WARRANTY; without even the implied warranty of
-// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
-// GNU Lesser General Public License for more details.
-//
-// You should have received a copy of the GNU Lesser General Public License
-// along with the go-ethereum library. If not, see .
-
-package miner
-
-import (
- "math/big"
- "sync/atomic"
- "testing"
- "time"
-
- "github.com/ethereum/go-ethereum/accounts"
- "github.com/ethereum/go-ethereum/common"
- "github.com/ethereum/go-ethereum/consensus"
- "github.com/ethereum/go-ethereum/consensus/clique"
- "github.com/ethereum/go-ethereum/consensus/ethash"
- "github.com/ethereum/go-ethereum/core"
- "github.com/ethereum/go-ethereum/core/rawdb"
- "github.com/ethereum/go-ethereum/core/txpool"
- "github.com/ethereum/go-ethereum/core/txpool/legacypool"
- "github.com/ethereum/go-ethereum/core/types"
- "github.com/ethereum/go-ethereum/core/vm"
- "github.com/ethereum/go-ethereum/crypto"
- "github.com/ethereum/go-ethereum/ethdb"
- "github.com/ethereum/go-ethereum/event"
- "github.com/ethereum/go-ethereum/params"
-)
-
-const (
- // testCode is the testing contract binary code which will initialises some
- // variables in constructor
- testCode = "0x60806040527fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff0060005534801561003457600080fd5b5060fc806100436000396000f3fe6080604052348015600f57600080fd5b506004361060325760003560e01c80630c4dae8814603757806398a213cf146053575b600080fd5b603d607e565b6040518082815260200191505060405180910390f35b607c60048036036020811015606757600080fd5b81019080803590602001909291905050506084565b005b60005481565b806000819055507fe9e44f9f7da8c559de847a3232b57364adc0354f15a2cd8dc636d54396f9587a6000546040518082815260200191505060405180910390a15056fea265627a7a723058208ae31d9424f2d0bc2a3da1a5dd659db2d71ec322a17db8f87e19e209e3a1ff4a64736f6c634300050a0032"
-
- // testGas is the gas required for contract deployment.
- testGas = 144109
-)
-
-var (
- // Test chain configurations
- testTxPoolConfig legacypool.Config
- ethashChainConfig *params.ChainConfig
- cliqueChainConfig *params.ChainConfig
-
- // Test accounts
- testBankKey, _ = crypto.GenerateKey()
- testBankAddress = crypto.PubkeyToAddress(testBankKey.PublicKey)
- testBankFunds = big.NewInt(1000000000000000000)
-
- testUserKey, _ = crypto.GenerateKey()
- testUserAddress = crypto.PubkeyToAddress(testUserKey.PublicKey)
-
- // Test transactions
- pendingTxs []*types.Transaction
- newTxs []*types.Transaction
-
- testConfig = &Config{
- Recommit: time.Second,
- GasCeil: params.GenesisGasLimit,
- }
-)
-
-func init() {
- testTxPoolConfig = legacypool.DefaultConfig
- testTxPoolConfig.Journal = ""
- ethashChainConfig = new(params.ChainConfig)
- *ethashChainConfig = *params.TestChainConfig
- cliqueChainConfig = new(params.ChainConfig)
- *cliqueChainConfig = *params.TestChainConfig
- cliqueChainConfig.Clique = ¶ms.CliqueConfig{
- Period: 10,
- Epoch: 30000,
- }
-
- signer := types.LatestSigner(params.TestChainConfig)
- tx1 := types.MustSignNewTx(testBankKey, signer, &types.AccessListTx{
- ChainID: params.TestChainConfig.ChainID,
- Nonce: 0,
- To: &testUserAddress,
- Value: big.NewInt(1000),
- Gas: params.TxGas,
- GasPrice: big.NewInt(params.InitialBaseFee),
- })
- pendingTxs = append(pendingTxs, tx1)
-
- tx2 := types.MustSignNewTx(testBankKey, signer, &types.LegacyTx{
- Nonce: 1,
- To: &testUserAddress,
- Value: big.NewInt(1000),
- Gas: params.TxGas,
- GasPrice: big.NewInt(params.InitialBaseFee),
- })
- newTxs = append(newTxs, tx2)
-}
-
-// testWorkerBackend implements worker.Backend interfaces and wraps all information needed during the testing.
-type testWorkerBackend struct {
- db ethdb.Database
- txPool *txpool.TxPool
- chain *core.BlockChain
- genesis *core.Genesis
-}
-
-func newTestWorkerBackend(t *testing.T, chainConfig *params.ChainConfig, engine consensus.Engine, db ethdb.Database, n int) *testWorkerBackend {
- var gspec = &core.Genesis{
- Config: chainConfig,
- Alloc: core.GenesisAlloc{testBankAddress: {Balance: testBankFunds}},
- }
- switch e := engine.(type) {
- case *clique.Clique:
- gspec.ExtraData = make([]byte, 32+common.AddressLength+crypto.SignatureLength)
- copy(gspec.ExtraData[32:32+common.AddressLength], testBankAddress.Bytes())
- e.Authorize(testBankAddress, func(account accounts.Account, s string, data []byte) ([]byte, error) {
- return crypto.Sign(crypto.Keccak256(data), testBankKey)
- })
- case *ethash.Ethash:
- default:
- t.Fatalf("unexpected consensus engine type: %T", engine)
- }
- chain, err := core.NewBlockChain(db, &core.CacheConfig{TrieDirtyDisabled: true}, gspec, nil, engine, vm.Config{}, nil, nil)
- if err != nil {
- t.Fatalf("core.NewBlockChain failed: %v", err)
- }
- pool := legacypool.New(testTxPoolConfig, chain)
- txpool, _ := txpool.New(new(big.Int).SetUint64(testTxPoolConfig.PriceLimit), chain, []txpool.SubPool{pool})
-
- return &testWorkerBackend{
- db: db,
- chain: chain,
- txPool: txpool,
- genesis: gspec,
- }
-}
-
-func (b *testWorkerBackend) BlockChain() *core.BlockChain { return b.chain }
-func (b *testWorkerBackend) TxPool() *txpool.TxPool { return b.txPool }
-
-func (b *testWorkerBackend) newRandomTx(creation bool) *types.Transaction {
- var tx *types.Transaction
- gasPrice := big.NewInt(10 * params.InitialBaseFee)
- if creation {
- tx, _ = types.SignTx(types.NewContractCreation(b.txPool.Nonce(testBankAddress), big.NewInt(0), testGas, gasPrice, common.FromHex(testCode)), types.HomesteadSigner{}, testBankKey)
- } else {
- tx, _ = types.SignTx(types.NewTransaction(b.txPool.Nonce(testBankAddress), testUserAddress, big.NewInt(1000), params.TxGas, gasPrice, nil), types.HomesteadSigner{}, testBankKey)
- }
- return tx
-}
-
-func newTestWorker(t *testing.T, chainConfig *params.ChainConfig, engine consensus.Engine, db ethdb.Database, blocks int) (*worker, *testWorkerBackend) {
- backend := newTestWorkerBackend(t, chainConfig, engine, db, blocks)
- backend.txPool.Add(pendingTxs, true, false)
- w := newWorker(testConfig, chainConfig, engine, backend, new(event.TypeMux), nil, false)
- w.setEtherbase(testBankAddress)
- return w, backend
-}
-
-func TestGenerateAndImportBlock(t *testing.T) {
- t.Parallel()
- var (
- db = rawdb.NewMemoryDatabase()
- config = *params.AllCliqueProtocolChanges
- )
- config.Clique = ¶ms.CliqueConfig{Period: 1, Epoch: 30000}
- engine := clique.New(config.Clique, db)
-
- w, b := newTestWorker(t, &config, engine, db, 0)
- defer w.close()
-
- // This test chain imports the mined blocks.
- chain, _ := core.NewBlockChain(rawdb.NewMemoryDatabase(), nil, b.genesis, nil, engine, vm.Config{}, nil, nil)
- defer chain.Stop()
-
- // Ignore empty commit here for less noise.
- w.skipSealHook = func(task *task) bool {
- return len(task.receipts) == 0
- }
-
- // Wait for mined blocks.
- sub := w.mux.Subscribe(core.NewMinedBlockEvent{})
- defer sub.Unsubscribe()
-
- // Start mining!
- w.start()
-
- for i := 0; i < 5; i++ {
- b.txPool.Add([]*types.Transaction{b.newRandomTx(true)}, true, false)
- b.txPool.Add([]*types.Transaction{b.newRandomTx(false)}, true, false)
-
- select {
- case ev := <-sub.Chan():
- block := ev.Data.(core.NewMinedBlockEvent).Block
- if _, err := chain.InsertChain([]*types.Block{block}); err != nil {
- t.Fatalf("failed to insert new mined block %d: %v", block.NumberU64(), err)
- }
- case <-time.After(3 * time.Second): // Worker needs 1s to include new changes.
- t.Fatalf("timeout")
- }
- }
-}
-
-func TestEmptyWorkEthash(t *testing.T) {
- t.Parallel()
- testEmptyWork(t, ethashChainConfig, ethash.NewFaker())
-}
-func TestEmptyWorkClique(t *testing.T) {
- t.Parallel()
- testEmptyWork(t, cliqueChainConfig, clique.New(cliqueChainConfig.Clique, rawdb.NewMemoryDatabase()))
-}
-
-func testEmptyWork(t *testing.T, chainConfig *params.ChainConfig, engine consensus.Engine) {
- defer engine.Close()
-
- w, _ := newTestWorker(t, chainConfig, engine, rawdb.NewMemoryDatabase(), 0)
- defer w.close()
-
- taskCh := make(chan struct{}, 2)
- checkEqual := func(t *testing.T, task *task) {
- // The work should contain 1 tx
- receiptLen, balance := 1, big.NewInt(1000)
- if len(task.receipts) != receiptLen {
- t.Fatalf("receipt number mismatch: have %d, want %d", len(task.receipts), receiptLen)
- }
- if task.state.GetBalance(testUserAddress).Cmp(balance) != 0 {
- t.Fatalf("account balance mismatch: have %d, want %d", task.state.GetBalance(testUserAddress), balance)
- }
- }
- w.newTaskHook = func(task *task) {
- if task.block.NumberU64() == 1 {
- checkEqual(t, task)
- taskCh <- struct{}{}
- }
- }
- w.skipSealHook = func(task *task) bool { return true }
- w.fullTaskHook = func() {
- time.Sleep(100 * time.Millisecond)
- }
- w.start() // Start mining!
- select {
- case <-taskCh:
- case <-time.NewTimer(3 * time.Second).C:
- t.Error("new task timeout")
- }
-}
-
-func TestAdjustIntervalEthash(t *testing.T) {
- t.Parallel()
- testAdjustInterval(t, ethashChainConfig, ethash.NewFaker())
-}
-
-func TestAdjustIntervalClique(t *testing.T) {
- t.Parallel()
- testAdjustInterval(t, cliqueChainConfig, clique.New(cliqueChainConfig.Clique, rawdb.NewMemoryDatabase()))
-}
-
-func testAdjustInterval(t *testing.T, chainConfig *params.ChainConfig, engine consensus.Engine) {
- defer engine.Close()
-
- w, _ := newTestWorker(t, chainConfig, engine, rawdb.NewMemoryDatabase(), 0)
- defer w.close()
-
- w.skipSealHook = func(task *task) bool {
- return true
- }
- w.fullTaskHook = func() {
- time.Sleep(100 * time.Millisecond)
- }
- var (
- progress = make(chan struct{}, 10)
- result = make([]float64, 0, 10)
- index = 0
- start atomic.Bool
- )
- w.resubmitHook = func(minInterval time.Duration, recommitInterval time.Duration) {
- // Short circuit if interval checking hasn't started.
- if !start.Load() {
- return
- }
- var wantMinInterval, wantRecommitInterval time.Duration
-
- switch index {
- case 0:
- wantMinInterval, wantRecommitInterval = 3*time.Second, 3*time.Second
- case 1:
- origin := float64(3 * time.Second.Nanoseconds())
- estimate := origin*(1-intervalAdjustRatio) + intervalAdjustRatio*(origin/0.8+intervalAdjustBias)
- wantMinInterval, wantRecommitInterval = 3*time.Second, time.Duration(estimate)*time.Nanosecond
- case 2:
- estimate := result[index-1]
- min := float64(3 * time.Second.Nanoseconds())
- estimate = estimate*(1-intervalAdjustRatio) + intervalAdjustRatio*(min-intervalAdjustBias)
- wantMinInterval, wantRecommitInterval = 3*time.Second, time.Duration(estimate)*time.Nanosecond
- case 3:
- wantMinInterval, wantRecommitInterval = time.Second, time.Second
- }
-
- // Check interval
- if minInterval != wantMinInterval {
- t.Errorf("resubmit min interval mismatch: have %v, want %v ", minInterval, wantMinInterval)
- }
- if recommitInterval != wantRecommitInterval {
- t.Errorf("resubmit interval mismatch: have %v, want %v", recommitInterval, wantRecommitInterval)
- }
- result = append(result, float64(recommitInterval.Nanoseconds()))
- index += 1
- progress <- struct{}{}
- }
- w.start()
-
- time.Sleep(time.Second) // Ensure two tasks have been submitted due to start opt
- start.Store(true)
-
- w.setRecommitInterval(3 * time.Second)
- select {
- case <-progress:
- case <-time.NewTimer(time.Second).C:
- t.Error("interval reset timeout")
- }
-
- w.resubmitAdjustCh <- &intervalAdjust{inc: true, ratio: 0.8}
- select {
- case <-progress:
- case <-time.NewTimer(time.Second).C:
- t.Error("interval reset timeout")
- }
-
- w.resubmitAdjustCh <- &intervalAdjust{inc: false}
- select {
- case <-progress:
- case <-time.NewTimer(time.Second).C:
- t.Error("interval reset timeout")
- }
-
- w.setRecommitInterval(500 * time.Millisecond)
- select {
- case <-progress:
- case <-time.NewTimer(time.Second).C:
- t.Error("interval reset timeout")
- }
-}
-
-func TestGetSealingWorkEthash(t *testing.T) {
- t.Parallel()
- testGetSealingWork(t, ethashChainConfig, ethash.NewFaker())
-}
-
-func TestGetSealingWorkClique(t *testing.T) {
- t.Parallel()
- testGetSealingWork(t, cliqueChainConfig, clique.New(cliqueChainConfig.Clique, rawdb.NewMemoryDatabase()))
-}
-
-func TestGetSealingWorkPostMerge(t *testing.T) {
- t.Parallel()
- local := new(params.ChainConfig)
- *local = *ethashChainConfig
- local.TerminalTotalDifficulty = big.NewInt(0)
- testGetSealingWork(t, local, ethash.NewFaker())
-}
-
-func testGetSealingWork(t *testing.T, chainConfig *params.ChainConfig, engine consensus.Engine) {
- defer engine.Close()
-
- w, b := newTestWorker(t, chainConfig, engine, rawdb.NewMemoryDatabase(), 0)
- defer w.close()
-
- w.setExtra([]byte{0x01, 0x02})
-
- w.skipSealHook = func(task *task) bool {
- return true
- }
- w.fullTaskHook = func() {
- time.Sleep(100 * time.Millisecond)
- }
- timestamp := uint64(time.Now().Unix())
- assertBlock := func(block *types.Block, number uint64, coinbase common.Address, random common.Hash) {
- if block.Time() != timestamp {
- // Sometime the timestamp will be mutated if the timestamp
- // is even smaller than parent block's. It's OK.
- t.Logf("Invalid timestamp, want %d, get %d", timestamp, block.Time())
- }
- _, isClique := engine.(*clique.Clique)
- if !isClique {
- if len(block.Extra()) != 2 {
- t.Error("Unexpected extra field")
- }
- if block.Coinbase() != coinbase {
- t.Errorf("Unexpected coinbase got %x want %x", block.Coinbase(), coinbase)
- }
- } else {
- if block.Coinbase() != (common.Address{}) {
- t.Error("Unexpected coinbase")
- }
- }
- if !isClique {
- if block.MixDigest() != random {
- t.Error("Unexpected mix digest")
- }
- }
- if block.Nonce() != 0 {
- t.Error("Unexpected block nonce")
- }
- if block.NumberU64() != number {
- t.Errorf("Mismatched block number, want %d got %d", number, block.NumberU64())
- }
- }
- var cases = []struct {
- parent common.Hash
- coinbase common.Address
- random common.Hash
- expectNumber uint64
- expectErr bool
- }{
- {
- b.chain.Genesis().Hash(),
- common.HexToAddress("0xdeadbeef"),
- common.HexToHash("0xcafebabe"),
- uint64(1),
- false,
- },
- {
- b.chain.CurrentBlock().Hash(),
- common.HexToAddress("0xdeadbeef"),
- common.HexToHash("0xcafebabe"),
- b.chain.CurrentBlock().Number.Uint64() + 1,
- false,
- },
- {
- b.chain.CurrentBlock().Hash(),
- common.Address{},
- common.HexToHash("0xcafebabe"),
- b.chain.CurrentBlock().Number.Uint64() + 1,
- false,
- },
- {
- b.chain.CurrentBlock().Hash(),
- common.Address{},
- common.Hash{},
- b.chain.CurrentBlock().Number.Uint64() + 1,
- false,
- },
- {
- common.HexToHash("0xdeadbeef"),
- common.HexToAddress("0xdeadbeef"),
- common.HexToHash("0xcafebabe"),
- 0,
- true,
- },
- }
-
- // This API should work even when the automatic sealing is not enabled
- for _, c := range cases {
- r := w.getSealingBlock(&generateParams{
- parentHash: c.parent,
- timestamp: timestamp,
- coinbase: c.coinbase,
- random: c.random,
- withdrawals: nil,
- beaconRoot: nil,
- noTxs: false,
- forceTime: true,
- })
- if c.expectErr {
- if r.err == nil {
- t.Error("Expect error but get nil")
- }
- } else {
- if r.err != nil {
- t.Errorf("Unexpected error %v", r.err)
- }
- assertBlock(r.block, c.expectNumber, c.coinbase, c.random)
- }
- }
-
- // This API should work even when the automatic sealing is enabled
- w.start()
- for _, c := range cases {
- r := w.getSealingBlock(&generateParams{
- parentHash: c.parent,
- timestamp: timestamp,
- coinbase: c.coinbase,
- random: c.random,
- withdrawals: nil,
- beaconRoot: nil,
- noTxs: false,
- forceTime: true,
- })
- if c.expectErr {
- if r.err == nil {
- t.Error("Expect error but get nil")
- }
- } else {
- if r.err != nil {
- t.Errorf("Unexpected error %v", r.err)
- }
- assertBlock(r.block, c.expectNumber, c.coinbase, c.random)
- }
- }
-}