mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-26 06:36:43 +00:00
Shivam/txpool tracing (#604)
* lock, unlock to rlock, runlock * add : tracing Pending() and Locals() * Log time spent in committing a tx during mining * Remove data from logging * Move log into case where a tx completes without error * profile fillTransactions * fix conflict * bug fixes * add logs * txpool: add tracing in Pending() * rearrange tracing * add attributes * fix * fix * log error in profiling * update file mode and file path for profiling * full profiling * fix log * fix log * less wait * fix * fix * logs * worker: use block number for prof files * initial * txList add * fix gas calculation * fix * green tests * linters * prettify * allocate less * no locks between pending and reorg * no locks * no locks on locals * more tests * linters * less allocs * comment * optimize errors * linters * fix * fix * Linters * linters * linters * simplify errors * atomics for transactions * fix * optimize workers * fix copy * linters * txpool tracing * linters * fix tracing * duration in mcs * locks * metrics * fix * cache hit/miss * less locks on evict * remove once * remove once * switch off pprof * fix data race * fix data race * add : sealed total/empty blocks metric gauge * add : RPC debug_getTraceStack * fix : RPC debug_getTraceStack * fix : RPC debug_getTraceStack for all go-routines * linters * add data race test on txpool * fix concurrency * noleak * increase batch size * prettify * tests * baseFee mutex * panic fix * linters * fix gas fee data race * linters * more transactions * debug * debug * fix ticker * fix test * add cacheMu * more tests * fix test panic * linters * add statistics * add statistics * txitems data race * fix tx list Has * fix : lint Co-authored-by: Arpit Temani <temaniarpit27@gmail.com> Co-authored-by: Jerry <jerrycgh@gmail.com> Co-authored-by: Manav Darji <manavdarji.india@gmail.com> Co-authored-by: Evgeny Danienko <6655321@bk.ru>
This commit is contained in:
parent
50a778207e
commit
b1d86bd6ea
32 changed files with 3352 additions and 451 deletions
7
Makefile
7
Makefile
|
|
@ -59,7 +59,10 @@ ios:
|
||||||
@echo "Import \"$(GOBIN)/Geth.framework\" to use the library."
|
@echo "Import \"$(GOBIN)/Geth.framework\" to use the library."
|
||||||
|
|
||||||
test:
|
test:
|
||||||
$(GOTEST) --timeout 5m -shuffle=on -cover -coverprofile=cover.out $(TESTALL)
|
$(GOTEST) --timeout 5m -shuffle=on -cover -short -coverprofile=cover.out -covermode=atomic $(TESTALL)
|
||||||
|
|
||||||
|
test-txpool-race:
|
||||||
|
$(GOTEST) -run=TestPoolMiningDataRaces --timeout 600m -race -v ./core/
|
||||||
|
|
||||||
test-race:
|
test-race:
|
||||||
$(GOTEST) --timeout 15m -race -shuffle=on $(TESTALL)
|
$(GOTEST) --timeout 15m -race -shuffle=on $(TESTALL)
|
||||||
|
|
@ -75,7 +78,7 @@ lint:
|
||||||
|
|
||||||
lintci-deps:
|
lintci-deps:
|
||||||
rm -f ./build/bin/golangci-lint
|
rm -f ./build/bin/golangci-lint
|
||||||
curl -sSfL https://raw.githubusercontent.com/golangci/golangci-lint/master/install.sh | sh -s -- -b ./build/bin v1.48.0
|
curl -sSfL https://raw.githubusercontent.com/golangci/golangci-lint/master/install.sh | sh -s -- -b ./build/bin v1.50.1
|
||||||
|
|
||||||
goimports:
|
goimports:
|
||||||
goimports -local "$(PACKAGE)" -w .
|
goimports -local "$(PACKAGE)" -w .
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,8 @@ import (
|
||||||
"os"
|
"os"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
|
"gopkg.in/urfave/cli.v1"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/common/hexutil"
|
"github.com/ethereum/go-ethereum/common/hexutil"
|
||||||
"github.com/ethereum/go-ethereum/core"
|
"github.com/ethereum/go-ethereum/core"
|
||||||
|
|
@ -32,7 +34,6 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/params"
|
"github.com/ethereum/go-ethereum/params"
|
||||||
"github.com/ethereum/go-ethereum/rlp"
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
"github.com/ethereum/go-ethereum/tests"
|
"github.com/ethereum/go-ethereum/tests"
|
||||||
"gopkg.in/urfave/cli.v1"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type result struct {
|
type result struct {
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
package debug
|
package debug
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
"runtime"
|
"runtime"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -26,3 +27,26 @@ func Callers(show int) []string {
|
||||||
|
|
||||||
return callers
|
return callers
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func CodeLine() (string, string, int) {
|
||||||
|
pc, filename, line, _ := runtime.Caller(1)
|
||||||
|
return runtime.FuncForPC(pc).Name(), filename, line
|
||||||
|
}
|
||||||
|
|
||||||
|
func CodeLineStr() string {
|
||||||
|
pc, filename, line, _ := runtime.Caller(1)
|
||||||
|
return fmt.Sprintf("%s:%d - %s", filename, line, runtime.FuncForPC(pc).Name())
|
||||||
|
}
|
||||||
|
|
||||||
|
func Stack(all bool) []byte {
|
||||||
|
buf := make([]byte, 4096)
|
||||||
|
|
||||||
|
for {
|
||||||
|
n := runtime.Stack(buf, all)
|
||||||
|
if n < len(buf) {
|
||||||
|
return buf[:n]
|
||||||
|
}
|
||||||
|
|
||||||
|
buf = make([]byte, 2*len(buf))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,8 @@ package math
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/big"
|
"math/big"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Various big integer limit values.
|
// Various big integer limit values.
|
||||||
|
|
@ -132,6 +134,7 @@ func MustParseBig256(s string) *big.Int {
|
||||||
// BigPow returns a ** b as a big integer.
|
// BigPow returns a ** b as a big integer.
|
||||||
func BigPow(a, b int64) *big.Int {
|
func BigPow(a, b int64) *big.Int {
|
||||||
r := big.NewInt(a)
|
r := big.NewInt(a)
|
||||||
|
|
||||||
return r.Exp(r, big.NewInt(b), nil)
|
return r.Exp(r, big.NewInt(b), nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -140,6 +143,15 @@ func BigMax(x, y *big.Int) *big.Int {
|
||||||
if x.Cmp(y) < 0 {
|
if x.Cmp(y) < 0 {
|
||||||
return y
|
return y
|
||||||
}
|
}
|
||||||
|
|
||||||
|
return x
|
||||||
|
}
|
||||||
|
|
||||||
|
func BigMaxUint(x, y *uint256.Int) *uint256.Int {
|
||||||
|
if x.Lt(y) {
|
||||||
|
return y
|
||||||
|
}
|
||||||
|
|
||||||
return x
|
return x
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -148,6 +160,15 @@ func BigMin(x, y *big.Int) *big.Int {
|
||||||
if x.Cmp(y) > 0 {
|
if x.Cmp(y) > 0 {
|
||||||
return y
|
return y
|
||||||
}
|
}
|
||||||
|
|
||||||
|
return x
|
||||||
|
}
|
||||||
|
|
||||||
|
func BigMinUint256(x, y *uint256.Int) *uint256.Int {
|
||||||
|
if x.Gt(y) {
|
||||||
|
return y
|
||||||
|
}
|
||||||
|
|
||||||
return x
|
return x
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -227,10 +248,10 @@ func U256Bytes(n *big.Int) []byte {
|
||||||
// S256 interprets x as a two's complement number.
|
// S256 interprets x as a two's complement number.
|
||||||
// x must not exceed 256 bits (the result is undefined if it does) and is not modified.
|
// x must not exceed 256 bits (the result is undefined if it does) and is not modified.
|
||||||
//
|
//
|
||||||
// S256(0) = 0
|
// S256(0) = 0
|
||||||
// S256(1) = 1
|
// S256(1) = 1
|
||||||
// S256(2**255) = -2**255
|
// S256(2**255) = -2**255
|
||||||
// S256(2**256-1) = -1
|
// S256(2**256-1) = -1
|
||||||
func S256(x *big.Int) *big.Int {
|
func S256(x *big.Int) *big.Int {
|
||||||
if x.Cmp(tt255) < 0 {
|
if x.Cmp(tt255) < 0 {
|
||||||
return x
|
return x
|
||||||
|
|
|
||||||
23
common/math/uint.go
Normal file
23
common/math/uint.go
Normal file
|
|
@ -0,0 +1,23 @@
|
||||||
|
package math
|
||||||
|
|
||||||
|
import (
|
||||||
|
"math/big"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
U0 = uint256.NewInt(0)
|
||||||
|
U1 = uint256.NewInt(1)
|
||||||
|
U100 = uint256.NewInt(100)
|
||||||
|
)
|
||||||
|
|
||||||
|
func U256LTE(a, b *uint256.Int) bool {
|
||||||
|
return a.Lt(b) || a.Eq(b)
|
||||||
|
}
|
||||||
|
|
||||||
|
func FromBig(v *big.Int) *uint256.Int {
|
||||||
|
u, _ := uint256.FromBig(v)
|
||||||
|
|
||||||
|
return u
|
||||||
|
}
|
||||||
9
common/time.go
Normal file
9
common/time.go
Normal file
|
|
@ -0,0 +1,9 @@
|
||||||
|
package common
|
||||||
|
|
||||||
|
import "time"
|
||||||
|
|
||||||
|
const TimeMilliseconds = "15:04:05.000"
|
||||||
|
|
||||||
|
func NowMilliseconds() string {
|
||||||
|
return time.Now().Format(TimeMilliseconds)
|
||||||
|
}
|
||||||
|
|
@ -4,6 +4,7 @@ import (
|
||||||
"context"
|
"context"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"go.opentelemetry.io/otel"
|
||||||
"go.opentelemetry.io/otel/attribute"
|
"go.opentelemetry.io/otel/attribute"
|
||||||
"go.opentelemetry.io/otel/trace"
|
"go.opentelemetry.io/otel/trace"
|
||||||
)
|
)
|
||||||
|
|
@ -51,11 +52,16 @@ func Trace(ctx context.Context, spanName string) (context.Context, trace.Span) {
|
||||||
return tr.Start(ctx, spanName)
|
return tr.Start(ctx, spanName)
|
||||||
}
|
}
|
||||||
|
|
||||||
func Exec(ctx context.Context, spanName string, opts ...Option) {
|
func Exec(ctx context.Context, instrumentationName, spanName string, opts ...Option) {
|
||||||
var span trace.Span
|
var span trace.Span
|
||||||
|
|
||||||
tr := FromContext(ctx)
|
tr := FromContext(ctx)
|
||||||
|
|
||||||
|
if tr == nil && len(instrumentationName) != 0 {
|
||||||
|
tr = otel.GetTracerProvider().Tracer(instrumentationName)
|
||||||
|
ctx = WithTracer(ctx, tr)
|
||||||
|
}
|
||||||
|
|
||||||
if tr != nil {
|
if tr != nil {
|
||||||
ctx, span = tr.Start(ctx, spanName)
|
ctx, span = tr.Start(ctx, spanName)
|
||||||
}
|
}
|
||||||
|
|
@ -85,7 +91,7 @@ func ElapsedTime(ctx context.Context, span trace.Span, msg string, fn func(conte
|
||||||
fn(ctx, span)
|
fn(ctx, span)
|
||||||
|
|
||||||
if span != nil {
|
if span != nil {
|
||||||
span.SetAttributes(attribute.Int(msg, int(time.Since(now).Milliseconds())))
|
span.SetAttributes(attribute.Int(msg, int(time.Since(now).Microseconds())))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -821,7 +821,7 @@ func (c *Bor) FinalizeAndAssemble(ctx context.Context, chain consensus.ChainHead
|
||||||
if IsSprintStart(headerNumber, c.config.CalculateSprint(headerNumber)) {
|
if IsSprintStart(headerNumber, c.config.CalculateSprint(headerNumber)) {
|
||||||
cx := statefull.ChainContext{Chain: chain, Bor: c}
|
cx := statefull.ChainContext{Chain: chain, Bor: c}
|
||||||
|
|
||||||
tracing.Exec(finalizeCtx, "bor.checkAndCommitSpan", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(finalizeCtx, "", "bor.checkAndCommitSpan", func(ctx context.Context, span trace.Span) {
|
||||||
// check and commit span
|
// check and commit span
|
||||||
err = c.checkAndCommitSpan(finalizeCtx, state, header, cx)
|
err = c.checkAndCommitSpan(finalizeCtx, state, header, cx)
|
||||||
})
|
})
|
||||||
|
|
@ -832,7 +832,7 @@ func (c *Bor) FinalizeAndAssemble(ctx context.Context, chain consensus.ChainHead
|
||||||
}
|
}
|
||||||
|
|
||||||
if c.HeimdallClient != nil {
|
if c.HeimdallClient != nil {
|
||||||
tracing.Exec(finalizeCtx, "bor.checkAndCommitSpan", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(finalizeCtx, "", "bor.checkAndCommitSpan", func(ctx context.Context, span trace.Span) {
|
||||||
// commit states
|
// commit states
|
||||||
stateSyncData, err = c.CommitStates(finalizeCtx, state, header, cx)
|
stateSyncData, err = c.CommitStates(finalizeCtx, state, header, cx)
|
||||||
})
|
})
|
||||||
|
|
@ -844,7 +844,7 @@ func (c *Bor) FinalizeAndAssemble(ctx context.Context, chain consensus.ChainHead
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
tracing.Exec(finalizeCtx, "bor.changeContractCodeIfNeeded", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(finalizeCtx, "", "bor.changeContractCodeIfNeeded", func(ctx context.Context, span trace.Span) {
|
||||||
err = c.changeContractCodeIfNeeded(headerNumber, state)
|
err = c.changeContractCodeIfNeeded(headerNumber, state)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -854,7 +854,7 @@ func (c *Bor) FinalizeAndAssemble(ctx context.Context, chain consensus.ChainHead
|
||||||
}
|
}
|
||||||
|
|
||||||
// No block rewards in PoA, so the state remains as it is
|
// No block rewards in PoA, so the state remains as it is
|
||||||
tracing.Exec(finalizeCtx, "bor.IntermediateRoot", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(finalizeCtx, "", "bor.IntermediateRoot", func(ctx context.Context, span trace.Span) {
|
||||||
header.Root = state.IntermediateRoot(chain.Config().IsEIP158(header.Number))
|
header.Root = state.IntermediateRoot(chain.Config().IsEIP158(header.Number))
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,8 @@ import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/big"
|
"math/big"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/common/math"
|
"github.com/ethereum/go-ethereum/common/math"
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
|
|
@ -92,3 +94,54 @@ func CalcBaseFee(config *params.ChainConfig, parent *types.Header) *big.Int {
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// CalcBaseFee calculates the basefee of the header.
|
||||||
|
func CalcBaseFeeUint(config *params.ChainConfig, parent *types.Header) *uint256.Int {
|
||||||
|
var (
|
||||||
|
initialBaseFeeUint = uint256.NewInt(params.InitialBaseFee)
|
||||||
|
baseFeeChangeDenominatorUint64 = params.BaseFeeChangeDenominator(config.Bor, parent.Number)
|
||||||
|
baseFeeChangeDenominatorUint = uint256.NewInt(baseFeeChangeDenominatorUint64)
|
||||||
|
)
|
||||||
|
|
||||||
|
// If the current block is the first EIP-1559 block, return the InitialBaseFee.
|
||||||
|
if !config.IsLondon(parent.Number) {
|
||||||
|
return initialBaseFeeUint.Clone()
|
||||||
|
}
|
||||||
|
|
||||||
|
var (
|
||||||
|
parentGasTarget = parent.GasLimit / params.ElasticityMultiplier
|
||||||
|
parentGasTargetBig = uint256.NewInt(parentGasTarget)
|
||||||
|
)
|
||||||
|
|
||||||
|
// If the parent gasUsed is the same as the target, the baseFee remains unchanged.
|
||||||
|
if parent.GasUsed == parentGasTarget {
|
||||||
|
return math.FromBig(parent.BaseFee)
|
||||||
|
}
|
||||||
|
|
||||||
|
if parent.GasUsed > parentGasTarget {
|
||||||
|
// If the parent block used more gas than its target, the baseFee should increase.
|
||||||
|
gasUsedDelta := uint256.NewInt(parent.GasUsed - parentGasTarget)
|
||||||
|
|
||||||
|
parentBaseFee := math.FromBig(parent.BaseFee)
|
||||||
|
x := gasUsedDelta.Mul(parentBaseFee, gasUsedDelta)
|
||||||
|
y := x.Div(x, parentGasTargetBig)
|
||||||
|
baseFeeDelta := math.BigMaxUint(
|
||||||
|
x.Div(y, baseFeeChangeDenominatorUint),
|
||||||
|
math.U1,
|
||||||
|
)
|
||||||
|
|
||||||
|
return x.Add(parentBaseFee, baseFeeDelta)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Otherwise if the parent block used less gas than its target, the baseFee should decrease.
|
||||||
|
gasUsedDelta := uint256.NewInt(parentGasTarget - parent.GasUsed)
|
||||||
|
parentBaseFee := math.FromBig(parent.BaseFee)
|
||||||
|
x := gasUsedDelta.Mul(parentBaseFee, gasUsedDelta)
|
||||||
|
y := x.Div(x, parentGasTargetBig)
|
||||||
|
baseFeeDelta := x.Div(y, baseFeeChangeDenominatorUint)
|
||||||
|
|
||||||
|
return math.BigMaxUint(
|
||||||
|
x.Sub(parentBaseFee, baseFeeDelta),
|
||||||
|
math.U0.Clone(),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -61,11 +61,13 @@ func (journal *txJournal) load(add func([]*types.Transaction) []error) error {
|
||||||
if _, err := os.Stat(journal.path); os.IsNotExist(err) {
|
if _, err := os.Stat(journal.path); os.IsNotExist(err) {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Open the journal for loading any past transactions
|
// Open the journal for loading any past transactions
|
||||||
input, err := os.Open(journal.path)
|
input, err := os.Open(journal.path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
defer input.Close()
|
defer input.Close()
|
||||||
|
|
||||||
// Temporarily discard any journal additions (don't double add on load)
|
// Temporarily discard any journal additions (don't double add on load)
|
||||||
|
|
@ -80,29 +82,35 @@ func (journal *txJournal) load(add func([]*types.Transaction) []error) error {
|
||||||
// appropriate progress counters. Then use this method to load all the
|
// appropriate progress counters. Then use this method to load all the
|
||||||
// journaled transactions in small-ish batches.
|
// journaled transactions in small-ish batches.
|
||||||
loadBatch := func(txs types.Transactions) {
|
loadBatch := func(txs types.Transactions) {
|
||||||
|
errs := add(txs)
|
||||||
|
|
||||||
|
dropped = len(errs)
|
||||||
|
|
||||||
for _, err := range add(txs) {
|
for _, err := range add(txs) {
|
||||||
if err != nil {
|
log.Debug("Failed to add journaled transaction", "err", err)
|
||||||
log.Debug("Failed to add journaled transaction", "err", err)
|
|
||||||
dropped++
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
var (
|
var (
|
||||||
failure error
|
failure error
|
||||||
batch types.Transactions
|
batch types.Transactions
|
||||||
)
|
)
|
||||||
|
|
||||||
for {
|
for {
|
||||||
// Parse the next transaction and terminate on error
|
// Parse the next transaction and terminate on error
|
||||||
tx := new(types.Transaction)
|
tx := new(types.Transaction)
|
||||||
|
|
||||||
if err = stream.Decode(tx); err != nil {
|
if err = stream.Decode(tx); err != nil {
|
||||||
if err != io.EOF {
|
if err != io.EOF {
|
||||||
failure = err
|
failure = err
|
||||||
}
|
}
|
||||||
|
|
||||||
if batch.Len() > 0 {
|
if batch.Len() > 0 {
|
||||||
loadBatch(batch)
|
loadBatch(batch)
|
||||||
}
|
}
|
||||||
|
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|
||||||
// New transaction parsed, queue up for later, import if threshold is reached
|
// New transaction parsed, queue up for later, import if threshold is reached
|
||||||
total++
|
total++
|
||||||
|
|
||||||
|
|
@ -111,6 +119,7 @@ func (journal *txJournal) load(add func([]*types.Transaction) []error) error {
|
||||||
batch = batch[:0]
|
batch = batch[:0]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Info("Loaded local transaction journal", "transactions", total, "dropped", dropped)
|
log.Info("Loaded local transaction journal", "transactions", total, "dropped", dropped)
|
||||||
|
|
||||||
return failure
|
return failure
|
||||||
|
|
|
||||||
236
core/tx_list.go
236
core/tx_list.go
|
|
@ -19,13 +19,15 @@ package core
|
||||||
import (
|
import (
|
||||||
"container/heap"
|
"container/heap"
|
||||||
"math"
|
"math"
|
||||||
"math/big"
|
|
||||||
"sort"
|
"sort"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
cmath "github.com/ethereum/go-ethereum/common/math"
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -54,36 +56,67 @@ func (h *nonceHeap) Pop() interface{} {
|
||||||
type txSortedMap struct {
|
type txSortedMap struct {
|
||||||
items map[uint64]*types.Transaction // Hash map storing the transaction data
|
items map[uint64]*types.Transaction // Hash map storing the transaction data
|
||||||
index *nonceHeap // Heap of nonces of all the stored transactions (non-strict mode)
|
index *nonceHeap // Heap of nonces of all the stored transactions (non-strict mode)
|
||||||
cache types.Transactions // Cache of the transactions already sorted
|
m sync.RWMutex
|
||||||
|
|
||||||
|
cache types.Transactions // Cache of the transactions already sorted
|
||||||
|
isEmpty bool
|
||||||
|
cacheMu sync.RWMutex
|
||||||
}
|
}
|
||||||
|
|
||||||
// newTxSortedMap creates a new nonce-sorted transaction map.
|
// newTxSortedMap creates a new nonce-sorted transaction map.
|
||||||
func newTxSortedMap() *txSortedMap {
|
func newTxSortedMap() *txSortedMap {
|
||||||
return &txSortedMap{
|
return &txSortedMap{
|
||||||
items: make(map[uint64]*types.Transaction),
|
items: make(map[uint64]*types.Transaction),
|
||||||
index: new(nonceHeap),
|
index: new(nonceHeap),
|
||||||
|
isEmpty: true,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get retrieves the current transactions associated with the given nonce.
|
// Get retrieves the current transactions associated with the given nonce.
|
||||||
func (m *txSortedMap) Get(nonce uint64) *types.Transaction {
|
func (m *txSortedMap) Get(nonce uint64) *types.Transaction {
|
||||||
|
m.m.RLock()
|
||||||
|
defer m.m.RUnlock()
|
||||||
|
|
||||||
return m.items[nonce]
|
return m.items[nonce]
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (m *txSortedMap) Has(nonce uint64) bool {
|
||||||
|
if m == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
m.m.RLock()
|
||||||
|
defer m.m.RUnlock()
|
||||||
|
|
||||||
|
return m.items[nonce] != nil
|
||||||
|
}
|
||||||
|
|
||||||
// Put inserts a new transaction into the map, also updating the map's nonce
|
// Put inserts a new transaction into the map, also updating the map's nonce
|
||||||
// index. If a transaction already exists with the same nonce, it's overwritten.
|
// index. If a transaction already exists with the same nonce, it's overwritten.
|
||||||
func (m *txSortedMap) Put(tx *types.Transaction) {
|
func (m *txSortedMap) Put(tx *types.Transaction) {
|
||||||
|
m.m.Lock()
|
||||||
|
defer m.m.Unlock()
|
||||||
|
|
||||||
nonce := tx.Nonce()
|
nonce := tx.Nonce()
|
||||||
if m.items[nonce] == nil {
|
if m.items[nonce] == nil {
|
||||||
heap.Push(m.index, nonce)
|
heap.Push(m.index, nonce)
|
||||||
}
|
}
|
||||||
m.items[nonce], m.cache = tx, nil
|
|
||||||
|
m.items[nonce] = tx
|
||||||
|
|
||||||
|
m.cacheMu.Lock()
|
||||||
|
m.isEmpty = true
|
||||||
|
m.cache = nil
|
||||||
|
m.cacheMu.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Forward removes all transactions from the map with a nonce lower than the
|
// Forward removes all transactions from the map with a nonce lower than the
|
||||||
// provided threshold. Every removed transaction is returned for any post-removal
|
// provided threshold. Every removed transaction is returned for any post-removal
|
||||||
// maintenance.
|
// maintenance.
|
||||||
func (m *txSortedMap) Forward(threshold uint64) types.Transactions {
|
func (m *txSortedMap) Forward(threshold uint64) types.Transactions {
|
||||||
|
m.m.Lock()
|
||||||
|
defer m.m.Unlock()
|
||||||
|
|
||||||
var removed types.Transactions
|
var removed types.Transactions
|
||||||
|
|
||||||
// Pop off heap items until the threshold is reached
|
// Pop off heap items until the threshold is reached
|
||||||
|
|
@ -92,10 +125,15 @@ func (m *txSortedMap) Forward(threshold uint64) types.Transactions {
|
||||||
removed = append(removed, m.items[nonce])
|
removed = append(removed, m.items[nonce])
|
||||||
delete(m.items, nonce)
|
delete(m.items, nonce)
|
||||||
}
|
}
|
||||||
|
|
||||||
// If we had a cached order, shift the front
|
// If we had a cached order, shift the front
|
||||||
|
m.cacheMu.Lock()
|
||||||
if m.cache != nil {
|
if m.cache != nil {
|
||||||
|
hitCacheCounter.Inc(1)
|
||||||
m.cache = m.cache[len(removed):]
|
m.cache = m.cache[len(removed):]
|
||||||
}
|
}
|
||||||
|
m.cacheMu.Unlock()
|
||||||
|
|
||||||
return removed
|
return removed
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -105,6 +143,9 @@ func (m *txSortedMap) Forward(threshold uint64) types.Transactions {
|
||||||
// If you want to do several consecutive filterings, it's therefore better to first
|
// If you want to do several consecutive filterings, it's therefore better to first
|
||||||
// do a .filter(func1) followed by .Filter(func2) or reheap()
|
// do a .filter(func1) followed by .Filter(func2) or reheap()
|
||||||
func (m *txSortedMap) Filter(filter func(*types.Transaction) bool) types.Transactions {
|
func (m *txSortedMap) Filter(filter func(*types.Transaction) bool) types.Transactions {
|
||||||
|
m.m.Lock()
|
||||||
|
defer m.m.Unlock()
|
||||||
|
|
||||||
removed := m.filter(filter)
|
removed := m.filter(filter)
|
||||||
// If transactions were removed, the heap and cache are ruined
|
// If transactions were removed, the heap and cache are ruined
|
||||||
if len(removed) > 0 {
|
if len(removed) > 0 {
|
||||||
|
|
@ -115,11 +156,19 @@ func (m *txSortedMap) Filter(filter func(*types.Transaction) bool) types.Transac
|
||||||
|
|
||||||
func (m *txSortedMap) reheap() {
|
func (m *txSortedMap) reheap() {
|
||||||
*m.index = make([]uint64, 0, len(m.items))
|
*m.index = make([]uint64, 0, len(m.items))
|
||||||
|
|
||||||
for nonce := range m.items {
|
for nonce := range m.items {
|
||||||
*m.index = append(*m.index, nonce)
|
*m.index = append(*m.index, nonce)
|
||||||
}
|
}
|
||||||
|
|
||||||
heap.Init(m.index)
|
heap.Init(m.index)
|
||||||
|
|
||||||
|
m.cacheMu.Lock()
|
||||||
m.cache = nil
|
m.cache = nil
|
||||||
|
m.isEmpty = true
|
||||||
|
m.cacheMu.Unlock()
|
||||||
|
|
||||||
|
resetCacheGauge.Inc(1)
|
||||||
}
|
}
|
||||||
|
|
||||||
// filter is identical to Filter, but **does not** regenerate the heap. This method
|
// filter is identical to Filter, but **does not** regenerate the heap. This method
|
||||||
|
|
@ -135,7 +184,12 @@ func (m *txSortedMap) filter(filter func(*types.Transaction) bool) types.Transac
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if len(removed) > 0 {
|
if len(removed) > 0 {
|
||||||
|
m.cacheMu.Lock()
|
||||||
m.cache = nil
|
m.cache = nil
|
||||||
|
m.isEmpty = true
|
||||||
|
m.cacheMu.Unlock()
|
||||||
|
|
||||||
|
resetCacheGauge.Inc(1)
|
||||||
}
|
}
|
||||||
return removed
|
return removed
|
||||||
}
|
}
|
||||||
|
|
@ -143,45 +197,66 @@ func (m *txSortedMap) filter(filter func(*types.Transaction) bool) types.Transac
|
||||||
// Cap places a hard limit on the number of items, returning all transactions
|
// Cap places a hard limit on the number of items, returning all transactions
|
||||||
// exceeding that limit.
|
// exceeding that limit.
|
||||||
func (m *txSortedMap) Cap(threshold int) types.Transactions {
|
func (m *txSortedMap) Cap(threshold int) types.Transactions {
|
||||||
|
m.m.Lock()
|
||||||
|
defer m.m.Unlock()
|
||||||
|
|
||||||
// Short circuit if the number of items is under the limit
|
// Short circuit if the number of items is under the limit
|
||||||
if len(m.items) <= threshold {
|
if len(m.items) <= threshold {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Otherwise gather and drop the highest nonce'd transactions
|
// Otherwise gather and drop the highest nonce'd transactions
|
||||||
var drops types.Transactions
|
var drops types.Transactions
|
||||||
|
|
||||||
sort.Sort(*m.index)
|
sort.Sort(*m.index)
|
||||||
|
|
||||||
for size := len(m.items); size > threshold; size-- {
|
for size := len(m.items); size > threshold; size-- {
|
||||||
drops = append(drops, m.items[(*m.index)[size-1]])
|
drops = append(drops, m.items[(*m.index)[size-1]])
|
||||||
delete(m.items, (*m.index)[size-1])
|
delete(m.items, (*m.index)[size-1])
|
||||||
}
|
}
|
||||||
|
|
||||||
*m.index = (*m.index)[:threshold]
|
*m.index = (*m.index)[:threshold]
|
||||||
heap.Init(m.index)
|
heap.Init(m.index)
|
||||||
|
|
||||||
// If we had a cache, shift the back
|
// If we had a cache, shift the back
|
||||||
|
m.cacheMu.Lock()
|
||||||
if m.cache != nil {
|
if m.cache != nil {
|
||||||
m.cache = m.cache[:len(m.cache)-len(drops)]
|
m.cache = m.cache[:len(m.cache)-len(drops)]
|
||||||
}
|
}
|
||||||
|
m.cacheMu.Unlock()
|
||||||
|
|
||||||
return drops
|
return drops
|
||||||
}
|
}
|
||||||
|
|
||||||
// Remove deletes a transaction from the maintained map, returning whether the
|
// Remove deletes a transaction from the maintained map, returning whether the
|
||||||
// transaction was found.
|
// transaction was found.
|
||||||
func (m *txSortedMap) Remove(nonce uint64) bool {
|
func (m *txSortedMap) Remove(nonce uint64) bool {
|
||||||
|
m.m.Lock()
|
||||||
|
defer m.m.Unlock()
|
||||||
|
|
||||||
// Short circuit if no transaction is present
|
// Short circuit if no transaction is present
|
||||||
_, ok := m.items[nonce]
|
_, ok := m.items[nonce]
|
||||||
if !ok {
|
if !ok {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
// Otherwise delete the transaction and fix the heap index
|
// Otherwise delete the transaction and fix the heap index
|
||||||
for i := 0; i < m.index.Len(); i++ {
|
for i := 0; i < m.index.Len(); i++ {
|
||||||
if (*m.index)[i] == nonce {
|
if (*m.index)[i] == nonce {
|
||||||
heap.Remove(m.index, i)
|
heap.Remove(m.index, i)
|
||||||
|
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
delete(m.items, nonce)
|
delete(m.items, nonce)
|
||||||
|
|
||||||
|
m.cacheMu.Lock()
|
||||||
m.cache = nil
|
m.cache = nil
|
||||||
|
m.isEmpty = true
|
||||||
|
m.cacheMu.Unlock()
|
||||||
|
|
||||||
|
resetCacheGauge.Inc(1)
|
||||||
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
@ -194,55 +269,125 @@ func (m *txSortedMap) Remove(nonce uint64) bool {
|
||||||
// prevent getting into and invalid state. This is not something that should ever
|
// prevent getting into and invalid state. This is not something that should ever
|
||||||
// happen but better to be self correcting than failing!
|
// happen but better to be self correcting than failing!
|
||||||
func (m *txSortedMap) Ready(start uint64) types.Transactions {
|
func (m *txSortedMap) Ready(start uint64) types.Transactions {
|
||||||
|
m.m.Lock()
|
||||||
|
defer m.m.Unlock()
|
||||||
|
|
||||||
// Short circuit if no transactions are available
|
// Short circuit if no transactions are available
|
||||||
if m.index.Len() == 0 || (*m.index)[0] > start {
|
if m.index.Len() == 0 || (*m.index)[0] > start {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Otherwise start accumulating incremental transactions
|
// Otherwise start accumulating incremental transactions
|
||||||
var ready types.Transactions
|
var ready types.Transactions
|
||||||
|
|
||||||
for next := (*m.index)[0]; m.index.Len() > 0 && (*m.index)[0] == next; next++ {
|
for next := (*m.index)[0]; m.index.Len() > 0 && (*m.index)[0] == next; next++ {
|
||||||
ready = append(ready, m.items[next])
|
ready = append(ready, m.items[next])
|
||||||
delete(m.items, next)
|
delete(m.items, next)
|
||||||
heap.Pop(m.index)
|
heap.Pop(m.index)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
m.cacheMu.Lock()
|
||||||
m.cache = nil
|
m.cache = nil
|
||||||
|
m.isEmpty = true
|
||||||
|
m.cacheMu.Unlock()
|
||||||
|
|
||||||
|
resetCacheGauge.Inc(1)
|
||||||
|
|
||||||
return ready
|
return ready
|
||||||
}
|
}
|
||||||
|
|
||||||
// Len returns the length of the transaction map.
|
// Len returns the length of the transaction map.
|
||||||
func (m *txSortedMap) Len() int {
|
func (m *txSortedMap) Len() int {
|
||||||
|
m.m.RLock()
|
||||||
|
defer m.m.RUnlock()
|
||||||
|
|
||||||
return len(m.items)
|
return len(m.items)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *txSortedMap) flatten() types.Transactions {
|
func (m *txSortedMap) flatten() types.Transactions {
|
||||||
// If the sorting was not cached yet, create and cache it
|
// If the sorting was not cached yet, create and cache it
|
||||||
if m.cache == nil {
|
m.cacheMu.Lock()
|
||||||
m.cache = make(types.Transactions, 0, len(m.items))
|
defer m.cacheMu.Unlock()
|
||||||
|
|
||||||
|
if m.isEmpty {
|
||||||
|
m.isEmpty = false // to simulate sync.Once
|
||||||
|
|
||||||
|
m.cacheMu.Unlock()
|
||||||
|
|
||||||
|
m.m.RLock()
|
||||||
|
|
||||||
|
cache := make(types.Transactions, 0, len(m.items))
|
||||||
|
|
||||||
for _, tx := range m.items {
|
for _, tx := range m.items {
|
||||||
m.cache = append(m.cache, tx)
|
cache = append(cache, tx)
|
||||||
}
|
}
|
||||||
sort.Sort(types.TxByNonce(m.cache))
|
|
||||||
|
m.m.RUnlock()
|
||||||
|
|
||||||
|
// exclude sorting from locks
|
||||||
|
sort.Sort(types.TxByNonce(cache))
|
||||||
|
|
||||||
|
m.cacheMu.Lock()
|
||||||
|
m.cache = cache
|
||||||
|
|
||||||
|
reinitCacheGauge.Inc(1)
|
||||||
|
missCacheCounter.Inc(1)
|
||||||
|
} else {
|
||||||
|
hitCacheCounter.Inc(1)
|
||||||
}
|
}
|
||||||
|
|
||||||
return m.cache
|
return m.cache
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (m *txSortedMap) lastElement() *types.Transaction {
|
||||||
|
// If the sorting was not cached yet, create and cache it
|
||||||
|
m.cacheMu.Lock()
|
||||||
|
defer m.cacheMu.Unlock()
|
||||||
|
|
||||||
|
cache := m.cache
|
||||||
|
|
||||||
|
if m.isEmpty {
|
||||||
|
m.isEmpty = false // to simulate sync.Once
|
||||||
|
|
||||||
|
m.cacheMu.Unlock()
|
||||||
|
|
||||||
|
cache = make(types.Transactions, 0, len(m.items))
|
||||||
|
|
||||||
|
m.m.RLock()
|
||||||
|
|
||||||
|
for _, tx := range m.items {
|
||||||
|
cache = append(cache, tx)
|
||||||
|
}
|
||||||
|
|
||||||
|
m.m.RUnlock()
|
||||||
|
|
||||||
|
// exclude sorting from locks
|
||||||
|
sort.Sort(types.TxByNonce(cache))
|
||||||
|
|
||||||
|
m.cacheMu.Lock()
|
||||||
|
m.cache = cache
|
||||||
|
|
||||||
|
reinitCacheGauge.Inc(1)
|
||||||
|
missCacheCounter.Inc(1)
|
||||||
|
} else {
|
||||||
|
hitCacheCounter.Inc(1)
|
||||||
|
}
|
||||||
|
|
||||||
|
return cache[len(cache)-1]
|
||||||
|
}
|
||||||
|
|
||||||
// Flatten creates a nonce-sorted slice of transactions based on the loosely
|
// Flatten creates a nonce-sorted slice of transactions based on the loosely
|
||||||
// sorted internal representation. The result of the sorting is cached in case
|
// sorted internal representation. The result of the sorting is cached in case
|
||||||
// it's requested again before any modifications are made to the contents.
|
// it's requested again before any modifications are made to the contents.
|
||||||
func (m *txSortedMap) Flatten() types.Transactions {
|
func (m *txSortedMap) Flatten() types.Transactions {
|
||||||
// Copy the cache to prevent accidental modifications
|
// Copy the cache to prevent accidental modifications
|
||||||
cache := m.flatten()
|
return m.flatten()
|
||||||
txs := make(types.Transactions, len(cache))
|
|
||||||
copy(txs, cache)
|
|
||||||
return txs
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// LastElement returns the last element of a flattened list, thus, the
|
// LastElement returns the last element of a flattened list, thus, the
|
||||||
// transaction with the highest nonce
|
// transaction with the highest nonce
|
||||||
func (m *txSortedMap) LastElement() *types.Transaction {
|
func (m *txSortedMap) LastElement() *types.Transaction {
|
||||||
cache := m.flatten()
|
return m.lastElement()
|
||||||
return cache[len(cache)-1]
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// txList is a "list" of transactions belonging to an account, sorted by account
|
// txList is a "list" of transactions belonging to an account, sorted by account
|
||||||
|
|
@ -253,17 +398,16 @@ type txList struct {
|
||||||
strict bool // Whether nonces are strictly continuous or not
|
strict bool // Whether nonces are strictly continuous or not
|
||||||
txs *txSortedMap // Heap indexed sorted hash map of the transactions
|
txs *txSortedMap // Heap indexed sorted hash map of the transactions
|
||||||
|
|
||||||
costcap *big.Int // Price of the highest costing transaction (reset only if exceeds balance)
|
costcap *uint256.Int // Price of the highest costing transaction (reset only if exceeds balance)
|
||||||
gascap uint64 // Gas limit of the highest spending transaction (reset only if exceeds block limit)
|
gascap uint64 // Gas limit of the highest spending transaction (reset only if exceeds block limit)
|
||||||
}
|
}
|
||||||
|
|
||||||
// newTxList create a new transaction list for maintaining nonce-indexable fast,
|
// newTxList create a new transaction list for maintaining nonce-indexable fast,
|
||||||
// gapped, sortable transaction lists.
|
// gapped, sortable transaction lists.
|
||||||
func newTxList(strict bool) *txList {
|
func newTxList(strict bool) *txList {
|
||||||
return &txList{
|
return &txList{
|
||||||
strict: strict,
|
strict: strict,
|
||||||
txs: newTxSortedMap(),
|
txs: newTxSortedMap(),
|
||||||
costcap: new(big.Int),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -285,31 +429,36 @@ func (l *txList) Add(tx *types.Transaction, priceBump uint64) (bool, *types.Tran
|
||||||
if old.GasFeeCapCmp(tx) >= 0 || old.GasTipCapCmp(tx) >= 0 {
|
if old.GasFeeCapCmp(tx) >= 0 || old.GasTipCapCmp(tx) >= 0 {
|
||||||
return false, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// thresholdFeeCap = oldFC * (100 + priceBump) / 100
|
// thresholdFeeCap = oldFC * (100 + priceBump) / 100
|
||||||
a := big.NewInt(100 + int64(priceBump))
|
a := uint256.NewInt(100 + priceBump)
|
||||||
aFeeCap := new(big.Int).Mul(a, old.GasFeeCap())
|
aFeeCap := uint256.NewInt(0).Mul(a, old.GasFeeCapUint())
|
||||||
aTip := a.Mul(a, old.GasTipCap())
|
aTip := a.Mul(a, old.GasTipCapUint())
|
||||||
|
|
||||||
// thresholdTip = oldTip * (100 + priceBump) / 100
|
// thresholdTip = oldTip * (100 + priceBump) / 100
|
||||||
b := big.NewInt(100)
|
b := cmath.U100
|
||||||
thresholdFeeCap := aFeeCap.Div(aFeeCap, b)
|
thresholdFeeCap := aFeeCap.Div(aFeeCap, b)
|
||||||
thresholdTip := aTip.Div(aTip, b)
|
thresholdTip := aTip.Div(aTip, b)
|
||||||
|
|
||||||
// We have to ensure that both the new fee cap and tip are higher than the
|
// We have to ensure that both the new fee cap and tip are higher than the
|
||||||
// old ones as well as checking the percentage threshold to ensure that
|
// old ones as well as checking the percentage threshold to ensure that
|
||||||
// this is accurate for low (Wei-level) gas price replacements.
|
// this is accurate for low (Wei-level) gas price replacements.
|
||||||
if tx.GasFeeCapIntCmp(thresholdFeeCap) < 0 || tx.GasTipCapIntCmp(thresholdTip) < 0 {
|
if tx.GasFeeCapUIntLt(thresholdFeeCap) || tx.GasTipCapUIntLt(thresholdTip) {
|
||||||
return false, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Otherwise overwrite the old transaction with the current one
|
// Otherwise overwrite the old transaction with the current one
|
||||||
l.txs.Put(tx)
|
l.txs.Put(tx)
|
||||||
if cost := tx.Cost(); l.costcap.Cmp(cost) < 0 {
|
|
||||||
|
if cost := tx.CostUint(); l.costcap == nil || l.costcap.Lt(cost) {
|
||||||
l.costcap = cost
|
l.costcap = cost
|
||||||
}
|
}
|
||||||
|
|
||||||
if gas := tx.Gas(); l.gascap < gas {
|
if gas := tx.Gas(); l.gascap < gas {
|
||||||
l.gascap = gas
|
l.gascap = gas
|
||||||
}
|
}
|
||||||
|
|
||||||
return true, old
|
return true, old
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -329,17 +478,20 @@ func (l *txList) Forward(threshold uint64) types.Transactions {
|
||||||
// a point in calculating all the costs or if the balance covers all. If the threshold
|
// a point in calculating all the costs or if the balance covers all. If the threshold
|
||||||
// is lower than the costgas cap, the caps will be reset to a new high after removing
|
// is lower than the costgas cap, the caps will be reset to a new high after removing
|
||||||
// the newly invalidated transactions.
|
// the newly invalidated transactions.
|
||||||
func (l *txList) Filter(costLimit *big.Int, gasLimit uint64) (types.Transactions, types.Transactions) {
|
func (l *txList) Filter(costLimit *uint256.Int, gasLimit uint64) (types.Transactions, types.Transactions) {
|
||||||
// If all transactions are below the threshold, short circuit
|
// If all transactions are below the threshold, short circuit
|
||||||
if l.costcap.Cmp(costLimit) <= 0 && l.gascap <= gasLimit {
|
if cmath.U256LTE(l.costcap, costLimit) && l.gascap <= gasLimit {
|
||||||
return nil, nil
|
return nil, nil
|
||||||
}
|
}
|
||||||
l.costcap = new(big.Int).Set(costLimit) // Lower the caps to the thresholds
|
|
||||||
|
l.costcap = costLimit.Clone() // Lower the caps to the thresholds
|
||||||
l.gascap = gasLimit
|
l.gascap = gasLimit
|
||||||
|
|
||||||
// Filter out all the transactions above the account's funds
|
// Filter out all the transactions above the account's funds
|
||||||
|
cost := uint256.NewInt(0)
|
||||||
removed := l.txs.Filter(func(tx *types.Transaction) bool {
|
removed := l.txs.Filter(func(tx *types.Transaction) bool {
|
||||||
return tx.Gas() > gasLimit || tx.Cost().Cmp(costLimit) > 0
|
cost.SetFromBig(tx.Cost())
|
||||||
|
return tx.Gas() > gasLimit || cost.Gt(costLimit)
|
||||||
})
|
})
|
||||||
|
|
||||||
if len(removed) == 0 {
|
if len(removed) == 0 {
|
||||||
|
|
@ -416,13 +568,18 @@ func (l *txList) LastElement() *types.Transaction {
|
||||||
return l.txs.LastElement()
|
return l.txs.LastElement()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (l *txList) Has(nonce uint64) bool {
|
||||||
|
return l != nil && l.txs.items[nonce] != nil
|
||||||
|
}
|
||||||
|
|
||||||
// priceHeap is a heap.Interface implementation over transactions for retrieving
|
// priceHeap is a heap.Interface implementation over transactions for retrieving
|
||||||
// price-sorted transactions to discard when the pool fills up. If baseFee is set
|
// price-sorted transactions to discard when the pool fills up. If baseFee is set
|
||||||
// then the heap is sorted based on the effective tip based on the given base fee.
|
// then the heap is sorted based on the effective tip based on the given base fee.
|
||||||
// If baseFee is nil then the sorting is based on gasFeeCap.
|
// If baseFee is nil then the sorting is based on gasFeeCap.
|
||||||
type priceHeap struct {
|
type priceHeap struct {
|
||||||
baseFee *big.Int // heap should always be re-sorted after baseFee is changed
|
baseFee *uint256.Int // heap should always be re-sorted after baseFee is changed
|
||||||
list []*types.Transaction
|
list []*types.Transaction
|
||||||
|
baseFeeMu sync.RWMutex
|
||||||
}
|
}
|
||||||
|
|
||||||
func (h *priceHeap) Len() int { return len(h.list) }
|
func (h *priceHeap) Len() int { return len(h.list) }
|
||||||
|
|
@ -440,16 +597,24 @@ func (h *priceHeap) Less(i, j int) bool {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (h *priceHeap) cmp(a, b *types.Transaction) int {
|
func (h *priceHeap) cmp(a, b *types.Transaction) int {
|
||||||
|
h.baseFeeMu.RLock()
|
||||||
|
|
||||||
if h.baseFee != nil {
|
if h.baseFee != nil {
|
||||||
// Compare effective tips if baseFee is specified
|
// Compare effective tips if baseFee is specified
|
||||||
if c := a.EffectiveGasTipCmp(b, h.baseFee); c != 0 {
|
if c := a.EffectiveGasTipTxUintCmp(b, h.baseFee); c != 0 {
|
||||||
|
h.baseFeeMu.RUnlock()
|
||||||
|
|
||||||
return c
|
return c
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
h.baseFeeMu.RUnlock()
|
||||||
|
|
||||||
// Compare fee caps if baseFee is not specified or effective tips are equal
|
// Compare fee caps if baseFee is not specified or effective tips are equal
|
||||||
if c := a.GasFeeCapCmp(b); c != 0 {
|
if c := a.GasFeeCapCmp(b); c != 0 {
|
||||||
return c
|
return c
|
||||||
}
|
}
|
||||||
|
|
||||||
// Compare tips if effective tips and fee caps are equal
|
// Compare tips if effective tips and fee caps are equal
|
||||||
return a.GasTipCapCmp(b)
|
return a.GasTipCapCmp(b)
|
||||||
}
|
}
|
||||||
|
|
@ -629,7 +794,10 @@ func (l *txPricedList) Reheap() {
|
||||||
|
|
||||||
// SetBaseFee updates the base fee and triggers a re-heap. Note that Removed is not
|
// SetBaseFee updates the base fee and triggers a re-heap. Note that Removed is not
|
||||||
// necessary to call right before SetBaseFee when processing a new block.
|
// necessary to call right before SetBaseFee when processing a new block.
|
||||||
func (l *txPricedList) SetBaseFee(baseFee *big.Int) {
|
func (l *txPricedList) SetBaseFee(baseFee *uint256.Int) {
|
||||||
|
l.urgent.baseFeeMu.Lock()
|
||||||
l.urgent.baseFee = baseFee
|
l.urgent.baseFee = baseFee
|
||||||
|
l.urgent.baseFeeMu.Unlock()
|
||||||
|
|
||||||
l.Reheap()
|
l.Reheap()
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -17,10 +17,11 @@
|
||||||
package core
|
package core
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"math/big"
|
|
||||||
"math/rand"
|
"math/rand"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
)
|
)
|
||||||
|
|
@ -59,11 +60,15 @@ func BenchmarkTxListAdd(b *testing.B) {
|
||||||
for i := 0; i < len(txs); i++ {
|
for i := 0; i < len(txs); i++ {
|
||||||
txs[i] = transaction(uint64(i), 0, key)
|
txs[i] = transaction(uint64(i), 0, key)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Insert the transactions in a random order
|
// Insert the transactions in a random order
|
||||||
priceLimit := big.NewInt(int64(DefaultTxPoolConfig.PriceLimit))
|
priceLimit := uint256.NewInt(DefaultTxPoolConfig.PriceLimit)
|
||||||
b.ResetTimer()
|
b.ResetTimer()
|
||||||
|
b.ReportAllocs()
|
||||||
|
|
||||||
for i := 0; i < b.N; i++ {
|
for i := 0; i < b.N; i++ {
|
||||||
list := newTxList(true)
|
list := newTxList(true)
|
||||||
|
|
||||||
for _, v := range rand.Perm(len(txs)) {
|
for _, v := range rand.Perm(len(txs)) {
|
||||||
list.Add(txs[v], DefaultTxPoolConfig.PriceBump)
|
list.Add(txs[v], DefaultTxPoolConfig.PriceBump)
|
||||||
list.Filter(priceLimit, DefaultTxPoolConfig.PriceBump)
|
list.Filter(priceLimit, DefaultTxPoolConfig.PriceBump)
|
||||||
|
|
|
||||||
1068
core/tx_pool.go
1068
core/tx_pool.go
File diff suppressed because it is too large
Load diff
1599
core/tx_pool_test.go
1599
core/tx_pool_test.go
File diff suppressed because it is too large
Load diff
|
|
@ -19,6 +19,8 @@ package types
|
||||||
import (
|
import (
|
||||||
"math/big"
|
"math/big"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -44,15 +46,16 @@ func (al AccessList) StorageKeys() int {
|
||||||
|
|
||||||
// AccessListTx is the data of EIP-2930 access list transactions.
|
// AccessListTx is the data of EIP-2930 access list transactions.
|
||||||
type AccessListTx struct {
|
type AccessListTx struct {
|
||||||
ChainID *big.Int // destination chain ID
|
ChainID *big.Int // destination chain ID
|
||||||
Nonce uint64 // nonce of sender account
|
Nonce uint64 // nonce of sender account
|
||||||
GasPrice *big.Int // wei per gas
|
GasPrice *big.Int // wei per gas
|
||||||
Gas uint64 // gas limit
|
gasPriceUint256 *uint256.Int // wei per gas
|
||||||
To *common.Address `rlp:"nil"` // nil means contract creation
|
Gas uint64 // gas limit
|
||||||
Value *big.Int // wei amount
|
To *common.Address `rlp:"nil"` // nil means contract creation
|
||||||
Data []byte // contract invocation input data
|
Value *big.Int // wei amount
|
||||||
AccessList AccessList // EIP-2930 access list
|
Data []byte // contract invocation input data
|
||||||
V, R, S *big.Int // signature values
|
AccessList AccessList // EIP-2930 access list
|
||||||
|
V, R, S *big.Int // signature values
|
||||||
}
|
}
|
||||||
|
|
||||||
// copy creates a deep copy of the transaction data and initializes all fields.
|
// copy creates a deep copy of the transaction data and initializes all fields.
|
||||||
|
|
@ -80,6 +83,12 @@ func (tx *AccessListTx) copy() TxData {
|
||||||
}
|
}
|
||||||
if tx.GasPrice != nil {
|
if tx.GasPrice != nil {
|
||||||
cpy.GasPrice.Set(tx.GasPrice)
|
cpy.GasPrice.Set(tx.GasPrice)
|
||||||
|
|
||||||
|
if cpy.gasPriceUint256 != nil {
|
||||||
|
cpy.gasPriceUint256.Set(tx.gasPriceUint256)
|
||||||
|
} else {
|
||||||
|
cpy.gasPriceUint256, _ = uint256.FromBig(tx.GasPrice)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if tx.V != nil {
|
if tx.V != nil {
|
||||||
cpy.V.Set(tx.V)
|
cpy.V.Set(tx.V)
|
||||||
|
|
@ -100,11 +109,39 @@ func (tx *AccessListTx) accessList() AccessList { return tx.AccessList }
|
||||||
func (tx *AccessListTx) data() []byte { return tx.Data }
|
func (tx *AccessListTx) data() []byte { return tx.Data }
|
||||||
func (tx *AccessListTx) gas() uint64 { return tx.Gas }
|
func (tx *AccessListTx) gas() uint64 { return tx.Gas }
|
||||||
func (tx *AccessListTx) gasPrice() *big.Int { return tx.GasPrice }
|
func (tx *AccessListTx) gasPrice() *big.Int { return tx.GasPrice }
|
||||||
func (tx *AccessListTx) gasTipCap() *big.Int { return tx.GasPrice }
|
func (tx *AccessListTx) gasPriceU256() *uint256.Int {
|
||||||
func (tx *AccessListTx) gasFeeCap() *big.Int { return tx.GasPrice }
|
if tx.gasPriceUint256 != nil {
|
||||||
func (tx *AccessListTx) value() *big.Int { return tx.Value }
|
return tx.gasPriceUint256
|
||||||
func (tx *AccessListTx) nonce() uint64 { return tx.Nonce }
|
}
|
||||||
func (tx *AccessListTx) to() *common.Address { return tx.To }
|
|
||||||
|
tx.gasPriceUint256, _ = uint256.FromBig(tx.GasPrice)
|
||||||
|
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
|
||||||
|
func (tx *AccessListTx) gasTipCap() *big.Int { return tx.GasPrice }
|
||||||
|
func (tx *AccessListTx) gasTipCapU256() *uint256.Int {
|
||||||
|
if tx.gasPriceUint256 != nil {
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
|
||||||
|
tx.gasPriceUint256, _ = uint256.FromBig(tx.GasPrice)
|
||||||
|
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
func (tx *AccessListTx) gasFeeCap() *big.Int { return tx.GasPrice }
|
||||||
|
func (tx *AccessListTx) gasFeeCapU256() *uint256.Int {
|
||||||
|
if tx.gasPriceUint256 != nil {
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
|
||||||
|
tx.gasPriceUint256, _ = uint256.FromBig(tx.GasPrice)
|
||||||
|
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
func (tx *AccessListTx) value() *big.Int { return tx.Value }
|
||||||
|
func (tx *AccessListTx) nonce() uint64 { return tx.Nonce }
|
||||||
|
func (tx *AccessListTx) to() *common.Address { return tx.To }
|
||||||
|
|
||||||
func (tx *AccessListTx) rawSignatureValues() (v, r, s *big.Int) {
|
func (tx *AccessListTx) rawSignatureValues() (v, r, s *big.Int) {
|
||||||
return tx.V, tx.R, tx.S
|
return tx.V, tx.R, tx.S
|
||||||
|
|
|
||||||
|
|
@ -19,19 +19,23 @@ package types
|
||||||
import (
|
import (
|
||||||
"math/big"
|
"math/big"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
)
|
)
|
||||||
|
|
||||||
type DynamicFeeTx struct {
|
type DynamicFeeTx struct {
|
||||||
ChainID *big.Int
|
ChainID *big.Int
|
||||||
Nonce uint64
|
Nonce uint64
|
||||||
GasTipCap *big.Int // a.k.a. maxPriorityFeePerGas
|
GasTipCap *big.Int // a.k.a. maxPriorityFeePerGas
|
||||||
GasFeeCap *big.Int // a.k.a. maxFeePerGas
|
gasTipCapUint256 *uint256.Int // a.k.a. maxPriorityFeePerGas
|
||||||
Gas uint64
|
GasFeeCap *big.Int // a.k.a. maxFeePerGas
|
||||||
To *common.Address `rlp:"nil"` // nil means contract creation
|
gasFeeCapUint256 *uint256.Int // a.k.a. maxFeePerGas
|
||||||
Value *big.Int
|
Gas uint64
|
||||||
Data []byte
|
To *common.Address `rlp:"nil"` // nil means contract creation
|
||||||
AccessList AccessList
|
Value *big.Int
|
||||||
|
Data []byte
|
||||||
|
AccessList AccessList
|
||||||
|
|
||||||
// Signature values
|
// Signature values
|
||||||
V *big.Int `json:"v" gencodec:"required"`
|
V *big.Int `json:"v" gencodec:"required"`
|
||||||
|
|
@ -65,9 +69,21 @@ func (tx *DynamicFeeTx) copy() TxData {
|
||||||
}
|
}
|
||||||
if tx.GasTipCap != nil {
|
if tx.GasTipCap != nil {
|
||||||
cpy.GasTipCap.Set(tx.GasTipCap)
|
cpy.GasTipCap.Set(tx.GasTipCap)
|
||||||
|
|
||||||
|
if cpy.gasTipCapUint256 != nil {
|
||||||
|
cpy.gasTipCapUint256.Set(tx.gasTipCapUint256)
|
||||||
|
} else {
|
||||||
|
cpy.gasTipCapUint256, _ = uint256.FromBig(tx.GasTipCap)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if tx.GasFeeCap != nil {
|
if tx.GasFeeCap != nil {
|
||||||
cpy.GasFeeCap.Set(tx.GasFeeCap)
|
cpy.GasFeeCap.Set(tx.GasFeeCap)
|
||||||
|
|
||||||
|
if cpy.gasFeeCapUint256 != nil {
|
||||||
|
cpy.gasFeeCapUint256.Set(tx.gasFeeCapUint256)
|
||||||
|
} else {
|
||||||
|
cpy.gasFeeCapUint256, _ = uint256.FromBig(tx.GasFeeCap)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if tx.V != nil {
|
if tx.V != nil {
|
||||||
cpy.V.Set(tx.V)
|
cpy.V.Set(tx.V)
|
||||||
|
|
@ -88,11 +104,38 @@ func (tx *DynamicFeeTx) accessList() AccessList { return tx.AccessList }
|
||||||
func (tx *DynamicFeeTx) data() []byte { return tx.Data }
|
func (tx *DynamicFeeTx) data() []byte { return tx.Data }
|
||||||
func (tx *DynamicFeeTx) gas() uint64 { return tx.Gas }
|
func (tx *DynamicFeeTx) gas() uint64 { return tx.Gas }
|
||||||
func (tx *DynamicFeeTx) gasFeeCap() *big.Int { return tx.GasFeeCap }
|
func (tx *DynamicFeeTx) gasFeeCap() *big.Int { return tx.GasFeeCap }
|
||||||
func (tx *DynamicFeeTx) gasTipCap() *big.Int { return tx.GasTipCap }
|
func (tx *DynamicFeeTx) gasFeeCapU256() *uint256.Int {
|
||||||
func (tx *DynamicFeeTx) gasPrice() *big.Int { return tx.GasFeeCap }
|
if tx.gasFeeCapUint256 != nil {
|
||||||
func (tx *DynamicFeeTx) value() *big.Int { return tx.Value }
|
return tx.gasFeeCapUint256
|
||||||
func (tx *DynamicFeeTx) nonce() uint64 { return tx.Nonce }
|
}
|
||||||
func (tx *DynamicFeeTx) to() *common.Address { return tx.To }
|
|
||||||
|
tx.gasFeeCapUint256, _ = uint256.FromBig(tx.GasFeeCap)
|
||||||
|
|
||||||
|
return tx.gasFeeCapUint256
|
||||||
|
}
|
||||||
|
func (tx *DynamicFeeTx) gasTipCap() *big.Int { return tx.GasTipCap }
|
||||||
|
func (tx *DynamicFeeTx) gasTipCapU256() *uint256.Int {
|
||||||
|
if tx.gasTipCapUint256 != nil {
|
||||||
|
return tx.gasTipCapUint256
|
||||||
|
}
|
||||||
|
|
||||||
|
tx.gasTipCapUint256, _ = uint256.FromBig(tx.GasTipCap)
|
||||||
|
|
||||||
|
return tx.gasTipCapUint256
|
||||||
|
}
|
||||||
|
func (tx *DynamicFeeTx) gasPrice() *big.Int { return tx.GasFeeCap }
|
||||||
|
func (tx *DynamicFeeTx) gasPriceU256() *uint256.Int {
|
||||||
|
if tx.gasFeeCapUint256 != nil {
|
||||||
|
return tx.gasTipCapUint256
|
||||||
|
}
|
||||||
|
|
||||||
|
tx.gasFeeCapUint256, _ = uint256.FromBig(tx.GasFeeCap)
|
||||||
|
|
||||||
|
return tx.gasFeeCapUint256
|
||||||
|
}
|
||||||
|
func (tx *DynamicFeeTx) value() *big.Int { return tx.Value }
|
||||||
|
func (tx *DynamicFeeTx) nonce() uint64 { return tx.Nonce }
|
||||||
|
func (tx *DynamicFeeTx) to() *common.Address { return tx.To }
|
||||||
|
|
||||||
func (tx *DynamicFeeTx) rawSignatureValues() (v, r, s *big.Int) {
|
func (tx *DynamicFeeTx) rawSignatureValues() (v, r, s *big.Int) {
|
||||||
return tx.V, tx.R, tx.S
|
return tx.V, tx.R, tx.S
|
||||||
|
|
|
||||||
|
|
@ -19,18 +19,21 @@ package types
|
||||||
import (
|
import (
|
||||||
"math/big"
|
"math/big"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
)
|
)
|
||||||
|
|
||||||
// LegacyTx is the transaction data of regular Ethereum transactions.
|
// LegacyTx is the transaction data of regular Ethereum transactions.
|
||||||
type LegacyTx struct {
|
type LegacyTx struct {
|
||||||
Nonce uint64 // nonce of sender account
|
Nonce uint64 // nonce of sender account
|
||||||
GasPrice *big.Int // wei per gas
|
GasPrice *big.Int // wei per gas
|
||||||
Gas uint64 // gas limit
|
gasPriceUint256 *uint256.Int // wei per gas
|
||||||
To *common.Address `rlp:"nil"` // nil means contract creation
|
Gas uint64 // gas limit
|
||||||
Value *big.Int // wei amount
|
To *common.Address `rlp:"nil"` // nil means contract creation
|
||||||
Data []byte // contract invocation input data
|
Value *big.Int // wei amount
|
||||||
V, R, S *big.Int // signature values
|
Data []byte // contract invocation input data
|
||||||
|
V, R, S *big.Int // signature values
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewTransaction creates an unsigned legacy transaction.
|
// NewTransaction creates an unsigned legacy transaction.
|
||||||
|
|
@ -77,6 +80,12 @@ func (tx *LegacyTx) copy() TxData {
|
||||||
}
|
}
|
||||||
if tx.GasPrice != nil {
|
if tx.GasPrice != nil {
|
||||||
cpy.GasPrice.Set(tx.GasPrice)
|
cpy.GasPrice.Set(tx.GasPrice)
|
||||||
|
|
||||||
|
if cpy.gasPriceUint256 != nil {
|
||||||
|
cpy.gasPriceUint256.Set(tx.gasPriceUint256)
|
||||||
|
} else {
|
||||||
|
cpy.gasPriceUint256, _ = uint256.FromBig(tx.GasPrice)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if tx.V != nil {
|
if tx.V != nil {
|
||||||
cpy.V.Set(tx.V)
|
cpy.V.Set(tx.V)
|
||||||
|
|
@ -97,11 +106,38 @@ func (tx *LegacyTx) accessList() AccessList { return nil }
|
||||||
func (tx *LegacyTx) data() []byte { return tx.Data }
|
func (tx *LegacyTx) data() []byte { return tx.Data }
|
||||||
func (tx *LegacyTx) gas() uint64 { return tx.Gas }
|
func (tx *LegacyTx) gas() uint64 { return tx.Gas }
|
||||||
func (tx *LegacyTx) gasPrice() *big.Int { return tx.GasPrice }
|
func (tx *LegacyTx) gasPrice() *big.Int { return tx.GasPrice }
|
||||||
func (tx *LegacyTx) gasTipCap() *big.Int { return tx.GasPrice }
|
func (tx *LegacyTx) gasPriceU256() *uint256.Int {
|
||||||
func (tx *LegacyTx) gasFeeCap() *big.Int { return tx.GasPrice }
|
if tx.gasPriceUint256 != nil {
|
||||||
func (tx *LegacyTx) value() *big.Int { return tx.Value }
|
return tx.gasPriceUint256
|
||||||
func (tx *LegacyTx) nonce() uint64 { return tx.Nonce }
|
}
|
||||||
func (tx *LegacyTx) to() *common.Address { return tx.To }
|
|
||||||
|
tx.gasPriceUint256, _ = uint256.FromBig(tx.GasPrice)
|
||||||
|
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
func (tx *LegacyTx) gasTipCap() *big.Int { return tx.GasPrice }
|
||||||
|
func (tx *LegacyTx) gasTipCapU256() *uint256.Int {
|
||||||
|
if tx.gasPriceUint256 != nil {
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
|
||||||
|
tx.gasPriceUint256, _ = uint256.FromBig(tx.GasPrice)
|
||||||
|
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
func (tx *LegacyTx) gasFeeCap() *big.Int { return tx.GasPrice }
|
||||||
|
func (tx *LegacyTx) gasFeeCapU256() *uint256.Int {
|
||||||
|
if tx.gasPriceUint256 != nil {
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
|
||||||
|
tx.gasPriceUint256, _ = uint256.FromBig(tx.GasPrice)
|
||||||
|
|
||||||
|
return tx.gasPriceUint256
|
||||||
|
}
|
||||||
|
func (tx *LegacyTx) value() *big.Int { return tx.Value }
|
||||||
|
func (tx *LegacyTx) nonce() uint64 { return tx.Nonce }
|
||||||
|
func (tx *LegacyTx) to() *common.Address { return tx.To }
|
||||||
|
|
||||||
func (tx *LegacyTx) rawSignatureValues() (v, r, s *big.Int) {
|
func (tx *LegacyTx) rawSignatureValues() (v, r, s *big.Int) {
|
||||||
return tx.V, tx.R, tx.S
|
return tx.V, tx.R, tx.S
|
||||||
|
|
|
||||||
|
|
@ -25,6 +25,8 @@ import (
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/common/math"
|
"github.com/ethereum/go-ethereum/common/math"
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
|
|
@ -53,9 +55,9 @@ type Transaction struct {
|
||||||
time time.Time // Time first seen locally (spam avoidance)
|
time time.Time // Time first seen locally (spam avoidance)
|
||||||
|
|
||||||
// caches
|
// caches
|
||||||
hash atomic.Value
|
hash atomic.Pointer[common.Hash]
|
||||||
size atomic.Value
|
size atomic.Pointer[common.StorageSize]
|
||||||
from atomic.Value
|
from atomic.Pointer[sigCache]
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewTx creates a new transaction.
|
// NewTx creates a new transaction.
|
||||||
|
|
@ -77,8 +79,11 @@ type TxData interface {
|
||||||
data() []byte
|
data() []byte
|
||||||
gas() uint64
|
gas() uint64
|
||||||
gasPrice() *big.Int
|
gasPrice() *big.Int
|
||||||
|
gasPriceU256() *uint256.Int
|
||||||
gasTipCap() *big.Int
|
gasTipCap() *big.Int
|
||||||
|
gasTipCapU256() *uint256.Int
|
||||||
gasFeeCap() *big.Int
|
gasFeeCap() *big.Int
|
||||||
|
gasFeeCapU256() *uint256.Int
|
||||||
value() *big.Int
|
value() *big.Int
|
||||||
nonce() uint64
|
nonce() uint64
|
||||||
to() *common.Address
|
to() *common.Address
|
||||||
|
|
@ -194,7 +199,8 @@ func (tx *Transaction) setDecoded(inner TxData, size int) {
|
||||||
tx.inner = inner
|
tx.inner = inner
|
||||||
tx.time = time.Now()
|
tx.time = time.Now()
|
||||||
if size > 0 {
|
if size > 0 {
|
||||||
tx.size.Store(common.StorageSize(size))
|
v := float64(size)
|
||||||
|
tx.size.Store((*common.StorageSize)(&v))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -265,16 +271,23 @@ func (tx *Transaction) AccessList() AccessList { return tx.inner.accessList() }
|
||||||
func (tx *Transaction) Gas() uint64 { return tx.inner.gas() }
|
func (tx *Transaction) Gas() uint64 { return tx.inner.gas() }
|
||||||
|
|
||||||
// GasPrice returns the gas price of the transaction.
|
// GasPrice returns the gas price of the transaction.
|
||||||
func (tx *Transaction) GasPrice() *big.Int { return new(big.Int).Set(tx.inner.gasPrice()) }
|
func (tx *Transaction) GasPrice() *big.Int { return new(big.Int).Set(tx.inner.gasPrice()) }
|
||||||
|
func (tx *Transaction) GasPriceRef() *big.Int { return tx.inner.gasPrice() }
|
||||||
|
func (tx *Transaction) GasPriceUint() *uint256.Int { return tx.inner.gasPriceU256() }
|
||||||
|
|
||||||
// GasTipCap returns the gasTipCap per gas of the transaction.
|
// GasTipCap returns the gasTipCap per gas of the transaction.
|
||||||
func (tx *Transaction) GasTipCap() *big.Int { return new(big.Int).Set(tx.inner.gasTipCap()) }
|
func (tx *Transaction) GasTipCap() *big.Int { return new(big.Int).Set(tx.inner.gasTipCap()) }
|
||||||
|
func (tx *Transaction) GasTipCapRef() *big.Int { return tx.inner.gasTipCap() }
|
||||||
|
func (tx *Transaction) GasTipCapUint() *uint256.Int { return tx.inner.gasTipCapU256() }
|
||||||
|
|
||||||
// GasFeeCap returns the fee cap per gas of the transaction.
|
// GasFeeCap returns the fee cap per gas of the transaction.
|
||||||
func (tx *Transaction) GasFeeCap() *big.Int { return new(big.Int).Set(tx.inner.gasFeeCap()) }
|
func (tx *Transaction) GasFeeCap() *big.Int { return new(big.Int).Set(tx.inner.gasFeeCap()) }
|
||||||
|
func (tx *Transaction) GasFeeCapRef() *big.Int { return tx.inner.gasFeeCap() }
|
||||||
|
func (tx *Transaction) GasFeeCapUint() *uint256.Int { return tx.inner.gasFeeCapU256() }
|
||||||
|
|
||||||
// Value returns the ether amount of the transaction.
|
// Value returns the ether amount of the transaction.
|
||||||
func (tx *Transaction) Value() *big.Int { return new(big.Int).Set(tx.inner.value()) }
|
func (tx *Transaction) Value() *big.Int { return new(big.Int).Set(tx.inner.value()) }
|
||||||
|
func (tx *Transaction) ValueRef() *big.Int { return tx.inner.value() }
|
||||||
|
|
||||||
// Nonce returns the sender account nonce of the transaction.
|
// Nonce returns the sender account nonce of the transaction.
|
||||||
func (tx *Transaction) Nonce() uint64 { return tx.inner.nonce() }
|
func (tx *Transaction) Nonce() uint64 { return tx.inner.nonce() }
|
||||||
|
|
@ -287,9 +300,19 @@ func (tx *Transaction) To() *common.Address {
|
||||||
|
|
||||||
// Cost returns gas * gasPrice + value.
|
// Cost returns gas * gasPrice + value.
|
||||||
func (tx *Transaction) Cost() *big.Int {
|
func (tx *Transaction) Cost() *big.Int {
|
||||||
total := new(big.Int).Mul(tx.GasPrice(), new(big.Int).SetUint64(tx.Gas()))
|
gasPrice, _ := uint256.FromBig(tx.GasPriceRef())
|
||||||
total.Add(total, tx.Value())
|
gasPrice.Mul(gasPrice, uint256.NewInt(tx.Gas()))
|
||||||
return total
|
value, _ := uint256.FromBig(tx.ValueRef())
|
||||||
|
|
||||||
|
return gasPrice.Add(gasPrice, value).ToBig()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) CostUint() *uint256.Int {
|
||||||
|
gasPrice, _ := uint256.FromBig(tx.GasPriceRef())
|
||||||
|
gasPrice.Mul(gasPrice, uint256.NewInt(tx.Gas()))
|
||||||
|
value, _ := uint256.FromBig(tx.ValueRef())
|
||||||
|
|
||||||
|
return gasPrice.Add(gasPrice, value)
|
||||||
}
|
}
|
||||||
|
|
||||||
// RawSignatureValues returns the V, R, S signature values of the transaction.
|
// RawSignatureValues returns the V, R, S signature values of the transaction.
|
||||||
|
|
@ -303,11 +326,18 @@ func (tx *Transaction) GasFeeCapCmp(other *Transaction) int {
|
||||||
return tx.inner.gasFeeCap().Cmp(other.inner.gasFeeCap())
|
return tx.inner.gasFeeCap().Cmp(other.inner.gasFeeCap())
|
||||||
}
|
}
|
||||||
|
|
||||||
// GasFeeCapIntCmp compares the fee cap of the transaction against the given fee cap.
|
|
||||||
func (tx *Transaction) GasFeeCapIntCmp(other *big.Int) int {
|
func (tx *Transaction) GasFeeCapIntCmp(other *big.Int) int {
|
||||||
return tx.inner.gasFeeCap().Cmp(other)
|
return tx.inner.gasFeeCap().Cmp(other)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) GasFeeCapUIntCmp(other *uint256.Int) int {
|
||||||
|
return tx.inner.gasFeeCapU256().Cmp(other)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) GasFeeCapUIntLt(other *uint256.Int) bool {
|
||||||
|
return tx.inner.gasFeeCapU256().Lt(other)
|
||||||
|
}
|
||||||
|
|
||||||
// GasTipCapCmp compares the gasTipCap of two transactions.
|
// GasTipCapCmp compares the gasTipCap of two transactions.
|
||||||
func (tx *Transaction) GasTipCapCmp(other *Transaction) int {
|
func (tx *Transaction) GasTipCapCmp(other *Transaction) int {
|
||||||
return tx.inner.gasTipCap().Cmp(other.inner.gasTipCap())
|
return tx.inner.gasTipCap().Cmp(other.inner.gasTipCap())
|
||||||
|
|
@ -318,6 +348,14 @@ func (tx *Transaction) GasTipCapIntCmp(other *big.Int) int {
|
||||||
return tx.inner.gasTipCap().Cmp(other)
|
return tx.inner.gasTipCap().Cmp(other)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) GasTipCapUIntCmp(other *uint256.Int) int {
|
||||||
|
return tx.inner.gasTipCapU256().Cmp(other)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) GasTipCapUIntLt(other *uint256.Int) bool {
|
||||||
|
return tx.inner.gasTipCapU256().Lt(other)
|
||||||
|
}
|
||||||
|
|
||||||
// EffectiveGasTip returns the effective miner gasTipCap for the given base fee.
|
// EffectiveGasTip returns the effective miner gasTipCap for the given base fee.
|
||||||
// Note: if the effective gasTipCap is negative, this method returns both error
|
// Note: if the effective gasTipCap is negative, this method returns both error
|
||||||
// the actual negative value, _and_ ErrGasFeeCapTooLow
|
// the actual negative value, _and_ ErrGasFeeCapTooLow
|
||||||
|
|
@ -356,10 +394,73 @@ func (tx *Transaction) EffectiveGasTipIntCmp(other *big.Int, baseFee *big.Int) i
|
||||||
return tx.EffectiveGasTipValue(baseFee).Cmp(other)
|
return tx.EffectiveGasTipValue(baseFee).Cmp(other)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) EffectiveGasTipUintCmp(other *uint256.Int, baseFee *uint256.Int) int {
|
||||||
|
if baseFee == nil {
|
||||||
|
return tx.GasTipCapUIntCmp(other)
|
||||||
|
}
|
||||||
|
|
||||||
|
return tx.EffectiveGasTipValueUint(baseFee).Cmp(other)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) EffectiveGasTipUintLt(other *uint256.Int, baseFee *uint256.Int) bool {
|
||||||
|
if baseFee == nil {
|
||||||
|
return tx.GasTipCapUIntLt(other)
|
||||||
|
}
|
||||||
|
|
||||||
|
return tx.EffectiveGasTipValueUint(baseFee).Lt(other)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) EffectiveGasTipTxUintCmp(other *Transaction, baseFee *uint256.Int) int {
|
||||||
|
if baseFee == nil {
|
||||||
|
return tx.inner.gasTipCapU256().Cmp(other.inner.gasTipCapU256())
|
||||||
|
}
|
||||||
|
|
||||||
|
return tx.EffectiveGasTipValueUint(baseFee).Cmp(other.EffectiveGasTipValueUint(baseFee))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) EffectiveGasTipValueUint(baseFee *uint256.Int) *uint256.Int {
|
||||||
|
effectiveTip, _ := tx.EffectiveGasTipUnit(baseFee)
|
||||||
|
return effectiveTip
|
||||||
|
}
|
||||||
|
|
||||||
|
func (tx *Transaction) EffectiveGasTipUnit(baseFee *uint256.Int) (*uint256.Int, error) {
|
||||||
|
if baseFee == nil {
|
||||||
|
return tx.GasFeeCapUint(), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
var err error
|
||||||
|
|
||||||
|
gasFeeCap := tx.GasFeeCapUint().Clone()
|
||||||
|
|
||||||
|
if gasFeeCap.Lt(baseFee) {
|
||||||
|
err = ErrGasFeeCapTooLow
|
||||||
|
}
|
||||||
|
|
||||||
|
gasTipCapUint := tx.GasTipCapUint()
|
||||||
|
|
||||||
|
if gasFeeCap.Lt(gasTipCapUint) {
|
||||||
|
return gasFeeCap, err
|
||||||
|
}
|
||||||
|
|
||||||
|
if gasFeeCap.Lt(gasTipCapUint) && baseFee.IsZero() {
|
||||||
|
return gasFeeCap, err
|
||||||
|
}
|
||||||
|
|
||||||
|
gasFeeCap.Sub(gasFeeCap, baseFee)
|
||||||
|
|
||||||
|
if gasFeeCap.Gt(gasTipCapUint) || gasFeeCap.Eq(gasTipCapUint) {
|
||||||
|
gasFeeCap.Add(gasFeeCap, baseFee)
|
||||||
|
|
||||||
|
return gasTipCapUint, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return gasFeeCap, err
|
||||||
|
}
|
||||||
|
|
||||||
// Hash returns the transaction hash.
|
// Hash returns the transaction hash.
|
||||||
func (tx *Transaction) Hash() common.Hash {
|
func (tx *Transaction) Hash() common.Hash {
|
||||||
if hash := tx.hash.Load(); hash != nil {
|
if hash := tx.hash.Load(); hash != nil {
|
||||||
return hash.(common.Hash)
|
return *hash
|
||||||
}
|
}
|
||||||
|
|
||||||
var h common.Hash
|
var h common.Hash
|
||||||
|
|
@ -368,7 +469,9 @@ func (tx *Transaction) Hash() common.Hash {
|
||||||
} else {
|
} else {
|
||||||
h = prefixedRlpHash(tx.Type(), tx.inner)
|
h = prefixedRlpHash(tx.Type(), tx.inner)
|
||||||
}
|
}
|
||||||
tx.hash.Store(h)
|
|
||||||
|
tx.hash.Store(&h)
|
||||||
|
|
||||||
return h
|
return h
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -376,11 +479,14 @@ func (tx *Transaction) Hash() common.Hash {
|
||||||
// encoding and returning it, or returning a previously cached value.
|
// encoding and returning it, or returning a previously cached value.
|
||||||
func (tx *Transaction) Size() common.StorageSize {
|
func (tx *Transaction) Size() common.StorageSize {
|
||||||
if size := tx.size.Load(); size != nil {
|
if size := tx.size.Load(); size != nil {
|
||||||
return size.(common.StorageSize)
|
return *size
|
||||||
}
|
}
|
||||||
|
|
||||||
c := writeCounter(0)
|
c := writeCounter(0)
|
||||||
|
|
||||||
rlp.Encode(&c, &tx.inner)
|
rlp.Encode(&c, &tx.inner)
|
||||||
tx.size.Store(common.StorageSize(c))
|
tx.size.Store((*common.StorageSize)(&c))
|
||||||
|
|
||||||
return common.StorageSize(c)
|
return common.StorageSize(c)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -444,14 +550,14 @@ func (s TxByNonce) Swap(i, j int) { s[i], s[j] = s[j], s[i] }
|
||||||
// TxWithMinerFee wraps a transaction with its gas price or effective miner gasTipCap
|
// TxWithMinerFee wraps a transaction with its gas price or effective miner gasTipCap
|
||||||
type TxWithMinerFee struct {
|
type TxWithMinerFee struct {
|
||||||
tx *Transaction
|
tx *Transaction
|
||||||
minerFee *big.Int
|
minerFee *uint256.Int
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewTxWithMinerFee creates a wrapped transaction, calculating the effective
|
// NewTxWithMinerFee creates a wrapped transaction, calculating the effective
|
||||||
// miner gasTipCap if a base fee is provided.
|
// miner gasTipCap if a base fee is provided.
|
||||||
// Returns error in case of a negative effective miner gasTipCap.
|
// Returns error in case of a negative effective miner gasTipCap.
|
||||||
func NewTxWithMinerFee(tx *Transaction, baseFee *big.Int) (*TxWithMinerFee, error) {
|
func NewTxWithMinerFee(tx *Transaction, baseFee *uint256.Int) (*TxWithMinerFee, error) {
|
||||||
minerFee, err := tx.EffectiveGasTip(baseFee)
|
minerFee, err := tx.EffectiveGasTipUnit(baseFee)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
@ -496,7 +602,7 @@ type TransactionsByPriceAndNonce struct {
|
||||||
txs map[common.Address]Transactions // Per account nonce-sorted list of transactions
|
txs map[common.Address]Transactions // Per account nonce-sorted list of transactions
|
||||||
heads TxByPriceAndTime // Next transaction for each unique account (price heap)
|
heads TxByPriceAndTime // Next transaction for each unique account (price heap)
|
||||||
signer Signer // Signer for the set of transactions
|
signer Signer // Signer for the set of transactions
|
||||||
baseFee *big.Int // Current base fee
|
baseFee *uint256.Int // Current base fee
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewTransactionsByPriceAndNonce creates a transaction set that can retrieve
|
// NewTransactionsByPriceAndNonce creates a transaction set that can retrieve
|
||||||
|
|
@ -504,6 +610,7 @@ type TransactionsByPriceAndNonce struct {
|
||||||
//
|
//
|
||||||
// Note, the input map is reowned so the caller should not interact any more with
|
// Note, the input map is reowned so the caller should not interact any more with
|
||||||
// if after providing it to the constructor.
|
// if after providing it to the constructor.
|
||||||
|
/*
|
||||||
func NewTransactionsByPriceAndNonce(signer Signer, txs map[common.Address]Transactions, baseFee *big.Int) *TransactionsByPriceAndNonce {
|
func NewTransactionsByPriceAndNonce(signer Signer, txs map[common.Address]Transactions, baseFee *big.Int) *TransactionsByPriceAndNonce {
|
||||||
// Initialize a price and received time based heap with the head transactions
|
// Initialize a price and received time based heap with the head transactions
|
||||||
heads := make(TxByPriceAndTime, 0, len(txs))
|
heads := make(TxByPriceAndTime, 0, len(txs))
|
||||||
|
|
@ -524,6 +631,39 @@ func NewTransactionsByPriceAndNonce(signer Signer, txs map[common.Address]Transa
|
||||||
}
|
}
|
||||||
heap.Init(&heads)
|
heap.Init(&heads)
|
||||||
|
|
||||||
|
// Assemble and return the transaction set
|
||||||
|
return &TransactionsByPriceAndNonce{
|
||||||
|
txs: txs,
|
||||||
|
heads: heads,
|
||||||
|
signer: signer,
|
||||||
|
baseFee: baseFee,
|
||||||
|
}
|
||||||
|
}*/
|
||||||
|
|
||||||
|
func NewTransactionsByPriceAndNonce(signer Signer, txs map[common.Address]Transactions, baseFee *uint256.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 {
|
||||||
|
if len(accTxs) == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
acc, _ := Sender(signer, accTxs[0])
|
||||||
|
wrapped, err := NewTxWithMinerFee(accTxs[0], baseFee)
|
||||||
|
|
||||||
|
// Remove transaction if sender doesn't match from, or if wrapping fails.
|
||||||
|
if acc != from || err != nil {
|
||||||
|
delete(txs, from)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
heads = append(heads, wrapped)
|
||||||
|
txs[from] = accTxs[1:]
|
||||||
|
}
|
||||||
|
|
||||||
|
heap.Init(&heads)
|
||||||
|
|
||||||
// Assemble and return the transaction set
|
// Assemble and return the transaction set
|
||||||
return &TransactionsByPriceAndNonce{
|
return &TransactionsByPriceAndNonce{
|
||||||
txs: txs,
|
txs: txs,
|
||||||
|
|
|
||||||
|
|
@ -130,12 +130,11 @@ func MustSignNewTx(prv *ecdsa.PrivateKey, s Signer, txdata TxData) *Transaction
|
||||||
// not match the signer used in the current call.
|
// not match the signer used in the current call.
|
||||||
func Sender(signer Signer, tx *Transaction) (common.Address, error) {
|
func Sender(signer Signer, tx *Transaction) (common.Address, error) {
|
||||||
if sc := tx.from.Load(); sc != nil {
|
if sc := tx.from.Load(); sc != nil {
|
||||||
sigCache := sc.(sigCache)
|
|
||||||
// If the signer used to derive from in a previous
|
// If the signer used to derive from in a previous
|
||||||
// call is not the same as used current, invalidate
|
// call is not the same as used current, invalidate
|
||||||
// the cache.
|
// the cache.
|
||||||
if sigCache.signer.Equal(signer) {
|
if sc.signer.Equal(signer) {
|
||||||
return sigCache.from, nil
|
return sc.from, nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -143,7 +142,9 @@ func Sender(signer Signer, tx *Transaction) (common.Address, error) {
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return common.Address{}, err
|
return common.Address{}, err
|
||||||
}
|
}
|
||||||
tx.from.Store(sigCache{signer: signer, from: addr})
|
|
||||||
|
tx.from.Store(&sigCache{signer: signer, from: addr})
|
||||||
|
|
||||||
return addr, nil
|
return addr, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -461,10 +462,10 @@ func (fs FrontierSigner) SignatureValues(tx *Transaction, sig []byte) (r, s, v *
|
||||||
func (fs FrontierSigner) Hash(tx *Transaction) common.Hash {
|
func (fs FrontierSigner) Hash(tx *Transaction) common.Hash {
|
||||||
return rlpHash([]interface{}{
|
return rlpHash([]interface{}{
|
||||||
tx.Nonce(),
|
tx.Nonce(),
|
||||||
tx.GasPrice(),
|
tx.GasPriceRef(),
|
||||||
tx.Gas(),
|
tx.Gas(),
|
||||||
tx.To(),
|
tx.To(),
|
||||||
tx.Value(),
|
tx.ValueRef(),
|
||||||
tx.Data(),
|
tx.Data(),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -27,7 +27,10 @@ import (
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/holiman/uint256"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
cmath "github.com/ethereum/go-ethereum/common/math"
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
"github.com/ethereum/go-ethereum/rlp"
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
)
|
)
|
||||||
|
|
@ -272,14 +275,22 @@ func TestTransactionPriceNonceSort1559(t *testing.T) {
|
||||||
// Tests that transactions can be correctly sorted according to their price in
|
// 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
|
// decreasing order, but at the same time with increasing nonces when issued by
|
||||||
// the same account.
|
// the same account.
|
||||||
func testTransactionPriceNonceSort(t *testing.T, baseFee *big.Int) {
|
//
|
||||||
|
//nolint:gocognit,thelper
|
||||||
|
func testTransactionPriceNonceSort(t *testing.T, baseFeeBig *big.Int) {
|
||||||
// Generate a batch of accounts to start with
|
// Generate a batch of accounts to start with
|
||||||
keys := make([]*ecdsa.PrivateKey, 25)
|
keys := make([]*ecdsa.PrivateKey, 25)
|
||||||
for i := 0; i < len(keys); i++ {
|
for i := 0; i < len(keys); i++ {
|
||||||
keys[i], _ = crypto.GenerateKey()
|
keys[i], _ = crypto.GenerateKey()
|
||||||
}
|
}
|
||||||
|
|
||||||
signer := LatestSignerForChainID(common.Big1)
|
signer := LatestSignerForChainID(common.Big1)
|
||||||
|
|
||||||
|
var baseFee *uint256.Int
|
||||||
|
if baseFeeBig != nil {
|
||||||
|
baseFee = cmath.FromBig(baseFeeBig)
|
||||||
|
}
|
||||||
|
|
||||||
// Generate a batch of transactions with overlapping values, but shifted nonces
|
// Generate a batch of transactions with overlapping values, but shifted nonces
|
||||||
groups := map[common.Address]Transactions{}
|
groups := map[common.Address]Transactions{}
|
||||||
expectedCount := 0
|
expectedCount := 0
|
||||||
|
|
@ -308,7 +319,7 @@ func testTransactionPriceNonceSort(t *testing.T, baseFee *big.Int) {
|
||||||
GasTipCap: big.NewInt(int64(rand.Intn(gasFeeCap + 1))),
|
GasTipCap: big.NewInt(int64(rand.Intn(gasFeeCap + 1))),
|
||||||
Data: nil,
|
Data: nil,
|
||||||
})
|
})
|
||||||
if count == 25 && int64(gasFeeCap) < baseFee.Int64() {
|
if count == 25 && uint64(gasFeeCap) < baseFee.Uint64() {
|
||||||
count = i
|
count = i
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -341,12 +352,25 @@ func testTransactionPriceNonceSort(t *testing.T, baseFee *big.Int) {
|
||||||
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())
|
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 the next tx has different from account, the price must be lower than the current one
|
||||||
if i+1 < len(txs) {
|
if i+1 < len(txs) {
|
||||||
next := txs[i+1]
|
next := txs[i+1]
|
||||||
fromNext, _ := Sender(signer, next)
|
fromNext, _ := Sender(signer, next)
|
||||||
tip, err := txi.EffectiveGasTip(baseFee)
|
tip, err := txi.EffectiveGasTipUnit(baseFee)
|
||||||
nextTip, nextErr := next.EffectiveGasTip(baseFee)
|
nextTip, nextErr := next.EffectiveGasTipUnit(baseFee)
|
||||||
|
|
||||||
|
tipBig, _ := txi.EffectiveGasTip(baseFeeBig)
|
||||||
|
nextTipBig, _ := next.EffectiveGasTip(baseFeeBig)
|
||||||
|
|
||||||
|
if tip.Cmp(cmath.FromBig(tipBig)) != 0 {
|
||||||
|
t.Fatalf("EffectiveGasTip incorrect. uint256 %q, big.Int %q, baseFee %q, baseFeeBig %q", tip.String(), tipBig.String(), baseFee.String(), baseFeeBig.String())
|
||||||
|
}
|
||||||
|
|
||||||
|
if nextTip.Cmp(cmath.FromBig(nextTipBig)) != 0 {
|
||||||
|
t.Fatalf("EffectiveGasTip next incorrect. uint256 %q, big.Int %q, baseFee %q, baseFeeBig %q", nextTip.String(), nextTipBig.String(), baseFee.String(), baseFeeBig.String())
|
||||||
|
}
|
||||||
|
|
||||||
if err != nil || nextErr != nil {
|
if err != nil || nextErr != nil {
|
||||||
t.Errorf("error calculating effective tip")
|
t.Errorf("error calculating effective tip")
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -236,11 +236,18 @@ func (b *EthAPIBackend) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscri
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *EthAPIBackend) SendTx(ctx context.Context, signedTx *types.Transaction) error {
|
func (b *EthAPIBackend) SendTx(ctx context.Context, signedTx *types.Transaction) error {
|
||||||
return b.eth.txPool.AddLocal(signedTx)
|
err := b.eth.txPool.AddLocal(signedTx)
|
||||||
|
if err != nil {
|
||||||
|
if unwrapped := errors.Unwrap(err); unwrapped != nil {
|
||||||
|
return unwrapped
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *EthAPIBackend) GetPoolTransactions() (types.Transactions, error) {
|
func (b *EthAPIBackend) GetPoolTransactions() (types.Transactions, error) {
|
||||||
pending := b.eth.txPool.Pending(false)
|
pending := b.eth.txPool.Pending(context.Background(), false)
|
||||||
var txs types.Transactions
|
var txs types.Transactions
|
||||||
for _, batch := range pending {
|
for _, batch := range pending {
|
||||||
txs = append(txs, batch...)
|
txs = append(txs, batch...)
|
||||||
|
|
|
||||||
|
|
@ -26,6 +26,7 @@ func newCheckpointVerifier(verifyFn func(ctx context.Context, handler *ethHandle
|
||||||
)
|
)
|
||||||
|
|
||||||
// check if we have the checkpoint blocks
|
// check if we have the checkpoint blocks
|
||||||
|
//nolint:contextcheck
|
||||||
head := handler.ethAPI.BlockNumber()
|
head := handler.ethAPI.BlockNumber()
|
||||||
if head < hexutil.Uint64(endBlock) {
|
if head < hexutil.Uint64(endBlock) {
|
||||||
log.Debug("Head block behind checkpoint block", "head", head, "checkpoint end block", endBlock)
|
log.Debug("Head block behind checkpoint block", "head", head, "checkpoint end block", endBlock)
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package eth
|
package eth
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"math"
|
"math"
|
||||||
"math/big"
|
"math/big"
|
||||||
|
|
@ -69,7 +70,7 @@ type txPool interface {
|
||||||
|
|
||||||
// Pending should return pending transactions.
|
// Pending should return pending transactions.
|
||||||
// The slice should be modifiable by the caller.
|
// The slice should be modifiable by the caller.
|
||||||
Pending(enforceTips bool) map[common.Address]types.Transactions
|
Pending(ctx context.Context, enforceTips bool) map[common.Address]types.Transactions
|
||||||
|
|
||||||
// SubscribeNewTxsEvent should return an event subscription of
|
// SubscribeNewTxsEvent should return an event subscription of
|
||||||
// NewTxsEvent and send events to the given channel.
|
// NewTxsEvent and send events to the given channel.
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package eth
|
package eth
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"math/big"
|
"math/big"
|
||||||
"sort"
|
"sort"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
@ -92,7 +93,7 @@ func (p *testTxPool) AddRemotes(txs []*types.Transaction) []error {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Pending returns all the transactions known to the pool
|
// Pending returns all the transactions known to the pool
|
||||||
func (p *testTxPool) Pending(enforceTips bool) map[common.Address]types.Transactions {
|
func (p *testTxPool) Pending(ctx context.Context, enforceTips bool) map[common.Address]types.Transactions {
|
||||||
p.lock.RLock()
|
p.lock.RLock()
|
||||||
defer p.lock.RUnlock()
|
defer p.lock.RUnlock()
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package eth
|
package eth
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"math/big"
|
"math/big"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
|
|
@ -44,20 +45,24 @@ func (h *handler) syncTransactions(p *eth.Peer) {
|
||||||
//
|
//
|
||||||
// TODO(karalabe): Figure out if we could get away with random order somehow
|
// TODO(karalabe): Figure out if we could get away with random order somehow
|
||||||
var txs types.Transactions
|
var txs types.Transactions
|
||||||
pending := h.txpool.Pending(false)
|
|
||||||
|
pending := h.txpool.Pending(context.Background(), false)
|
||||||
for _, batch := range pending {
|
for _, batch := range pending {
|
||||||
txs = append(txs, batch...)
|
txs = append(txs, batch...)
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(txs) == 0 {
|
if len(txs) == 0 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// The eth/65 protocol introduces proper transaction announcements, so instead
|
// The eth/65 protocol introduces proper transaction announcements, so instead
|
||||||
// of dripping transactions across multiple peers, just send the entire list as
|
// of dripping transactions across multiple peers, just send the entire list as
|
||||||
// an announcement and let the remote side decide what they need (likely nothing).
|
// an announcement and let the remote side decide what they need (likely nothing).
|
||||||
|
|
||||||
hashes := make([]common.Hash, len(txs))
|
hashes := make([]common.Hash, len(txs))
|
||||||
for i, tx := range txs {
|
for i, tx := range txs {
|
||||||
hashes[i] = tx.Hash()
|
hashes[i] = tx.Hash()
|
||||||
}
|
}
|
||||||
|
|
||||||
p.AsyncSendPooledTransactionHashes(hashes)
|
p.AsyncSendPooledTransactionHashes(hashes)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -61,6 +61,28 @@ func CPUProfile(ctx context.Context, sec int) ([]byte, map[string]string, error)
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// CPUProfile generates a CPU Profile for a given duration
|
||||||
|
func CPUProfileWithChannel(done chan bool) ([]byte, map[string]string, error) {
|
||||||
|
var buf bytes.Buffer
|
||||||
|
if err := pprof.StartCPUProfile(&buf); err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-done:
|
||||||
|
case <-time.After(30 * time.Second):
|
||||||
|
}
|
||||||
|
|
||||||
|
pprof.StopCPUProfile()
|
||||||
|
|
||||||
|
return buf.Bytes(),
|
||||||
|
map[string]string{
|
||||||
|
"X-Content-Type-Options": "nosniff",
|
||||||
|
"Content-Type": "application/octet-stream",
|
||||||
|
"Content-Disposition": `attachment; filename="profile"`,
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
|
|
||||||
// Trace runs a trace profile for a given duration
|
// Trace runs a trace profile for a given duration
|
||||||
func Trace(ctx context.Context, sec int) ([]byte, map[string]string, error) {
|
func Trace(ctx context.Context, sec int) ([]byte, map[string]string, error) {
|
||||||
if sec <= 0 {
|
if sec <= 0 {
|
||||||
|
|
|
||||||
|
|
@ -21,6 +21,7 @@ import (
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/big"
|
"math/big"
|
||||||
|
"runtime"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
|
@ -2229,6 +2230,21 @@ func (api *PrivateDebugAPI) PurgeCheckpointWhitelist() {
|
||||||
api.b.PurgeCheckpointWhitelist()
|
api.b.PurgeCheckpointWhitelist()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetTraceStack returns the current trace stack
|
||||||
|
func (api *PrivateDebugAPI) GetTraceStack() string {
|
||||||
|
buf := make([]byte, 1024)
|
||||||
|
|
||||||
|
for {
|
||||||
|
n := runtime.Stack(buf, true)
|
||||||
|
|
||||||
|
if n < len(buf) {
|
||||||
|
return string(buf)
|
||||||
|
}
|
||||||
|
|
||||||
|
buf = make([]byte, 2*len(buf))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// PublicNetAPI offers network related RPC methods
|
// PublicNetAPI offers network related RPC methods
|
||||||
type PublicNetAPI struct {
|
type PublicNetAPI struct {
|
||||||
net *p2p.Server
|
net *p2p.Server
|
||||||
|
|
|
||||||
|
|
@ -484,6 +484,11 @@ web3._extend({
|
||||||
call: 'debug_purgeCheckpointWhitelist',
|
call: 'debug_purgeCheckpointWhitelist',
|
||||||
params: 0,
|
params: 0,
|
||||||
}),
|
}),
|
||||||
|
new web3._extend.Method({
|
||||||
|
name: 'getTraceStack',
|
||||||
|
call: 'debug_getTraceStack',
|
||||||
|
params: 0,
|
||||||
|
}),
|
||||||
],
|
],
|
||||||
properties: []
|
properties: []
|
||||||
});
|
});
|
||||||
|
|
|
||||||
|
|
@ -617,7 +617,7 @@ func testTransactionStatus(t *testing.T, protocol int) {
|
||||||
sendRequest(rawPeer.app, GetTxStatusMsg, reqID, []common.Hash{tx.Hash()})
|
sendRequest(rawPeer.app, GetTxStatusMsg, reqID, []common.Hash{tx.Hash()})
|
||||||
}
|
}
|
||||||
if err := expectResponse(rawPeer.app, TxStatusMsg, reqID, testBufLimit, []light.TxStatus{expStatus}); err != nil {
|
if err := expectResponse(rawPeer.app, TxStatusMsg, reqID, testBufLimit, []light.TxStatus{expStatus}); err != nil {
|
||||||
t.Errorf("transaction status mismatch")
|
t.Error("transaction status mismatch", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
signer := types.HomesteadSigner{}
|
signer := types.HomesteadSigner{}
|
||||||
|
|
|
||||||
|
|
@ -507,25 +507,39 @@ func handleSendTx(msg Decoder) (serveRequestFn, uint64, uint64, error) {
|
||||||
if err := msg.Decode(&r); err != nil {
|
if err := msg.Decode(&r); err != nil {
|
||||||
return nil, 0, 0, err
|
return nil, 0, 0, err
|
||||||
}
|
}
|
||||||
|
|
||||||
amount := uint64(len(r.Txs))
|
amount := uint64(len(r.Txs))
|
||||||
|
|
||||||
return func(backend serverBackend, p *clientPeer, waitOrStop func() bool) *reply {
|
return func(backend serverBackend, p *clientPeer, waitOrStop func() bool) *reply {
|
||||||
stats := make([]light.TxStatus, len(r.Txs))
|
stats := make([]light.TxStatus, len(r.Txs))
|
||||||
|
|
||||||
|
var (
|
||||||
|
err error
|
||||||
|
addFn func(transaction *types.Transaction) error
|
||||||
|
)
|
||||||
|
|
||||||
for i, tx := range r.Txs {
|
for i, tx := range r.Txs {
|
||||||
if i != 0 && !waitOrStop() {
|
if i != 0 && !waitOrStop() {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
hash := tx.Hash()
|
hash := tx.Hash()
|
||||||
stats[i] = txStatus(backend, hash)
|
stats[i] = txStatus(backend, hash)
|
||||||
|
|
||||||
if stats[i].Status == core.TxStatusUnknown {
|
if stats[i].Status == core.TxStatusUnknown {
|
||||||
addFn := backend.TxPool().AddRemotes
|
addFn = backend.TxPool().AddRemote
|
||||||
|
|
||||||
// Add txs synchronously for testing purpose
|
// Add txs synchronously for testing purpose
|
||||||
if backend.AddTxsSync() {
|
if backend.AddTxsSync() {
|
||||||
addFn = backend.TxPool().AddRemotesSync
|
addFn = backend.TxPool().AddRemoteSync
|
||||||
}
|
}
|
||||||
if errs := addFn([]*types.Transaction{tx}); errs[0] != nil {
|
|
||||||
stats[i].Error = errs[0].Error()
|
if err = addFn(tx); err != nil {
|
||||||
|
stats[i].Error = err.Error()
|
||||||
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
stats[i] = txStatus(backend, hash)
|
stats[i] = txStatus(backend, hash)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
192
miner/worker.go
192
miner/worker.go
|
|
@ -17,10 +17,15 @@
|
||||||
package miner
|
package miner
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/big"
|
"math/big"
|
||||||
|
"os"
|
||||||
|
"runtime"
|
||||||
|
"runtime/pprof"
|
||||||
|
ptrace "runtime/trace"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
@ -31,6 +36,7 @@ import (
|
||||||
"go.opentelemetry.io/otel/trace"
|
"go.opentelemetry.io/otel/trace"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
cmath "github.com/ethereum/go-ethereum/common/math"
|
||||||
"github.com/ethereum/go-ethereum/common/tracing"
|
"github.com/ethereum/go-ethereum/common/tracing"
|
||||||
"github.com/ethereum/go-ethereum/consensus"
|
"github.com/ethereum/go-ethereum/consensus"
|
||||||
"github.com/ethereum/go-ethereum/consensus/misc"
|
"github.com/ethereum/go-ethereum/consensus/misc"
|
||||||
|
|
@ -39,6 +45,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
"github.com/ethereum/go-ethereum/event"
|
"github.com/ethereum/go-ethereum/event"
|
||||||
"github.com/ethereum/go-ethereum/log"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
|
"github.com/ethereum/go-ethereum/metrics"
|
||||||
"github.com/ethereum/go-ethereum/params"
|
"github.com/ethereum/go-ethereum/params"
|
||||||
"github.com/ethereum/go-ethereum/trie"
|
"github.com/ethereum/go-ethereum/trie"
|
||||||
)
|
)
|
||||||
|
|
@ -83,6 +90,12 @@ const (
|
||||||
staleThreshold = 7
|
staleThreshold = 7
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// metrics gauge to track total and empty blocks sealed by a miner
|
||||||
|
var (
|
||||||
|
sealedBlocksCounter = metrics.NewRegisteredCounter("worker/sealedBlocks", nil)
|
||||||
|
sealedEmptyBlocksCounter = metrics.NewRegisteredCounter("worker/sealedEmptyBlocks", nil)
|
||||||
|
)
|
||||||
|
|
||||||
// environment is the worker's current environment and holds all
|
// environment is the worker's current environment and holds all
|
||||||
// information of the sealing block generation.
|
// information of the sealing block generation.
|
||||||
type environment struct {
|
type environment struct {
|
||||||
|
|
@ -257,6 +270,8 @@ type worker struct {
|
||||||
skipSealHook func(*task) bool // Method to decide whether skipping the sealing.
|
skipSealHook func(*task) bool // Method to decide whether skipping the sealing.
|
||||||
fullTaskHook func() // Method to call before pushing the full sealing task.
|
fullTaskHook func() // Method to call before pushing the full sealing task.
|
||||||
resubmitHook func(time.Duration, time.Duration) // Method to call upon updating resubmitting interval.
|
resubmitHook func(time.Duration, time.Duration) // Method to call upon updating resubmitting interval.
|
||||||
|
|
||||||
|
profileCount *int32 // Global count for profiling
|
||||||
}
|
}
|
||||||
|
|
||||||
//nolint:staticcheck
|
//nolint:staticcheck
|
||||||
|
|
@ -285,6 +300,7 @@ func newWorker(config *Config, chainConfig *params.ChainConfig, engine consensus
|
||||||
resubmitIntervalCh: make(chan time.Duration),
|
resubmitIntervalCh: make(chan time.Duration),
|
||||||
resubmitAdjustCh: make(chan *intervalAdjust, resubmitAdjustChanSize),
|
resubmitAdjustCh: make(chan *intervalAdjust, resubmitAdjustChanSize),
|
||||||
}
|
}
|
||||||
|
worker.profileCount = new(int32)
|
||||||
// Subscribe NewTxsEvent for tx pool
|
// Subscribe NewTxsEvent for tx pool
|
||||||
worker.txsSub = eth.TxPool().SubscribeNewTxsEvent(worker.txsCh)
|
worker.txsSub = eth.TxPool().SubscribeNewTxsEvent(worker.txsCh)
|
||||||
// Subscribe events for blockchain
|
// Subscribe events for blockchain
|
||||||
|
|
@ -560,9 +576,11 @@ func (w *worker) mainLoop(ctx context.Context) {
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case req := <-w.newWorkCh:
|
case req := <-w.newWorkCh:
|
||||||
|
//nolint:contextcheck
|
||||||
w.commitWork(req.ctx, req.interrupt, req.noempty, req.timestamp)
|
w.commitWork(req.ctx, req.interrupt, req.noempty, req.timestamp)
|
||||||
|
|
||||||
case req := <-w.getWorkCh:
|
case req := <-w.getWorkCh:
|
||||||
|
//nolint:contextcheck
|
||||||
block, err := w.generateWork(req.ctx, req.params)
|
block, err := w.generateWork(req.ctx, req.params)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
req.err = err
|
req.err = err
|
||||||
|
|
@ -622,13 +640,17 @@ func (w *worker) mainLoop(ctx context.Context) {
|
||||||
if gp := w.current.gasPool; gp != nil && gp.Gas() < params.TxGas {
|
if gp := w.current.gasPool; gp != nil && gp.Gas() < params.TxGas {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
txs := make(map[common.Address]types.Transactions)
|
txs := make(map[common.Address]types.Transactions)
|
||||||
|
|
||||||
for _, tx := range ev.Txs {
|
for _, tx := range ev.Txs {
|
||||||
acc, _ := types.Sender(w.current.signer, tx)
|
acc, _ := types.Sender(w.current.signer, tx)
|
||||||
txs[acc] = append(txs[acc], tx)
|
txs[acc] = append(txs[acc], tx)
|
||||||
}
|
}
|
||||||
txset := types.NewTransactionsByPriceAndNonce(w.current.signer, txs, w.current.header.BaseFee)
|
|
||||||
|
txset := types.NewTransactionsByPriceAndNonce(w.current.signer, txs, cmath.FromBig(w.current.header.BaseFee))
|
||||||
tcount := w.current.tcount
|
tcount := w.current.tcount
|
||||||
|
|
||||||
w.commitTransactions(w.current, txset, nil)
|
w.commitTransactions(w.current, txset, nil)
|
||||||
|
|
||||||
// Only update the snapshot if any new transactions were added
|
// Only update the snapshot if any new transactions were added
|
||||||
|
|
@ -758,7 +780,7 @@ func (w *worker) resultLoop() {
|
||||||
err error
|
err error
|
||||||
)
|
)
|
||||||
|
|
||||||
tracing.Exec(task.ctx, "resultLoop", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(task.ctx, "", "resultLoop", func(ctx context.Context, span trace.Span) {
|
||||||
for i, taskReceipt := range task.receipts {
|
for i, taskReceipt := range task.receipts {
|
||||||
receipt := new(types.Receipt)
|
receipt := new(types.Receipt)
|
||||||
receipts[i] = receipt
|
receipts[i] = receipt
|
||||||
|
|
@ -808,6 +830,12 @@ func (w *worker) resultLoop() {
|
||||||
// Broadcast the block and announce chain insertion event
|
// Broadcast the block and announce chain insertion event
|
||||||
w.mux.Post(core.NewMinedBlockEvent{Block: block})
|
w.mux.Post(core.NewMinedBlockEvent{Block: block})
|
||||||
|
|
||||||
|
sealedBlocksCounter.Inc(1)
|
||||||
|
|
||||||
|
if block.Transactions().Len() == 0 {
|
||||||
|
sealedEmptyBlocksCounter.Inc(1)
|
||||||
|
}
|
||||||
|
|
||||||
// Insert the block into the set of pending ones to resultLoop for confirmations
|
// Insert the block into the set of pending ones to resultLoop for confirmations
|
||||||
w.unconfirmed.Insert(block.NumberU64(), block.Hash())
|
w.unconfirmed.Insert(block.NumberU64(), block.Hash())
|
||||||
case <-w.exitCh:
|
case <-w.exitCh:
|
||||||
|
|
@ -965,7 +993,10 @@ func (w *worker) commitTransactions(env *environment, txs *types.TransactionsByP
|
||||||
// Start executing the transaction
|
// Start executing the transaction
|
||||||
env.state.Prepare(tx.Hash(), env.tcount)
|
env.state.Prepare(tx.Hash(), env.tcount)
|
||||||
|
|
||||||
|
start := time.Now()
|
||||||
|
|
||||||
logs, err := w.commitTransaction(env, tx)
|
logs, err := w.commitTransaction(env, tx)
|
||||||
|
|
||||||
switch {
|
switch {
|
||||||
case errors.Is(err, core.ErrGasLimitReached):
|
case errors.Is(err, core.ErrGasLimitReached):
|
||||||
// Pop the current out-of-gas transaction without shifting in the next from the account
|
// Pop the current out-of-gas transaction without shifting in the next from the account
|
||||||
|
|
@ -987,6 +1018,7 @@ func (w *worker) commitTransactions(env *environment, txs *types.TransactionsByP
|
||||||
coalescedLogs = append(coalescedLogs, logs...)
|
coalescedLogs = append(coalescedLogs, logs...)
|
||||||
env.tcount++
|
env.tcount++
|
||||||
txs.Shift()
|
txs.Shift()
|
||||||
|
log.Info("Committed new tx", "tx hash", tx.Hash(), "from", from, "to", tx.To(), "nonce", tx.Nonce(), "gas", tx.Gas(), "gasPrice", tx.GasPrice(), "value", tx.Value(), "time spent", time.Since(start))
|
||||||
|
|
||||||
case errors.Is(err, core.ErrTxTypeNotSupported):
|
case errors.Is(err, core.ErrTxTypeNotSupported):
|
||||||
// Pop the unsupported transaction without shifting in the next from the account
|
// Pop the unsupported transaction without shifting in the next from the account
|
||||||
|
|
@ -1077,7 +1109,7 @@ func (w *worker) prepareWork(genParams *generateParams) (*environment, error) {
|
||||||
}
|
}
|
||||||
// Set baseFee and GasLimit if we are on an EIP-1559 chain
|
// Set baseFee and GasLimit if we are on an EIP-1559 chain
|
||||||
if w.chainConfig.IsLondon(header.Number) {
|
if w.chainConfig.IsLondon(header.Number) {
|
||||||
header.BaseFee = misc.CalcBaseFee(w.chainConfig, parent.Header())
|
header.BaseFee = misc.CalcBaseFeeUint(w.chainConfig, parent.Header()).ToBig()
|
||||||
if !w.chainConfig.IsLondon(parent.Number()) {
|
if !w.chainConfig.IsLondon(parent.Number()) {
|
||||||
parentGasLimit := parent.GasLimit() * params.ElasticityMultiplier
|
parentGasLimit := parent.GasLimit() * params.ElasticityMultiplier
|
||||||
header.GasLimit = core.CalcGasLimit(parentGasLimit, w.config.GasCeil)
|
header.GasLimit = core.CalcGasLimit(parentGasLimit, w.config.GasCeil)
|
||||||
|
|
@ -1117,9 +1149,75 @@ func (w *worker) prepareWork(genParams *generateParams) (*environment, error) {
|
||||||
return env, nil
|
return env, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func startProfiler(profile string, filepath string, number uint64) (func() error, error) {
|
||||||
|
var (
|
||||||
|
buf bytes.Buffer
|
||||||
|
err error
|
||||||
|
)
|
||||||
|
|
||||||
|
closeFn := func() {}
|
||||||
|
|
||||||
|
switch profile {
|
||||||
|
case "cpu":
|
||||||
|
err = pprof.StartCPUProfile(&buf)
|
||||||
|
|
||||||
|
if err == nil {
|
||||||
|
closeFn = func() {
|
||||||
|
pprof.StopCPUProfile()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
case "trace":
|
||||||
|
err = ptrace.Start(&buf)
|
||||||
|
|
||||||
|
if err == nil {
|
||||||
|
closeFn = func() {
|
||||||
|
ptrace.Stop()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
case "heap":
|
||||||
|
runtime.GC()
|
||||||
|
|
||||||
|
err = pprof.WriteHeapProfile(&buf)
|
||||||
|
default:
|
||||||
|
log.Info("Incorrect profile name")
|
||||||
|
}
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
return func() error {
|
||||||
|
closeFn()
|
||||||
|
return nil
|
||||||
|
}, err
|
||||||
|
}
|
||||||
|
|
||||||
|
closeFnNew := func() error {
|
||||||
|
var err error
|
||||||
|
|
||||||
|
closeFn()
|
||||||
|
|
||||||
|
if buf.Len() == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
f, err := os.Create(filepath + "/" + profile + "-" + fmt.Sprint(number) + ".prof")
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
defer f.Close()
|
||||||
|
|
||||||
|
_, err = f.Write(buf.Bytes())
|
||||||
|
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return closeFnNew, nil
|
||||||
|
}
|
||||||
|
|
||||||
// fillTransactions retrieves the pending transactions from the txpool and fills them
|
// fillTransactions retrieves the pending transactions from the txpool and fills them
|
||||||
// into the given sealing block. The transaction selection and ordering strategy can
|
// into the given sealing block. The transaction selection and ordering strategy can
|
||||||
// be customized with the plugin in the future.
|
// be customized with the plugin in the future.
|
||||||
|
//
|
||||||
|
//nolint:gocognit
|
||||||
func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *environment) {
|
func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *environment) {
|
||||||
ctx, span := tracing.StartSpan(ctx, "fillTransactions")
|
ctx, span := tracing.StartSpan(ctx, "fillTransactions")
|
||||||
defer tracing.EndSpan(span)
|
defer tracing.EndSpan(span)
|
||||||
|
|
@ -1134,10 +1232,76 @@ func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *en
|
||||||
remoteTxs map[common.Address]types.Transactions
|
remoteTxs map[common.Address]types.Transactions
|
||||||
)
|
)
|
||||||
|
|
||||||
tracing.Exec(ctx, "worker.SplittingTransactions", func(ctx context.Context, span trace.Span) {
|
// TODO: move to config or RPC
|
||||||
pending := w.eth.TxPool().Pending(true)
|
const profiling = false
|
||||||
|
|
||||||
|
if profiling {
|
||||||
|
doneCh := make(chan struct{})
|
||||||
|
|
||||||
|
defer func() {
|
||||||
|
close(doneCh)
|
||||||
|
}()
|
||||||
|
|
||||||
|
go func(number uint64) {
|
||||||
|
closeFn := func() error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-time.After(150 * time.Millisecond):
|
||||||
|
// Check if we've not crossed limit
|
||||||
|
if attempt := atomic.AddInt32(w.profileCount, 1); attempt >= 10 {
|
||||||
|
log.Info("Completed profiling", "attempt", attempt)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
log.Info("Starting profiling in fill transactions", "number", number)
|
||||||
|
|
||||||
|
dir, err := os.MkdirTemp("", fmt.Sprintf("bor-traces-%s-", time.Now().UTC().Format("2006-01-02-150405Z")))
|
||||||
|
if err != nil {
|
||||||
|
log.Error("Error in profiling", "path", dir, "number", number, "err", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// grab the cpu profile
|
||||||
|
closeFnInternal, err := startProfiler("cpu", dir, number)
|
||||||
|
if err != nil {
|
||||||
|
log.Error("Error in profiling", "path", dir, "number", number, "err", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
closeFn = func() error {
|
||||||
|
err := closeFnInternal()
|
||||||
|
|
||||||
|
log.Info("Completed profiling", "path", dir, "number", number, "error", err)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
case <-doneCh:
|
||||||
|
err := closeFn()
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
log.Info("closing fillTransactions", "number", number, "error", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}(env.header.Number.Uint64())
|
||||||
|
}
|
||||||
|
|
||||||
|
tracing.Exec(ctx, "", "worker.SplittingTransactions", func(ctx context.Context, span trace.Span) {
|
||||||
|
|
||||||
|
prePendingTime := time.Now()
|
||||||
|
|
||||||
|
pending := w.eth.TxPool().Pending(ctx, true)
|
||||||
remoteTxs = pending
|
remoteTxs = pending
|
||||||
|
|
||||||
|
postPendingTime := time.Now()
|
||||||
|
|
||||||
for _, account := range w.eth.TxPool().Locals() {
|
for _, account := range w.eth.TxPool().Locals() {
|
||||||
if txs := remoteTxs[account]; len(txs) > 0 {
|
if txs := remoteTxs[account]; len(txs) > 0 {
|
||||||
delete(remoteTxs, account)
|
delete(remoteTxs, account)
|
||||||
|
|
@ -1145,6 +1309,8 @@ func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *en
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
postLocalsTime := time.Now()
|
||||||
|
|
||||||
localTxsCount = len(localTxs)
|
localTxsCount = len(localTxs)
|
||||||
remoteTxsCount = len(remoteTxs)
|
remoteTxsCount = len(remoteTxs)
|
||||||
|
|
||||||
|
|
@ -1152,6 +1318,8 @@ func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *en
|
||||||
span,
|
span,
|
||||||
attribute.Int("len of local txs", localTxsCount),
|
attribute.Int("len of local txs", localTxsCount),
|
||||||
attribute.Int("len of remote txs", remoteTxsCount),
|
attribute.Int("len of remote txs", remoteTxsCount),
|
||||||
|
attribute.String("time taken by Pending()", fmt.Sprintf("%v", postPendingTime.Sub(prePendingTime))),
|
||||||
|
attribute.String("time taken by Locals()", fmt.Sprintf("%v", postLocalsTime.Sub(postPendingTime))),
|
||||||
)
|
)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -1164,8 +1332,8 @@ func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *en
|
||||||
if localTxsCount > 0 {
|
if localTxsCount > 0 {
|
||||||
var txs *types.TransactionsByPriceAndNonce
|
var txs *types.TransactionsByPriceAndNonce
|
||||||
|
|
||||||
tracing.Exec(ctx, "worker.LocalTransactionsByPriceAndNonce", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(ctx, "", "worker.LocalTransactionsByPriceAndNonce", func(ctx context.Context, span trace.Span) {
|
||||||
txs = types.NewTransactionsByPriceAndNonce(env.signer, localTxs, env.header.BaseFee)
|
txs = types.NewTransactionsByPriceAndNonce(env.signer, localTxs, cmath.FromBig(env.header.BaseFee))
|
||||||
|
|
||||||
tracing.SetAttributes(
|
tracing.SetAttributes(
|
||||||
span,
|
span,
|
||||||
|
|
@ -1173,7 +1341,7 @@ func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *en
|
||||||
)
|
)
|
||||||
})
|
})
|
||||||
|
|
||||||
tracing.Exec(ctx, "worker.LocalCommitTransactions", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(ctx, "", "worker.LocalCommitTransactions", func(ctx context.Context, span trace.Span) {
|
||||||
committed = w.commitTransactions(env, txs, interrupt)
|
committed = w.commitTransactions(env, txs, interrupt)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -1187,8 +1355,8 @@ func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *en
|
||||||
if remoteTxsCount > 0 {
|
if remoteTxsCount > 0 {
|
||||||
var txs *types.TransactionsByPriceAndNonce
|
var txs *types.TransactionsByPriceAndNonce
|
||||||
|
|
||||||
tracing.Exec(ctx, "worker.RemoteTransactionsByPriceAndNonce", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(ctx, "", "worker.RemoteTransactionsByPriceAndNonce", func(ctx context.Context, span trace.Span) {
|
||||||
txs = types.NewTransactionsByPriceAndNonce(env.signer, remoteTxs, env.header.BaseFee)
|
txs = types.NewTransactionsByPriceAndNonce(env.signer, remoteTxs, cmath.FromBig(env.header.BaseFee))
|
||||||
|
|
||||||
tracing.SetAttributes(
|
tracing.SetAttributes(
|
||||||
span,
|
span,
|
||||||
|
|
@ -1196,7 +1364,7 @@ func (w *worker) fillTransactions(ctx context.Context, interrupt *int32, env *en
|
||||||
)
|
)
|
||||||
})
|
})
|
||||||
|
|
||||||
tracing.Exec(ctx, "worker.RemoteCommitTransactions", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(ctx, "", "worker.RemoteCommitTransactions", func(ctx context.Context, span trace.Span) {
|
||||||
committed = w.commitTransactions(env, txs, interrupt)
|
committed = w.commitTransactions(env, txs, interrupt)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -1237,7 +1405,7 @@ func (w *worker) commitWork(ctx context.Context, interrupt *int32, noempty bool,
|
||||||
err error
|
err error
|
||||||
)
|
)
|
||||||
|
|
||||||
tracing.Exec(ctx, "worker.prepareWork", func(ctx context.Context, span trace.Span) {
|
tracing.Exec(ctx, "", "worker.prepareWork", func(ctx context.Context, span trace.Span) {
|
||||||
// Set the coinbase if the worker is running or it's required
|
// Set the coinbase if the worker is running or it's required
|
||||||
var coinbase common.Address
|
var coinbase common.Address
|
||||||
if w.isRunning() {
|
if w.isRunning() {
|
||||||
|
|
|
||||||
|
|
@ -141,9 +141,6 @@ func (tm *testMatcher) findSkip(name string) (reason string, skipload bool) {
|
||||||
isWin32 := runtime.GOARCH == "386" && runtime.GOOS == "windows"
|
isWin32 := runtime.GOARCH == "386" && runtime.GOOS == "windows"
|
||||||
for _, re := range tm.slowpat {
|
for _, re := range tm.slowpat {
|
||||||
if re.MatchString(name) {
|
if re.MatchString(name) {
|
||||||
if testing.Short() {
|
|
||||||
return "skipped in -short mode", false
|
|
||||||
}
|
|
||||||
if isWin32 {
|
if isWin32 {
|
||||||
return "skipped on 32bit windows", false
|
return "skipped on 32bit windows", false
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue