Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 20 additions & 1 deletion giga/evmonly/executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"math/big"
"sync"
"sync/atomic"
"time"

"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core"
Expand All @@ -29,6 +30,7 @@ type Executor struct {
cfg Config
resultSink ResultSink
occPool *occWorkerPool
parseSizer *parseSizer
resultPool *blockResultPool
stateDBPool sync.Pool
storeMu sync.Mutex
Expand Down Expand Up @@ -102,6 +104,7 @@ func NewExecutor(cfg Config, opts ...Option) *Executor {
resultPool: newBlockResultPool(cfg.BlockResultPoolSize),
blockPhases: seidbmetrics.NewPhaseTimer(otel.Meter(executorMeterName), "evmonly_block"),
}
e.parseSizer = newParseSizer(e.cfg.ParseWorkers)
if e.cfg.OCCWorkers > 1 {
e.occPool = newOCCWorkerPool(e.cfg.OCCWorkers)
}
Expand Down Expand Up @@ -176,7 +179,18 @@ func (e *Executor) ExecuteBlock(ctx context.Context, req BlockRequest) (*BlockRe
return result, nil
}

// PrepareBlock decodes the block's transactions and recovers their senders on
// every parse worker.
func (e *Executor) PrepareBlock(ctx context.Context, req BlockRequest) (PreparedBlock, error) {
return e.PrepareBlockWithin(ctx, req, 0)
}

// PrepareBlockWithin decodes the block's transactions and recovers their senders
// on as few parse workers as the decode is expected to fit in budget on, leaving
// the rest of the processors to whatever runs alongside. A budget of 0 uses every
// parse worker and does not inform the expectation, which comes from the budgeted
// decodes before this one; the first of those is decoded on every worker.
func (e *Executor) PrepareBlockWithin(ctx context.Context, req BlockRequest, budget time.Duration) (PreparedBlock, error) {
chainConfig := e.chainConfig(req.Context)
if err := validateBlockContext(chainConfig, req.Context); err != nil {
return PreparedBlock{}, err
Expand All @@ -185,10 +199,15 @@ func (e *Executor) PrepareBlock(ctx context.Context, req BlockRequest) (Prepared
if len(req.Senders) != 0 && len(req.Senders) != len(req.Txs) {
return PreparedBlock{}, fmt.Errorf("block request has %d senders for %d txs", len(req.Senders), len(req.Txs))
}
parsed, err := parseBlockTxs(ctx, req.Txs, signer, req.Senders, e.cfg.ParseWorkers)
workers := e.parseSizer.workers(len(req.Txs), budget)
start := time.Now()
parsed, err := parseBlockTxs(ctx, req.Txs, signer, req.Senders, workers)
if err != nil {
return PreparedBlock{}, err
}
if budget > 0 {
e.parseSizer.observe(len(req.Txs), workers, time.Since(start))
}
return PreparedBlock{
Context: req.Context,
Txs: parsed,
Expand Down
65 changes: 65 additions & 0 deletions giga/evmonly/parse_sizer.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
package evmonly

import (
"sync/atomic"
"time"
)

// perTxCostDecay is the denominator of the exponential moving average the
// per-transaction decode cost falls by; a block that measures below the estimate
// moves it 1/perTxCostDecay of the way down. A block that measures above it
// replaces it outright: a decode sized too small holds up the block that needs it,
// one sized too large only spends processors.
const perTxCostDecay = 8

// parseSizer picks how many workers decode a block from the block's size, the
// time available to decode it, and a running estimate of the per-transaction
// decode cost measured on the blocks before it.
type parseSizer struct {
maxWorkers int
// perTx is the estimated worker time to decode one transaction, in
// nanoseconds: the wall time of a decode times the workers it ran on, per
// transaction, so it includes the share of the processors those workers
// were given. 0 until the first block has been measured.
perTx atomic.Int64
}

func newParseSizer(maxWorkers int) *parseSizer {
return &parseSizer{maxWorkers: max(maxWorkers, 1)}
}

// workers returns how many workers decode txs transactions within budget: the
// fewest the estimated cost fits in, and every worker when there is no budget
// or no estimate yet.
func (s *parseSizer) workers(txs int, budget time.Duration) int {
if txs <= 1 {
return 1
}
perTx := s.perTx.Load()
if budget <= 0 || perTx <= 0 {
return min(s.maxWorkers, txs)
}
needed := (int64(txs)*perTx + int64(budget) - 1) / int64(budget)
return int(max(1, min(needed, int64(min(s.maxWorkers, txs)))))
}

// observe folds a decode of txs transactions on workers workers that took
// elapsed into the per-transaction cost estimate: a costlier decode than
// estimated raises the estimate to what it measured, a cheaper one lowers it
// gradually.
func (s *parseSizer) observe(txs, workers int, elapsed time.Duration) {
if txs <= 0 || workers <= 0 || elapsed <= 0 {
return
}
measured := int64(elapsed) * int64(workers) / int64(txs)
for {
current := s.perTx.Load()
next := measured
if current > measured {
next = current - (current-measured)/perTxCostDecay
}
if s.perTx.CompareAndSwap(current, next) {
return
}
}
}
55 changes: 55 additions & 0 deletions giga/evmonly/parse_sizer_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
package evmonly

import (
"testing"
"time"

"github.com/stretchr/testify/require"
)

func TestParseSizerUsesEveryWorkerUntilItHasAnEstimate(t *testing.T) {
sizer := newParseSizer(8)
require.Equal(t, 8, sizer.workers(1000, time.Millisecond))
require.Equal(t, 5, sizer.workers(5, time.Millisecond))
require.Equal(t, 1, sizer.workers(1, time.Millisecond))
require.Equal(t, 1, sizer.workers(0, time.Millisecond))
}

func TestParseSizerFitsTheEstimatedCostInTheBudget(t *testing.T) {
sizer := newParseSizer(32)
// 1000 txs on 4 workers in 12.5ms: 50µs per tx.
sizer.observe(1000, 4, 12500*time.Microsecond)

// 1000 txs is 50ms of work: 5 workers fit it in 10ms, 3 in 20ms, 1 in 50ms or more.
require.Equal(t, 5, sizer.workers(1000, 10*time.Millisecond))
require.Equal(t, 3, sizer.workers(1000, 20*time.Millisecond))
require.Equal(t, 1, sizer.workers(1000, 50*time.Millisecond))
require.Equal(t, 1, sizer.workers(1000, time.Second))
// A block too large for the budget gets every worker, but never more than it has txs.
require.Equal(t, 32, sizer.workers(5000, time.Millisecond))
require.Equal(t, 10, sizer.workers(10, time.Microsecond))
}

func TestParseSizerWithoutABudgetUsesEveryWorker(t *testing.T) {
sizer := newParseSizer(8)
sizer.observe(1000, 8, time.Millisecond)
require.Equal(t, 8, sizer.workers(1000, 0))
require.Equal(t, 8, sizer.workers(1000, -time.Second))
}

func TestParseSizerFollowsTheMeasuredCost(t *testing.T) {
sizer := newParseSizer(64)
sizer.observe(1000, 1, 10*time.Millisecond)
require.Equal(t, int64(10*time.Microsecond), sizer.perTx.Load())
// A costlier decode raises the estimate to what it measured at once.
sizer.observe(1000, 1, 90*time.Millisecond)
require.Equal(t, int64(90*time.Microsecond), sizer.perTx.Load())
// A cheaper one lowers it a step at a time.
sizer.observe(1000, 1, 10*time.Millisecond)
require.Equal(t, int64(80*time.Microsecond), sizer.perTx.Load())
// Zero-sized observations are ignored.
sizer.observe(0, 1, time.Millisecond)
sizer.observe(1000, 0, time.Millisecond)
sizer.observe(1000, 1, 0)
require.Equal(t, int64(80*time.Microsecond), sizer.perTx.Load())
}
48 changes: 46 additions & 2 deletions sei-tendermint/internal/evmonlyapp/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
"slices"
"sync"
"sync/atomic"
"time"

"go.opentelemetry.io/otel/attribute"
otelmetric "go.opentelemetry.io/otel/metric"
Expand Down Expand Up @@ -55,6 +56,36 @@ const checkedSendersCap = 1 << 18
// minTxsPerHashWorker is the minimum transaction count assigned to a hash worker.
const minTxsPerHashWorker = 64

// prepareBudgetShare is the fraction of the typical block execution time
// PrepareBlock is given to decode the next block in. Decoding runs alongside the
// current block's OCC speculation, which holds a worker per processor, so it is
// sized to finish within that block on as few processors as it can rather than
// contending for all of them; the share leaves room for the block being shorter
// than typical.
const prepareBudgetShare = 2

// executeEstimateDecay is the denominator of the exponential moving average of
// block execution time; each block moves the estimate 1/executeEstimateDecay of
// the way to what it took, so a single short block does not hand the next decode
// every processor.
const executeEstimateDecay = 8

// prepareBudget returns how long PrepareBlock has to decode the next block given
// the typical block execution time. 0 when no block has executed yet, which
// decodes on every worker.
func prepareBudget(executeEstimate time.Duration) time.Duration {
return executeEstimate / prepareBudgetShare
}

// nextExecuteEstimate folds the execution time of a block into the estimate of the
// typical one.
func nextExecuteEstimate(current, executed time.Duration) time.Duration {
if current <= 0 {
return executed
}
return current + (executed-current)/executeEstimateDecay
}

type evmOnlyApplication struct {
abci.BaseApplication

Expand Down Expand Up @@ -89,6 +120,14 @@ type evmOnlyApplication struct {
prepared utils.Mutex[*utils.Option[preparedBlock]]
// preparedBlocks counts finalized blocks by whether prepared held them.
preparedBlocks otelmetric.Int64Counter
// preparePhases times PrepareBlock's decode of the next block. PrepareBlock is
// called from the single block fetcher, so one timer serves the app.
preparePhases *seidbmetrics.PhaseTimer
// executeEstimate is the typical time a prepared FinalizeBlock spends executing,
// in nanoseconds, averaged over the recent ones; PrepareBlock's decode budget is
// derived from it. Only prepared blocks contribute: an unprepared one includes
// its own decode.
executeEstimate atomic.Int64
}

// preparedBlock is the stateless part of a FinalizeBlock request, computed before the
Expand Down Expand Up @@ -152,6 +191,7 @@ func NewEVMOnlyApplication(
finalizePhases: seidbmetrics.NewPhaseTimer(otel.Meter(finalizeMeterName), "evmonly_finalize"),
prepared: utils.NewMutex(new(utils.Option[preparedBlock])),
preparedBlocks: newPreparedBlocksCounter(otel.Meter(finalizeMeterName)),
preparePhases: seidbmetrics.NewPhaseTimer(otel.Meter(finalizeMeterName), "evmonly_prepare"),
settler: utils.NewAtomicSend(utils.None[*evmonly.Executor]()),
cursor: utils.NewMutex(&evmOnlyCursorState{}),
checkedSenders: utils.NewMutex(map[common.Hash]common.Address{}),
Expand Down Expand Up @@ -648,7 +688,8 @@ func (a *evmOnlyApplication) PrepareBlock(ctx context.Context, req *abci.Request
}
// Only Number and Time reach the decoded transactions (through the signer); the
// parent-derived fields are filled in by FinalizeBlock.
prepared, err := executor.PrepareBlock(ctx, evmonly.BlockRequest{
a.preparePhases.SetPhase("parse")
prepared, err := executor.PrepareBlockWithin(ctx, evmonly.BlockRequest{
Context: evmonly.BlockContext{
Number: block.number,
Time: block.timestamp,
Expand All @@ -659,7 +700,8 @@ func (a *evmOnlyApplication) PrepareBlock(ctx context.Context, req *abci.Request
},
Txs: req.Txs,
Senders: a.peekSenders(req.Txs),
})
}, prepareBudget(time.Duration(a.executeEstimate.Load())))
a.preparePhases.Reset()
if err != nil {
return ctx.Err()
}
Expand Down Expand Up @@ -734,7 +776,9 @@ func (a *evmOnlyApplication) finalizeBlockLocked(
a.finalizePhases.SetPhase("take_senders")
a.forgetSenders(req.Txs)
a.finalizePhases.SetPhase("execute")
start := time.Now()
result, err = executor.ExecutePreparedBlock(ctx, evmonly.PreparedBlock{Context: blockCtx, Txs: prepared})
a.executeEstimate.Store(int64(nextExecuteEstimate(time.Duration(a.executeEstimate.Load()), time.Since(start))))
} else {
a.finalizePhases.SetPhase("take_senders")
senders := a.takeSenders(req.Txs)
Expand Down
17 changes: 17 additions & 0 deletions sei-tendermint/internal/evmonlyapp/app_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -841,6 +841,23 @@ func TestEVMOnlyApplicationReportsAnUndecodableBlockFromFinalizeBlock(t *testing
require.Error(t, unpreparedErr)
}

func TestPrepareBudgetIsAShareOfTheTypicalExecution(t *testing.T) {
require.Equal(t, time.Duration(0), prepareBudget(0))
require.Equal(t, 10*time.Millisecond, prepareBudget(20*time.Millisecond))
}

func TestExecuteEstimateSmoothsOverBlocks(t *testing.T) {
// The first block sets the estimate; later ones move it a step.
estimate := nextExecuteEstimate(0, 20*time.Millisecond)
require.Equal(t, 20*time.Millisecond, estimate)
// A single short block barely dents the budget rather than zeroing it.
estimate = nextExecuteEstimate(estimate, time.Millisecond)
require.Equal(t, 17625*time.Microsecond, estimate)
require.Equal(t, 8812500*time.Nanosecond, prepareBudget(estimate))
estimate = nextExecuteEstimate(estimate, 40*time.Millisecond)
require.Equal(t, 20421875*time.Nanosecond, estimate)
}

// Preparing before InitChain is a no-op rather than a failure.
func TestEVMOnlyApplicationPrepareBlockBeforeInitChainIsANoOp(t *testing.T) {
app := newEVMOnlyTestApp(t, nil)
Expand Down
Loading