diff --git a/giga/evmonly/executor.go b/giga/evmonly/executor.go index ca03f6fb1f..5db07cbae0 100644 --- a/giga/evmonly/executor.go +++ b/giga/evmonly/executor.go @@ -6,6 +6,7 @@ import ( "math/big" "sync" "sync/atomic" + "time" "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/core" @@ -29,6 +30,7 @@ type Executor struct { cfg Config resultSink ResultSink occPool *occWorkerPool + parseSizer *parseSizer resultPool *blockResultPool stateDBPool sync.Pool storeMu sync.Mutex @@ -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) } @@ -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 @@ -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, diff --git a/giga/evmonly/parse_sizer.go b/giga/evmonly/parse_sizer.go new file mode 100644 index 0000000000..12d324d7b6 --- /dev/null +++ b/giga/evmonly/parse_sizer.go @@ -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 + } + } +} diff --git a/giga/evmonly/parse_sizer_test.go b/giga/evmonly/parse_sizer_test.go new file mode 100644 index 0000000000..cbf0139370 --- /dev/null +++ b/giga/evmonly/parse_sizer_test.go @@ -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()) +} diff --git a/sei-tendermint/internal/evmonlyapp/app.go b/sei-tendermint/internal/evmonlyapp/app.go index 9de73b706e..b7e7d43245 100644 --- a/sei-tendermint/internal/evmonlyapp/app.go +++ b/sei-tendermint/internal/evmonlyapp/app.go @@ -13,6 +13,7 @@ import ( "slices" "sync" "sync/atomic" + "time" "go.opentelemetry.io/otel/attribute" otelmetric "go.opentelemetry.io/otel/metric" @@ -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 @@ -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 @@ -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{}), @@ -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, @@ -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() } @@ -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) diff --git a/sei-tendermint/internal/evmonlyapp/app_test.go b/sei-tendermint/internal/evmonlyapp/app_test.go index 7eff0fa9ed..825eae2c65 100644 --- a/sei-tendermint/internal/evmonlyapp/app_test.go +++ b/sei-tendermint/internal/evmonlyapp/app_test.go @@ -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)