From 453eb65acbc7f0e9a2053f14f1e69d15fb44114b Mon Sep 17 00:00:00 2001 From: Brandon Chatham Date: Sun, 20 Sep 2026 22:33:29 +0000 Subject: [PATCH] Revert "Revert "autobahn/evmonlyapp: prepare block n+1 while block n executes (#4260)"" This reverts commit 11a429e8b5aecdba6dc7cff558b9cdcf3b0ef5f0. --- sei-tendermint/internal/evmonlyapp/app.go | 197 +++++++++++++++--- .../internal/evmonlyapp/app_test.go | 158 ++++++++++++++ .../internal/p2p/giga_router_common.go | 145 ++++++++----- .../p2p/giga_router_testhelper_test.go | 34 +++ .../p2p/giga_router_validator_test.go | 1 + sei-tendermint/internal/proxy/proxy.go | 17 ++ 6 files changed, 472 insertions(+), 80 deletions(-) diff --git a/sei-tendermint/internal/evmonlyapp/app.go b/sei-tendermint/internal/evmonlyapp/app.go index f89d43d972..324b99793e 100644 --- a/sei-tendermint/internal/evmonlyapp/app.go +++ b/sei-tendermint/internal/evmonlyapp/app.go @@ -14,6 +14,9 @@ import ( "sync" "sync/atomic" + "go.opentelemetry.io/otel/attribute" + otelmetric "go.opentelemetry.io/otel/metric" + "github.com/ethereum/go-ethereum/common" ethcore "github.com/ethereum/go-ethereum/core" @@ -82,12 +85,38 @@ type evmOnlyApplication struct { // FinalizeBlock is serialized by executor, so one timer serves the app; it // is only touched with that lock held. finalizePhases *seidbmetrics.PhaseTimer + // prepared holds the block PrepareBlock decoded ahead of FinalizeBlock, if any. + 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 +} + +// preparedBlock is the stateless part of a FinalizeBlock request, computed before the +// request arrives. FinalizeBlock uses it only for the block with this height and hash. +type preparedBlock struct { + height int64 + hash common.Hash + txs []evmonly.PreparedTx } // finalizeMeterName is the OTel meter FinalizeBlock's phase timer records to, // as evmonly_finalize_phase_duration_seconds_total. const finalizeMeterName = "evmonly_app" +func newPreparedBlocksCounter(meter otelmetric.Meter) otelmetric.Int64Counter { + counter, err := meter.Int64Counter( + "evmonly_finalize_prepared_blocks_total", + otelmetric.WithDescription("Finalized blocks by whether PrepareBlock had already decoded them"), + ) + if err != nil { + panic(fmt.Sprintf("evmonly_finalize_prepared_blocks_total: %v", err)) + } + return counter +} + // evmOnlyCursorState is the execution position: the block whose state is // committed to storage and the block finalized but not yet acknowledged by // Commit. @@ -124,6 +153,9 @@ func NewEVMOnlyApplication( validators: slices.Clone(validators), executor: utils.NewMutex(new(utils.Option[*evmonly.Executor])), 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{}), @@ -351,6 +383,20 @@ func (a *evmOnlyApplication) rememberSender(hash common.Hash, sender common.Addr // raw transaction is the keccak of its bytes for every transaction type, so no // decoding is needed. func (a *evmOnlyApplication) takeSenders(txs [][]byte) []utils.Option[common.Address] { + return a.checkedSendersOf(txs, true) +} + +// peekSenders is takeSenders without forgetting the entries. +func (a *evmOnlyApplication) peekSenders(txs [][]byte) []utils.Option[common.Address] { + return a.checkedSendersOf(txs, false) +} + +// forgetSenders drops the CheckTx-recovered senders of txs. +func (a *evmOnlyApplication) forgetSenders(txs [][]byte) { + a.checkedSendersOf(txs, true) +} + +func (a *evmOnlyApplication) checkedSendersOf(txs [][]byte, forget bool) []utils.Option[common.Address] { out := make([]utils.Option[common.Address], len(txs)) // Hashed outside the lock; CheckTx writes this map constantly. hashes := hashRawTxs(txs) @@ -358,7 +404,9 @@ func (a *evmOnlyApplication) takeSenders(txs [][]byte) []utils.Option[common.Add for i, hash := range hashes { if sender, ok := senders[hash]; ok { out[i] = utils.Some(sender) - delete(senders, hash) + if forget { + delete(senders, hash) + } } } } @@ -556,26 +604,106 @@ func (a *evmOnlyApplication) EvmCall(ctx context.Context, msg *ethcore.Message) } } -func (a *evmOnlyApplication) FinalizeBlock(ctx context.Context, req *abci.RequestFinalizeBlock) (*abci.ResponseFinalizeBlock, error) { +// finalizeRequest is the block identity FinalizeBlock and PrepareBlock derive from a request. +type finalizeRequest struct { + height int64 + number uint64 + timestamp uint64 + blockHash common.Hash +} + +func parseFinalizeRequest(req *abci.RequestFinalizeBlock) (finalizeRequest, error) { height := req.Header.Height if height <= 0 { - return nil, fmt.Errorf("EVM-only block height must be positive: %d", height) + return finalizeRequest{}, fmt.Errorf("EVM-only block height must be positive: %d", height) } number, ok := utils.SafeCast[uint64](height) if !ok { - return nil, fmt.Errorf("EVM-only block height exceeds uint64: %d", height) + return finalizeRequest{}, fmt.Errorf("EVM-only block height exceeds uint64: %d", height) } timestamp, ok := utils.SafeCast[uint64](req.Header.Time.Unix()) if !ok { - return nil, fmt.Errorf("EVM-only block timestamp is negative: %s", req.Header.Time) + return finalizeRequest{}, fmt.Errorf("EVM-only block timestamp is negative: %s", req.Header.Time) + } + return finalizeRequest{ + height: height, + number: number, + timestamp: timestamp, + blockHash: common.BytesToHash(req.Hash), + }, nil +} + +// PrepareBlock decodes a block's transactions and recovers their senders before +// FinalizeBlock is called for it, so that work runs while the previous block +// executes. It may run concurrently with FinalizeBlock. Only the most recent +// prepared block is kept, and FinalizeBlock uses it only for the same height and +// hash, so preparing the wrong block costs nothing but the work: the senders +// CheckTx cached stay cached until a prepared block is consumed. Anything that +// would fail the block is left for FinalizeBlock to report; the only error +// returned is ctx ending. +func (a *evmOnlyApplication) PrepareBlock(ctx context.Context, req *abci.RequestFinalizeBlock) error { + executor, ok := a.settler.Load().Get() + if !ok { + return nil + } + block, err := parseFinalizeRequest(req) + if err != nil { + return nil + } + // Only Number and Time reach the decoded transactions (through the signer); the + // parent-derived fields are filled in by FinalizeBlock. + a.preparePhases.SetPhase("parse") + prepared, err := executor.PrepareBlock(ctx, evmonly.BlockRequest{ + Context: evmonly.BlockContext{ + Number: block.number, + Time: block.timestamp, + GasLimit: a.EvmGasLimit(), + ChainID: new(big.Int).Set(a.chainID), + BaseFee: evmOnlyBaseFee(), + BlobBaseFee: new(big.Int), + }, + Txs: req.Txs, + Senders: a.peekSenders(req.Txs), + }) + a.preparePhases.Reset() + if err != nil { + return ctx.Err() + } + for slot := range a.prepared.Lock() { + *slot = utils.Some(preparedBlock{ + height: block.height, + hash: block.blockHash, + txs: prepared.Txs, + }) + } + return nil +} + +// takePrepared returns the prepared transactions of the given block and removes +// them. A prepared block for another block is left in place. +func (a *evmOnlyApplication) takePrepared(height int64, hash common.Hash) ([]evmonly.PreparedTx, bool) { + for slot := range a.prepared.Lock() { + prepared, ok := slot.Get() + if !ok || prepared.height != height || prepared.hash != hash { + return nil, false + } + *slot = utils.None[preparedBlock]() + return prepared.txs, true + } + panic("unreachable") +} + +func (a *evmOnlyApplication) FinalizeBlock(ctx context.Context, req *abci.RequestFinalizeBlock) (*abci.ResponseFinalizeBlock, error) { + block, err := parseFinalizeRequest(req) + if err != nil { + return nil, err } - blockHash := common.BytesToHash(req.Hash) for executor := range a.executor.Lock() { executor, ok := executor.Get() if !ok { return nil, fmt.Errorf("EVM-only block finalized before InitChain") } - return a.finalizeBlockLocked(ctx, executor, req, number, timestamp, blockHash) + return a.finalizeBlockLocked(ctx, executor, req, block) } panic("unreachable") } @@ -586,38 +714,47 @@ func (a *evmOnlyApplication) finalizeBlockLocked( ctx context.Context, executor *evmonly.Executor, req *abci.RequestFinalizeBlock, - number, timestamp uint64, - blockHash common.Hash, + block finalizeRequest, ) (*abci.ResponseFinalizeBlock, error) { - height := req.Header.Height - parent, err := a.beginBlock(height) + parent, err := a.beginBlock(block.height) if err != nil { return nil, err } + blockCtx := evmonly.BlockContext{ + Number: block.number, + Time: block.timestamp, + GasLimit: parent.gasLimit, + ChainID: new(big.Int).Set(a.chainID), + BaseFee: evmOnlyBaseFee(), + BlobBaseFee: new(big.Int), + ParentHash: parent.blockHash, + BlockHash: block.blockHash, + PrevRandao: parent.appHash, + } // Closes the stage in flight, so the gap until the next block is charged to neither. defer a.finalizePhases.Reset() - a.finalizePhases.SetPhase("take_senders") - senders := a.takeSenders(req.Txs) - result, err := executeBlockPipelined(ctx, executor, a.finalizePhases, evmonly.BlockRequest{ - Context: evmonly.BlockContext{ - Number: number, - Time: timestamp, - GasLimit: parent.gasLimit, - ChainID: new(big.Int).Set(a.chainID), - BaseFee: evmOnlyBaseFee(), - BlobBaseFee: new(big.Int), - ParentHash: parent.blockHash, - BlockHash: blockHash, - PrevRandao: parent.appHash, - }, - Txs: req.Txs, - Senders: senders, - }) + prepared, hit := a.takePrepared(block.height, block.blockHash) + a.preparedBlocks.Add(ctx, 1, otelmetric.WithAttributes(attribute.Bool("prepared", hit))) + var result *evmonly.BlockResult + if hit { + a.finalizePhases.SetPhase("take_senders") + a.forgetSenders(req.Txs) + a.finalizePhases.SetPhase("execute") + result, err = executor.ExecutePreparedBlock(ctx, evmonly.PreparedBlock{Context: blockCtx, Txs: prepared}) + } else { + a.finalizePhases.SetPhase("take_senders") + senders := a.takeSenders(req.Txs) + result, err = executeBlockPipelined(ctx, executor, a.finalizePhases, evmonly.BlockRequest{ + Context: blockCtx, + Txs: req.Txs, + Senders: senders, + }) + } if err != nil { - return nil, errors.Join(err, a.abandonPending(executor, height)) + return nil, errors.Join(err, a.abandonPending(executor, block.height)) } defer result.Release() - pending, err := a.pendingCursor(height) + pending, err := a.pendingCursor(block.height) if err != nil { return nil, err } diff --git a/sei-tendermint/internal/evmonlyapp/app_test.go b/sei-tendermint/internal/evmonlyapp/app_test.go index 74b9aeef2c..4243c9dbbb 100644 --- a/sei-tendermint/internal/evmonlyapp/app_test.go +++ b/sei-tendermint/internal/evmonlyapp/app_test.go @@ -25,6 +25,7 @@ import ( "github.com/sei-protocol/sei-chain/sei-db/proto" abci "github.com/sei-protocol/sei-chain/sei-tendermint/abci/types" "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils/require" + "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils/scope" tmproto "github.com/sei-protocol/sei-chain/sei-tendermint/proto/tendermint/types" ) @@ -738,3 +739,160 @@ func TestEVMOnlyApplicationTimesEveryFinalizeBlockPhase(t *testing.T) { require.True(t, ok, "phase %q not recorded", want) } } + +// preparedBlockCounts reads evmonly_finalize_prepared_blocks_total by its prepared label. +func preparedBlockCounts(t *testing.T, reader *sdkmetric.ManualReader) map[bool]int64 { + t.Helper() + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(t.Context(), &rm)) + counts := map[bool]int64{} + for _, scope := range rm.ScopeMetrics { + for _, m := range scope.Metrics { + if m.Name != "evmonly_finalize_prepared_blocks_total" { + continue + } + sum, ok := m.Data.(metricdata.Sum[int64]) + require.True(t, ok) + for _, point := range sum.DataPoints { + prepared, ok := point.Attributes.Value("prepared") + require.True(t, ok) + counts[prepared.AsBool()] = point.Value + } + } + } + return counts +} + +// A block decoded ahead of FinalizeBlock has to execute to the same result as one +// decoded inside it, and a block nobody prepared still has to execute. +func TestEVMOnlyApplicationExecutesAPreparedBlockLikeAnUnpreparedOne(t *testing.T) { + prepared := newInitializedEVMOnlyTestApp(t) + preparedApp, ok := prepared.(*evmOnlyApplication) + require.True(t, ok) + reader := sdkmetric.NewManualReader() + preparedApp.preparedBlocks = newPreparedBlocksCounter( + sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)).Meter(finalizeMeterName), + ) + unprepared := newInitializedEVMOnlyTestApp(t) + + first, _ := signedEVMOnlyTestTx(t, evmOnlyTestChainID, 0) + second, _ := signedEVMOnlyTestTx(t, evmOnlyTestChainID, 0) + for height, txs := range [][][]byte{{first}, {second}} { + block := evmOnlyTestBlock(int64(height)+1, txs...) + if height == 0 { + require.NoError(t, preparedApp.PrepareBlock(t.Context(), block)) + } + preparedResp, err := prepared.FinalizeBlock(t.Context(), block) + require.NoError(t, err) + unpreparedResp, err := unprepared.FinalizeBlock(t.Context(), block) + require.NoError(t, err) + require.Equal(t, unpreparedResp.AppHash, preparedResp.AppHash) + require.Equal(t, len(unpreparedResp.TxResults), len(preparedResp.TxResults)) + for i := range preparedResp.TxResults { + require.Equal(t, unpreparedResp.TxResults[i].Code, preparedResp.TxResults[i].Code) + require.Equal(t, unpreparedResp.TxResults[i].GasUsed, preparedResp.TxResults[i].GasUsed) + } + _, err = prepared.Commit(t.Context()) + require.NoError(t, err) + _, err = unprepared.Commit(t.Context()) + require.NoError(t, err) + } + counts := preparedBlockCounts(t, reader) + require.Equal(t, int64(1), counts[true]) + require.Equal(t, int64(1), counts[false]) +} + +// A prepared block is only used for the block it was prepared for: one with the +// same height but another hash is decoded again, and the prepared one is kept +// for its own block. +func TestEVMOnlyApplicationIgnoresAPreparedBlockForAnotherHash(t *testing.T) { + prepared := newInitializedEVMOnlyTestApp(t) + preparedApp, ok := prepared.(*evmOnlyApplication) + require.True(t, ok) + unprepared := newInitializedEVMOnlyTestApp(t) + + other, _ := signedEVMOnlyTestTx(t, evmOnlyTestChainID, 0) + otherBlock := evmOnlyTestBlock(1, other) + otherBlock.Hash = crypto.Keccak256([]byte("other-block")) + require.NoError(t, preparedApp.PrepareBlock(t.Context(), otherBlock)) + + raw, _ := signedEVMOnlyTestTx(t, evmOnlyTestChainID, 0) + block := evmOnlyTestBlock(1, raw) + preparedHash := finalizeAndCommitEVMOnlyTestBlock(t, prepared, block) + unpreparedHash := finalizeAndCommitEVMOnlyTestBlock(t, unprepared, block) + require.Equal(t, unpreparedHash, preparedHash) + _, ok = preparedApp.takePrepared(1, common.BytesToHash(otherBlock.Hash)) + require.True(t, ok) +} + +// A block PrepareBlock cannot decode is reported by FinalizeBlock, the same as +// when nobody prepared it. +func TestEVMOnlyApplicationReportsAnUndecodableBlockFromFinalizeBlock(t *testing.T) { + prepared := newInitializedEVMOnlyTestApp(t) + preparedApp, ok := prepared.(*evmOnlyApplication) + require.True(t, ok) + unprepared := newInitializedEVMOnlyTestApp(t) + + block := evmOnlyTestBlock(1, []byte("not a transaction")) + require.NoError(t, preparedApp.PrepareBlock(t.Context(), block)) + _, preparedErr := prepared.FinalizeBlock(t.Context(), block) + require.Error(t, preparedErr) + _, unpreparedErr := unprepared.FinalizeBlock(t.Context(), block) + require.Error(t, unpreparedErr) +} + +// Preparing before InitChain is a no-op rather than a failure. +func TestEVMOnlyApplicationPrepareBlockBeforeInitChainIsANoOp(t *testing.T) { + app := newEVMOnlyTestApp(t, nil) + evmOnlyApp, ok := app.(*evmOnlyApplication) + require.True(t, ok) + raw, _ := signedEVMOnlyTestTx(t, evmOnlyTestChainID, 0) + require.NoError(t, evmOnlyApp.PrepareBlock(t.Context(), evmOnlyTestBlock(1, raw))) +} + +// PrepareBlock for the next block runs while FinalizeBlock runs the current one, +// and both blocks come out as if they had been finalized alone. +func TestEVMOnlyApplicationPreparesTheNextBlockWhileFinalizingTheCurrentOne(t *testing.T) { + prepared := newInitializedEVMOnlyTestApp(t) + preparedApp, ok := prepared.(*evmOnlyApplication) + require.True(t, ok) + unprepared := newInitializedEVMOnlyTestApp(t) + + key, err := crypto.GenerateKey() + require.NoError(t, err) + var blocks []*abci.RequestFinalizeBlock + for height := range int64(4) { + var txs [][]byte + for i := range 8 { + txs = append(txs, signedEVMOnlyTestTxFrom(t, key, evmOnlyTestChainID, uint64(height)*8+uint64(i))) //nolint:gosec // G115: small test counters. + } + blocks = append(blocks, evmOnlyTestBlock(height+1, txs...)) + } + + hashes := make([][]byte, len(blocks)) + require.NoError(t, scope.Run(t.Context(), func(ctx context.Context, s scope.Scope) error { + for i, block := range blocks { + var next *abci.RequestFinalizeBlock + if i+1 < len(blocks) { + next = blocks[i+1] + } + resp, err := scope.Run1(ctx, func(ctx context.Context, s scope.Scope) (*abci.ResponseFinalizeBlock, error) { + if next != nil { + s.Spawn(func() error { return preparedApp.PrepareBlock(ctx, next) }) + } + return prepared.FinalizeBlock(ctx, block) + }) + if err != nil { + return err + } + if _, err := prepared.Commit(ctx); err != nil { + return err + } + hashes[i] = resp.AppHash + } + return nil + })) + for i, block := range blocks { + require.Equal(t, finalizeAndCommitEVMOnlyTestBlock(t, unprepared, block), hashes[i]) + } +} diff --git a/sei-tendermint/internal/p2p/giga_router_common.go b/sei-tendermint/internal/p2p/giga_router_common.go index f71ba4cf31..b263c5c63e 100644 --- a/sei-tendermint/internal/p2p/giga_router_common.go +++ b/sei-tendermint/internal/p2p/giga_router_common.go @@ -224,22 +224,11 @@ func (r *gigaRouterCommon) translateGlobalBlock(gb *atypes.GlobalBlock) *coretyp } } -func (r *gigaRouterCommon) executeBlock(ctx context.Context, b *atypes.GlobalBlock, hashVault hashvault.HashVault) (*abci.ResponseCommit, error) { - app := r.app +// finalizeRequest builds the FinalizeBlock request for a global block, except +// for the header's ProposerAddress, which depends on the committed state. +func (r *gigaRouterCommon) finalizeRequest(b *atypes.GlobalBlock) *abci.RequestFinalizeBlock { hash := b.Header.Hash() - var proposerAddress types.Address - if vals := app.GetValidators(); len(vals) > 0 { - // Deterministically select a proposer from the app's validator committee. - // We need it so that app does not emit error logs. - proposer := slices.MinFunc(vals, func(a, b abci.ValidatorUpdate) int { return a.PubKey.Compare(b.PubKey) }) - key, err := crypto.PubKeyFromProto(proposer.PubKey) - if err != nil { - return nil, fmt.Errorf("crypto.PubKeyFromProto(): %w", err) - } - proposerAddress = key.Address() - } - - resp, err := app.FinalizeBlock(ctx, &abci.RequestFinalizeBlock{ + return &abci.RequestFinalizeBlock{ Txs: b.Payload.Txs(), // Empty DecidedLastCommit does not indicate missing votes. DecidedLastCommit: abci.CommitInfo{}, @@ -251,11 +240,45 @@ func (r *gigaRouterCommon) executeBlock(ctx context.Context, b *atypes.GlobalBlo ChainID: r.cfg.GenDoc.ChainID, Height: int64(b.GlobalNumber), // nolint:gosec // different representations of the same value Time: b.Timestamp, - // WARNING: the reward distribution has corner cases where it forgets the proposer, - // because reward is distributed with a delay. This is not our problem here though. - ProposerAddress: proposerAddress, }).ToProto(), - }) + } +} + +// proposerAddress returns the proposer of the next block, selected from the +// app's current validator committee. +func (r *gigaRouterCommon) proposerAddress() (types.Address, error) { + vals := r.app.GetValidators() + if len(vals) == 0 { + return nil, nil + } + // Deterministically select a proposer from the app's validator committee. + // We need it so that app does not emit error logs. + proposer := slices.MinFunc(vals, func(a, b abci.ValidatorUpdate) int { return a.PubKey.Compare(b.PubKey) }) + key, err := crypto.PubKeyFromProto(proposer.PubKey) + if err != nil { + return nil, fmt.Errorf("crypto.PubKeyFromProto(): %w", err) + } + return key.Address(), nil +} + +// fetchedBlock is a global block with the FinalizeBlock request built for it. +type fetchedBlock struct { + block *atypes.GlobalBlock + req *abci.RequestFinalizeBlock +} + +func (r *gigaRouterCommon) executeBlock(ctx context.Context, f fetchedBlock, hashVault hashvault.HashVault) (*abci.ResponseCommit, error) { + app := r.app + b := f.block + // Read from the app's state on the execute loop, after the previous block's Commit. + // WARNING: the reward distribution has corner cases where it forgets the proposer, + // because reward is distributed with a delay. This is not our problem here though. + proposer, err := r.proposerAddress() + if err != nil { + return nil, err + } + f.req.Header.ProposerAddress = proposer + resp, err := app.FinalizeBlock(ctx, f.req) if err != nil { return nil, fmt.Errorf("app.FinalizeBlock(): %w", err) } @@ -475,39 +498,61 @@ func (r *gigaRouterCommon) runExecute(ctx context.Context) error { } } - for n := next; ; n += 1 { - gigametrics.SetPhase(gigametrics.PhaseConsensus) - b, err := r.data.GlobalBlock(ctx, n) - if err != nil { - return fmt.Errorf("r.data.GlobalBlock(%v): %w", n, err) - } - gigametrics.SetPhase(gigametrics.PhaseExecution) - commitResp, err := r.executeBlock(ctx, b, hashVault) - if err != nil { - return fmt.Errorf("r.executeBlock(%v): %w", n, err) - } - pruneBefore, ok := utils.SafeCast[atypes.GlobalBlockNumber](commitResp.RetainHeight) - if !ok { - return fmt.Errorf("invalid commitResp.RetainHeight = %v", commitResp.RetainHeight) - } - gigametrics.SetStoragePhase(gigametrics.StoragePhasePruneData) - if err := r.data.PruneBefore(pruneBefore); err != nil { - return fmt.Errorf("r.data.PruneBefore(%v): %w", pruneBefore, err) - } - // Align the vault's retention with the data layer's prune boundary. - gigametrics.SetStoragePhase(gigametrics.StoragePhasePruneVault) - if err := hashVault.Prune(ctx, uint64(pruneBefore)); err != nil { - // A canceled context just means we're shutting down between a successful executeBlock - // and this prune; that's benign, not a prune failure, so don't alarm operators. - if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { - logger.Info("hashvault prune aborted by context cancellation during shutdown", - "prune_before", pruneBefore, "err", err) - } else { - logger.Error("failed to prune hashvault", "prune_before", pruneBefore, "err", err) + return scope.Run(ctx, func(ctx context.Context, s scope.Scope) error { + // Blocks are fetched and prepared one ahead of execution: the unbuffered + // channel lets the fetcher hold block n+1, already prepared, while the + // loop below executes block n. + blocks := make(chan fetchedBlock) + s.Spawn(func() error { + for n := next; ; n += 1 { + b, err := r.data.GlobalBlock(ctx, n) + if err != nil { + return fmt.Errorf("r.data.GlobalBlock(%v): %w", n, err) + } + req := r.finalizeRequest(b) + if err := app.PrepareBlock(ctx, req); err != nil { + return fmt.Errorf("app.PrepareBlock(%v): %w", n, err) + } + if err := utils.Send(ctx, blocks, fetchedBlock{block: b, req: req}); err != nil { + return err + } + } + }) + for { + gigametrics.SetPhase(gigametrics.PhaseConsensus) + f, err := utils.Recv(ctx, blocks) + if err != nil { + return err + } + n := f.block.GlobalNumber + gigametrics.SetPhase(gigametrics.PhaseExecution) + commitResp, err := r.executeBlock(ctx, f, hashVault) + if err != nil { + return fmt.Errorf("r.executeBlock(%v): %w", n, err) + } + pruneBefore, ok := utils.SafeCast[atypes.GlobalBlockNumber](commitResp.RetainHeight) + if !ok { + return fmt.Errorf("invalid commitResp.RetainHeight = %v", commitResp.RetainHeight) } + gigametrics.SetStoragePhase(gigametrics.StoragePhasePruneData) + if err := r.data.PruneBefore(pruneBefore); err != nil { + return fmt.Errorf("r.data.PruneBefore(%v): %w", pruneBefore, err) + } + // Align the vault's retention with the data layer's prune boundary. + gigametrics.SetStoragePhase(gigametrics.StoragePhasePruneVault) + if err := hashVault.Prune(ctx, uint64(pruneBefore)); err != nil { + // A canceled context just means we're shutting down between a successful executeBlock + // and this prune; that's benign, not a prune failure, so don't alarm operators. + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + logger.Info("hashvault prune aborted by context cancellation during shutdown", + "prune_before", pruneBefore, "err", err) + } else { + logger.Error("failed to prune hashvault", "prune_before", pruneBefore, "err", err) + } + } + gigametrics.EndStoragePhase() } - gigametrics.EndStoragePhase() - } + }) } // dialAndRunConn dials a peer, handshakes as a SeiGiga connection, diff --git a/sei-tendermint/internal/p2p/giga_router_testhelper_test.go b/sei-tendermint/internal/p2p/giga_router_testhelper_test.go index b5bc6e3533..731f8d6093 100644 --- a/sei-tendermint/internal/p2p/giga_router_testhelper_test.go +++ b/sei-tendermint/internal/p2p/giga_router_testhelper_test.go @@ -1,6 +1,7 @@ package p2p import ( + "bytes" "context" "crypto/sha256" "encoding/json" @@ -152,6 +153,39 @@ func (a *testApp) FinalizeBlock(_ context.Context, req *abci.RequestFinalizeBloc panic("unreachable") } +// CheckBlocks verifies the FinalizeBlock requests the router delivered: heights +// are contiguous from InitialHeight and every header names a current validator +// as proposer. +func (s *testAppState) CheckBlocks() error { + init, ok := s.Init.Get() + if !ok { + return fmt.Errorf("app not initialized") + } + for i, b := range s.Blocks { + if want := init.InitialHeight + int64(i); b.Header.Height != want { + return fmt.Errorf("blocks[%v].Height = %v, want %v", i, b.Header.Height, want) + } + if err := checkProposer(b.Header.ProposerAddress, s.Validators); err != nil { + return fmt.Errorf("blocks[%v]: %w", i, err) + } + } + return nil +} + +// checkProposer verifies that the proposer is one of the given validators. +func checkProposer(proposer types.Address, vals []abci.ValidatorUpdate) error { + for _, val := range vals { + key, err := crypto.PubKeyFromProto(val.PubKey) + if err != nil { + return fmt.Errorf("crypto.PubKeyFromProto(): %w", err) + } + if bytes.Equal(key.Address(), proposer) { + return nil + } + } + return fmt.Errorf("proposer %X is not a current validator", proposer) +} + func (a *testApp) Commit(context.Context) (*abci.ResponseCommit, error) { for state, ctrl := range a.state.Lock() { if state.Committed { diff --git a/sei-tendermint/internal/p2p/giga_router_validator_test.go b/sei-tendermint/internal/p2p/giga_router_validator_test.go index aa216a0c84..cf3c71cdf4 100644 --- a/sei-tendermint/internal/p2p/giga_router_validator_test.go +++ b/sei-tendermint/internal/p2p/giga_router_validator_test.go @@ -149,6 +149,7 @@ func TestGigaRouter_FinalizeBlocks(t *testing.T) { } // Nodes should agree on the final state. want := apps[0].Snapshot() + require.NoError(t, want.CheckBlocks(), "CheckBlocks") for i, app := range apps { t.Logf("app[%v]", i) require.NoError(t, utils.TestDiff(want, app.Snapshot()), "state mismatch app[%v]", i) diff --git a/sei-tendermint/internal/proxy/proxy.go b/sei-tendermint/internal/proxy/proxy.go index 1fa1cb73bc..c8ba5b68a1 100644 --- a/sei-tendermint/internal/proxy/proxy.go +++ b/sei-tendermint/internal/proxy/proxy.go @@ -80,6 +80,23 @@ func (app *Proxy) AwaitCommits() error { return settler.AwaitCommits() } +// blockPreparer is implemented by applications that can do the stateless part of +// FinalizeBlock for a block before it is called. +type blockPreparer interface { + PrepareBlock(context.Context, *types.RequestFinalizeBlock) error +} + +// PrepareBlock lets the wrapped application decode a block ahead of its +// FinalizeBlock. It is a no-op for an application that does not prepare blocks. +func (app *Proxy) PrepareBlock(ctx context.Context, req *types.RequestFinalizeBlock) error { + defer addTimeSample(Global.MethodTimingAt("prepare_block", "sync"))() + preparer, ok := app.app.(blockPreparer) + if !ok { + return nil + } + return preparer.PrepareBlock(ctx, req) +} + // evmCaller is implemented by applications that can run a read-only EVM call // against their current state. type evmCaller interface {