diff --git a/giga/evmonly/README.md b/giga/evmonly/README.md index d471467f33..36c6f0c70e 100644 --- a/giga/evmonly/README.md +++ b/giga/evmonly/README.md @@ -78,34 +78,31 @@ receipt store, plus the `NamedChangeSetEncoder` for its state implementation. Unit tests can supply those dependencies independently. Execution fails closed if the state store or the encoder is missing. The receipt store is optional: a node configured without one (`enable_receipt_store = false`, meant for -validators that serve no receipt reads) skips receipt persistence entirely and -`AwaitReceipts` returns at once. +validators that serve no receipt reads) skips receipt persistence entirely. For each block the executor opens a current `giga.StateView`, executes against its EVM-native read methods, converts the resulting `StateChangeSet`, and calls `CommitStateChanges`. Execution and commit on an executor are serialized so blocks cannot share a stale snapshot or overlap commits; callers must still -submit block heights in order. The snapshot is closed once execution returns; -the commit reads from the changeset, not the snapshot. An empty block still -commits an encoded empty changeset so the store can advance its height. -Stateless preparation can continue concurrently with store-backed execution. +submit block heights in order. The snapshot stays open through the commit and +is always closed afterward. An empty block still commits an encoded empty +changeset so the store can advance its height. Stateless preparation can +continue concurrently with store-backed execution. The encoder is explicit because `giga.StateDB` defines the protobuf commit transport but does not define an on-disk key layout. In particular, an encoder must preserve `StorageClears` as prefix clears rather than silently dropping -persisted slots that were not read during execution. Encoding and both store -writes run behind the block, one block at a time in block order; `ResultSink` -runs on the block loop once that work has been handed off, so it may see a -result whose writes have not landed yet. Ethereum receipts are converted into +persisted slots that were not read during execution. Encoding, state commit, or +receipt-store failures release the block result and return an error without +invoking `ResultSink`. Ethereum receipts are converted into `receipt.ReceiptRecord` values and persisted through the shared -`receipt.ReceiptStore` interface, including for empty blocks. That write starts -as soon as execution returns, before the block encoder runs and the previous -commit is waited on, and the block's height-advancing state commit waits for it -to land; `AwaitReceipts` blocks until the newest block's receipts are readable. -A failure in either write latches: `AwaitReceipts` reports a receipt failure, -the next block's execution reports a state failure, and the executor accepts -no further blocks. A receipt failure leaves state unchanged; a state failure, -or a block that fails after its receipts were handed off, can leave receipts -behind, which re-executing the block after a restart overwrites. +`receipt.ReceiptStore` interface before the height-advancing state commit, +including for empty blocks. A store with an async write queue only accepts the +write here: the receipts land behind the block, so a reader that follows the +state head can briefly miss the newest block's receipts, and recovery replays +the tail a crash leaves unwritten. A receipt failure leaves state unchanged so +the block can be retried. A state failure can leave receipts behind, but +retrying the block overwrites them. `ResultSink` runs only after both stores +accept the block. `ExecuteBlock` advances the state store's version itself, independently of any ABCI `Commit`, so what the store holds after a restart is decided by the diff --git a/giga/evmonly/executor.go b/giga/evmonly/executor.go index d47c0e34f1..ca03f6fb1f 100644 --- a/giga/evmonly/executor.go +++ b/giga/evmonly/executor.go @@ -19,7 +19,6 @@ import ( "github.com/sei-protocol/sei-chain/sei-db/ledger_db/receipt" gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" "go.opentelemetry.io/otel" - "go.opentelemetry.io/otel/attribute" ) // executorMeterName is the OTel meter this package's instruments are created on. @@ -47,24 +46,13 @@ type Executor struct { // Breaks a store-backed block into its stages. That path is serialized by storeMu, so one timer // serves the executor. blockPhases *seidbmetrics.PhaseTimer - // Breaks the background persistence of a block's state into its stages. One block is - // committed at a time, so one timer serves it. - pipelinePhases *seidbmetrics.PhaseTimer - // Breaks the background receipt write into its stages. Receipt writes run one at a time, in - // block order, but overlap the state commit, so they have a timer of their own on the same - // metric, told apart by a stage label. - receiptPhases *seidbmetrics.PhaseTimer // The commit running behind the current block, and what it will write. A block reads the latter // through an overlay so it need not wait for the former. - pipelineMu sync.Mutex - // Closed once the block is fully persisted. - pipelineDone chan struct{} - pipelineErr error - // The most recent block's receipt write; kept after the block retires so a waiter that arrives - // late still finds its answer. - pipelineReceipts *receiptWrite - pipelineChanges *pendingChanges + pipelineMu sync.Mutex + pipelineDone chan struct{} + pipelineErr error + pipelineChanges *pendingChanges // Counts commits started, so a reader can tell that a block landed between two of its steps. pipelineGeneration uint64 // The first commit that failed, kept so no caller can miss it. @@ -109,13 +97,10 @@ func WithBlockChangeSetEncoder(encoder BlockChangeSetEncoder) Option { // NewExecutor constructs an EVM-only executor. Call Close to disable future OCC // execution on this executor. func NewExecutor(cfg Config, opts ...Option) *Executor { - pipelineTimers := seidbmetrics.NewPhaseTimerFactory(otel.Meter(executorMeterName), "evmonly_pipeline") e := &Executor{ - cfg: cfg.WithDefaults(), - resultPool: newBlockResultPool(cfg.BlockResultPoolSize), - blockPhases: seidbmetrics.NewPhaseTimer(otel.Meter(executorMeterName), "evmonly_block"), - pipelinePhases: pipelineTimers.Build(attribute.String("stage", "state")), - receiptPhases: pipelineTimers.Build(attribute.String("stage", "receipts")), + cfg: cfg.WithDefaults(), + resultPool: newBlockResultPool(cfg.BlockResultPoolSize), + blockPhases: seidbmetrics.NewPhaseTimer(otel.Meter(executorMeterName), "evmonly_block"), } if e.cfg.OCCWorkers > 1 { e.occPool = newOCCWorkerPool(e.cfg.OCCWorkers) @@ -132,10 +117,8 @@ func (e *Executor) Close() { } e.closed.Store(true) // Land the commit running behind the last block before the pool it may need goes away. The - // failure is kept rather than reported, for the next AwaitCommits to return. The receipt write - // is waited on separately: a block that failed after starting it has no commit to land. + // failure is kept rather than reported, for the next AwaitCommits to return. _ = e.awaitPipelineCommit() - _ = e.AwaitReceipts() if e.occPool != nil { e.occPool.Close() } diff --git a/giga/evmonly/executor_test.go b/giga/evmonly/executor_test.go index 4d682dcba6..87aaac9178 100644 --- a/giga/evmonly/executor_test.go +++ b/giga/evmonly/executor_test.go @@ -174,21 +174,19 @@ func TestExecutorReturnsReceiptStoreError(t *testing.T) { require.ErrorIs(t, err, storeErr) require.Nil(t, result) - // Receipts are written behind the block, so the sink has already seen the result. - require.Len(t, sink.results, 1) - sink.releases[0]() + require.Empty(t, sink.results) require.Equal(t, BlockResultPoolStats{Capacity: 1, Available: 1}, executor.ResultPoolStats()) view := stateStore.OpenView() require.Zero(t, view.GetBlockHeight()) view.Close() - // A failed write is latched: the executor refuses further blocks rather than run ahead of a - // store that is missing a block. receiptStore.err = nil result, err = executor.ExecuteBlock(t.Context(), request) - require.ErrorIs(t, err, storeErr) - require.Nil(t, result) + require.NoError(t, err) + require.NotNil(t, result) require.Len(t, sink.results, 1) + result.Release() + sink.releases[0]() } func TestExecutorPooledResultRelease(t *testing.T) { diff --git a/giga/evmonly/giga_store.go b/giga/evmonly/giga_store.go index 44b81887c9..5f6dfe44ea 100644 --- a/giga/evmonly/giga_store.go +++ b/giga/evmonly/giga_store.go @@ -10,7 +10,6 @@ import ( gigametrics "github.com/sei-protocol/sei-chain/giga/metrics" "github.com/sei-protocol/sei-chain/sei-db/common/keys" - seidbtypes "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" "github.com/sei-protocol/sei-chain/sei-db/proto" gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" ) @@ -26,12 +25,17 @@ var ( var _ StateReader = gigaSnapshotStateReader{} // NamedChangeSetEncoder converts an executor-native state result into the -// on-disk changesets understood by a giga store. It is called in the background, -// after the previous block's commit has landed and before this block's starts, on a -// copy of the block's changes that it must treat as immutable. +// on-disk changesets understood by a giga store. It is called synchronously +// while the block's read snapshot is still open. It must treat the input as +// immutable and must not retain references to it after returning. // -// It may read the store, which then holds every earlier block and none of this one; expanding a -// storage clear does so. +// What it returns must not alias the input either. The commit runs in the background, outliving +// the block result and its return to the pool, so an aliasing pair would be rewritten underneath +// the write by the next block. +// +// It must read only the changeset it is given. Encoding overlaps the previous block's commit, so +// an encoder that reads the store would see a store mid-write. Expanding a storage clear is the +// one exception, and the executor waits for that commit before encoding a block that has one. type NamedChangeSetEncoder func(StateChangeSet) ([]*proto.NamedChangeSet, error) // BlockChangeSetEncoder contributes named changesets that are committed in the @@ -106,28 +110,23 @@ func (e *Executor) executePreparedBlockWithStore(ctx context.Context, req Prepar return nil, err } gigametrics.SetPhase(gigametrics.PhaseStorage) - // The receipts need nothing but the result, so their write starts here and runs under the rest - // of the block's tail and the caller's, instead of after it. - e.blockPhases.SetPhase("start_receipts") - receipts := e.startReceiptWrite(ctx, blockNumber, result) - // One commit is in flight at a time, so the previous one lands before this block starts its - // own. It has had this block's whole execution to run, so it rarely still holds. A block encoder - // that reads the store needs it landed before it runs; one that reads only the result overlaps - // it. The state encoder runs in the background after this wait either way, so a storage clear - // it expands against the store sees every earlier block and none of this one. - settleBeforeBlockEncoder := e.blockChangeSetEncoder != nil && e.blockEncoderReadsStore - if settleBeforeBlockEncoder { + // Encoding that reads the store has to see a store holding every earlier block and none of + // this one, so the previous commit lands first. Encoding that reads only this block's own + // changes runs while that commit is still going, and waits below instead. + settleBeforeEncoding := e.encodingReadsTheStore(&result.ChangeSet) + if settleBeforeEncoding { e.blockPhases.SetPhase("await_commit") if err := e.awaitPipelineCommit(); err != nil { return nil, err } } - // The block encoder stays on the loop: what it stages is what the caller reports for the block, - // so it has to have run when this returns. - var extra []*proto.NamedChangeSet + e.blockPhases.SetPhase("encode_changesets") + changesets, err := e.changeSetEncoder(result.ChangeSet) + if err != nil { + return nil, fmt.Errorf("encode state changes for block %d: %w", req.Context.Number, err) + } if e.blockChangeSetEncoder != nil { - e.blockPhases.SetPhase("encode_block_changesets") - extra, err = e.blockChangeSetEncoder(req.Context, result) + extra, err := e.blockChangeSetEncoder(req.Context, result) if err != nil { return nil, fmt.Errorf("encode block changes for block %d: %w", req.Context.Number, err) } @@ -139,49 +138,48 @@ func (e *Executor) executePreparedBlockWithStore(ctx context.Context, req Prepar cs.Name, req.Context.Number, errBlockEncoderUsedEVMStoreKey) } } + changesets = append(changesets, extra...) } if err := ctx.Err(); err != nil { return nil, err } - if !settleBeforeBlockEncoder { + // An executor without a receipt store keeps none. + if e.receiptStore != nil { + e.blockPhases.SetPhase("encode_receipts") + records, err := e.receiptRecordsParallel(ctx, req.Context.Number, result) + if err != nil { + return nil, fmt.Errorf("encode receipts for block %d: %w", req.Context.Number, err) + } + e.blockPhases.SetPhase("write_receipts") + if err := e.receiptStore.SetReceipts(newReceiptContext(ctx, blockNumber), records); err != nil { + return nil, fmt.Errorf("store receipts for block %d: %w", req.Context.Number, err) + } + } + // One commit is in flight at a time, so the previous one lands before this block starts its + // own. It has had this block's whole execution to run, so it rarely still holds. + if !settleBeforeEncoding { e.blockPhases.SetPhase("await_commit") if err := e.awaitPipelineCommit(); err != nil { return nil, err } } - e.blockPhases.SetPhase("start_commit") - if err := e.startPipelineCommit(blockNumber, result, receipts, extra); err != nil { + e.blockPhases.SetPhase("commit_state") + if err := e.startPipelineCommit(blockNumber, changesets, &result.ChangeSet); err != nil { return nil, fmt.Errorf("commit state changes for block %d: %w", req.Context.Number, err) } ok = true return result, nil } -// receiptWrite is one block's background receipt write: done is closed once err is set. -type receiptWrite struct { - done chan struct{} - err error -} - -// AwaitReceipts blocks until the receipts of the last block this executor ran are in the receipt -// store, and reports the write's failure if it had one. It returns at once on an executor without -// a receipt store, which keeps none. +// encodingReadsTheStore reports whether encoding this block's changesets reads the store as well as +// the block's own changes, which decides whether encoding may overlap the previous block's commit. // -// Receipts are written in the background behind ExecutePreparedBlock. A caller that publishes a -// block to readers who expect its receipts has to wait here first. Unlike AwaitCommits it does not -// wait for the block's state commit, which keeps landing behind the next block. -func (e *Executor) AwaitReceipts() error { - if e == nil { - return nil - } - e.pipelineMu.Lock() - write := e.pipelineReceipts - e.pipelineMu.Unlock() - if write == nil { - return nil - } - <-write.done - return write.err +// Expanding a storage clear iterates the live store to find the slots to delete, so a block that +// clears one must not be encoded against a store mid-commit. A block encoder is caller-supplied and +// free to read whatever it likes, so one is assumed to read the store unless it was registered as +// store-independent: assuming otherwise would surrender the receipt-stage slack on every block. +func (e *Executor) encodingReadsTheStore(changes *StateChangeSet) bool { + return len(changes.StorageClears) > 0 || (e.blockChangeSetEncoder != nil && e.blockEncoderReadsStore) } // AwaitCommits blocks until every block this executor has run is committed, and reports the first @@ -271,7 +269,7 @@ func (e *Executor) pipelineFailureLocked() error { return e.pipelineFailure } if e.pipelineErr != nil { - return fmt.Errorf("persist block: %w", e.pipelineErr) + return fmt.Errorf("commit state changes: %w", e.pipelineErr) } return nil } @@ -298,7 +296,7 @@ func (e *Executor) awaitPipelineCommit() error { // Only the waiters on this commit retire it; a later one owns its own state. if e.pipelineDone == done { if e.pipelineErr != nil && e.pipelineFailure == nil { - e.pipelineFailure = fmt.Errorf("persist block: %w", e.pipelineErr) + e.pipelineFailure = fmt.Errorf("commit state changes: %w", e.pipelineErr) } e.pipelineDone = nil e.pipelineChanges = nil @@ -311,55 +309,13 @@ func (e *Executor) awaitPipelineCommit() error { return e.pipelineFailure } -// startReceiptWrite persists the block's receipts in the background, after the previous block's -// have landed, and returns the write to wait on. The block result is held until the write has -// landed. An executor without a receipt store keeps none, so its write is done on return. -// -// The write is recorded as the executor's newest, so AwaitReceipts finds it whether or not the -// block's commit is started afterwards. -func (e *Executor) startReceiptWrite(ctx context.Context, blockNumber int64, result *BlockResult) *receiptWrite { - receipts := &receiptWrite{done: make(chan struct{})} - if e.receiptStore == nil { - close(receipts.done) - return receipts - } - e.pipelineMu.Lock() - previous := e.pipelineReceipts - e.pipelineReceipts = receipts - e.pipelineMu.Unlock() - - // The write finishes even if the request that ran the block is cancelled: Close waits for it, - // and a failure is reported through the pipeline rather than by dropping the block. - bgCtx := context.WithoutCancel(ctx) - releaseResult := result.retain() - go func() { - defer close(receipts.done) - defer releaseResult() - // Receipts land in block order, so the store's version means every block up to it. The - // previous write owns the phase timer until it is done. - if previous != nil { - <-previous.done - if previous.err != nil { - receipts.err = fmt.Errorf("receipts for block %d not written after an earlier failure: %w", blockNumber, previous.err) - return - } - } - defer e.receiptPhases.Reset() - receipts.err = e.persistReceipts(bgCtx, blockNumber, result) - }() - return receipts -} - -// startPipelineCommit commits the block's encoded state changes in the background once its -// receipts have landed, and records what it changed, so the next block reads those changes through -// an overlay rather than waiting for the write. The state changes are encoded from a copy that -// outlives the block result. +// startPipelineCommit writes the block in the background and records what it changed, so the next +// block reads those changes through an overlay rather than waiting for the write. // // Commits stay ordered because only one is ever in flight: awaitPipelineCommit lands the previous // one before this is called. -func (e *Executor) startPipelineCommit(blockNumber int64, result *BlockResult, receipts *receiptWrite, extra []*proto.NamedChangeSet) error { - changes := result.ChangeSet.clone() - pending := newPendingChanges(changes) +func (e *Executor) startPipelineCommit(blockNumber int64, changesets []*proto.NamedChangeSet, changes *StateChangeSet) error { + pending := newPendingChanges(changes.clone()) done := make(chan struct{}) e.pipelineMu.Lock() if failure := e.pipelineFailure; failure != nil { @@ -373,67 +329,15 @@ func (e *Executor) startPipelineCommit(blockNumber int64, result *BlockResult, r e.pipelineMu.Unlock() go func() { - defer close(done) - defer e.pipelinePhases.Reset() - // A block whose receipts were lost is a failed block: its state is not committed, so the - // store never holds a block whose receipts cannot be read. - e.pipelinePhases.SetPhase("await_receipt_write") - <-receipts.done - if receipts.err != nil { - e.pipelineMu.Lock() - e.pipelineErr = receipts.err - e.pipelineMu.Unlock() - return - } - err := e.commitStateChanges(blockNumber, changes, extra) + err := e.stateStore.CommitStateChanges(blockNumber, changesets) e.pipelineMu.Lock() e.pipelineErr = err e.pipelineMu.Unlock() + close(done) }() return nil } -// persistReceipts encodes the block's receipts and writes them to the receipt store, returning once -// they are readable there. -func (e *Executor) persistReceipts(ctx context.Context, blockNumber int64, result *BlockResult) error { - e.receiptPhases.SetPhase("encode_receipts") - records, err := e.receiptRecordsParallel(ctx, uint64(blockNumber), result) //nolint:gosec // G115: non-negative, checked by the caller. - if err != nil { - return fmt.Errorf("encode receipts for block %d: %w", blockNumber, err) - } - e.receiptPhases.SetPhase("write_receipts") - if err := e.receiptStore.SetReceipts(newReceiptContext(ctx, blockNumber), records); err != nil { - return fmt.Errorf("store receipts for block %d: %w", blockNumber, err) - } - // A store that applies writes from its own queue reports the landing through its version; - // a write it dropped after accepting shows up as a version short of this block. - if waiter, ok := e.receiptStore.(seidbtypes.PendingWriteWaiter); ok { - e.receiptPhases.SetPhase("await_store") - waiter.WaitForPendingWrites() - if latest := e.receiptStore.LatestVersion(); latest < blockNumber { - return fmt.Errorf("receipts for block %d did not land: receipt store is at block %d", blockNumber, latest) - } - } - return nil -} - -// commitStateChanges encodes the block's state changes, appends the block encoder's, and commits -// them to the state store. -func (e *Executor) commitStateChanges(blockNumber int64, changes *StateChangeSet, extra []*proto.NamedChangeSet) error { - e.pipelinePhases.SetPhase("encode_changesets") - var stateChanges StateChangeSet - if changes != nil { - stateChanges = *changes - } - changesets, err := e.changeSetEncoder(stateChanges) - if err != nil { - return fmt.Errorf("encode state changes for block %d: %w", blockNumber, err) - } - changesets = append(changesets, extra...) - e.pipelinePhases.SetPhase("commit_state") - return e.stateStore.CommitStateChanges(blockNumber, changesets) -} - type gigaSnapshotStateReader struct { snapshot gigatypes.EVMStateView missingState StateReader diff --git a/giga/evmonly/giga_store_test.go b/giga/evmonly/giga_store_test.go index cd79f0cdc2..ed88661301 100644 --- a/giga/evmonly/giga_store_test.go +++ b/giga/evmonly/giga_store_test.go @@ -328,7 +328,7 @@ func TestExecutorGigaStoreFailuresDoNotCommitPartialState(t *testing.T) { require.Equal(t, 1, snapshot.closeCount) }) - t.Run("context canceled during encoding still persists the executed block", func(t *testing.T) { + t.Run("context canceled during encoding", func(t *testing.T) { snapshot := newMemoryGigaSnapshot(0) store := &recordingGigaStore{snapshot: snapshot} ctx, cancel := context.WithCancel(t.Context()) @@ -339,10 +339,9 @@ func TestExecutorGigaStoreFailuresDoNotCommitPartialState(t *testing.T) { result, err := executor.ExecuteBlock(ctx, BlockRequest{Context: blockContext(big.NewInt(testChainID))}) - require.NoError(t, err) - require.NotNil(t, result) - result.Release() - require.Len(t, store.commits, 1) + require.ErrorIs(t, err, context.Canceled) + require.Nil(t, result) + require.Empty(t, store.commits) require.Equal(t, 1, snapshot.closeCount) require.Equal(t, BlockResultPoolStats{Capacity: 1, Available: 1}, executor.ResultPoolStats()) }) diff --git a/giga/evmonly/pipeline_store_test.go b/giga/evmonly/pipeline_store_test.go index 9d0523456c..e0802ab80e 100644 --- a/giga/evmonly/pipeline_store_test.go +++ b/giga/evmonly/pipeline_store_test.go @@ -4,13 +4,10 @@ import ( "errors" "math/big" "testing" - "time" "github.com/ethereum/go-ethereum/crypto" "github.com/stretchr/testify/require" - sdk "github.com/sei-protocol/sei-chain/sei-cosmos/types" - "github.com/sei-protocol/sei-chain/sei-db/ledger_db/receipt" "github.com/sei-protocol/sei-chain/sei-db/proto" gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" ) @@ -327,72 +324,8 @@ func TestReadLatestAccountRestartsWhenItsPendingBlockRetiresUnderIt(t *testing.T require.Equal(t, uint64(2), got.Nonce, "the read must not replay retired block 41 over the landed state") } -// executePipelinedBlock runs one block through the pipelined path, which returns before the block's -// commit has landed. -func executePipelinedBlock(t *testing.T, executor *Executor, chainID *big.Int, number uint64, txs ...[]byte) *BlockResult { - t.Helper() - blockCtx := blockContext(chainID) - blockCtx.Number = number - prepared, err := executor.PrepareBlock(t.Context(), BlockRequest{Context: blockCtx, Txs: txs}) - require.NoError(t, err) - result, err := executor.ExecutePreparedBlock(t.Context(), prepared) - require.NoError(t, err) - return result -} - -// errTestReceiptWriteFailed is the failure a test receipt store reports from SetReceipts. -var errTestReceiptWriteFailed = errors.New("receipt write failed") - -// gatedCommitStore holds every state commit until the test releases it, so receipts can be -// observed landing while the state commit is still in flight. -type gatedCommitStore struct { - *recordingGigaStore - release chan struct{} -} - -func (s *gatedCommitStore) CommitStateChanges(blockNum int64, changeset []*proto.NamedChangeSet) error { - <-s.release - return s.recordingGigaStore.CommitStateChanges(blockNum, changeset) -} - -// A block's receipts are readable once AwaitReceipts returns, whether or not its state commit has -// landed: that is what lets the RPC head advance behind receipts alone. -func TestAwaitReceiptsReturnsBeforeTheStateCommitLands(t *testing.T) { - chainID := big.NewInt(testChainID) - key, err := crypto.GenerateKey() - require.NoError(t, err) - sender := crypto.PubkeyToAddress(key.PublicKey) - recipient := testAddress(0xa9) - - snapshot := newMemoryGigaSnapshot(40) - snapshot.setBalance(sender, big.NewInt(testFundedBalanceWei)) - store := &gatedCommitStore{recordingGigaStore: &recordingGigaStore{snapshot: snapshot}, release: make(chan struct{})} - receipts := NewMemoryReceiptStore() - executor := NewExecutor(Config{BlockResultPoolSize: 1}, withTestStores(store, receipts, noopChangeSetEncoder)) - defer executor.Close() - - rawTx := signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil) - result := executePipelinedBlock(t, executor, chainID, 41, rawTx) - txHash := result.Txs[0].Hash - result.Release() - - require.NoError(t, executor.AwaitReceipts()) - stored, err := receipts.GetReceipt(newReceiptContext(t.Context(), 41), txHash) - require.NoError(t, err) - require.Equal(t, uint64(41), stored.BlockNumber) - require.Empty(t, store.commits, "the state commit is still held") - // The result is needed only until its receipts are encoded, so the pool has it back while - // the commit is still running. - require.Equal(t, BlockResultPoolStats{Capacity: 1, Available: 1}, executor.ResultPoolStats()) - - close(store.release) - require.NoError(t, executor.AwaitCommits()) - require.Equal(t, []int64{41}, store.commitBlock) -} - -// An executor without a receipt store commits state only: AwaitReceipts has nothing to wait on -// and the result goes back to the pool as soon as the block returns. -func TestNoReceiptStoreCommitsStateWithoutWaitingOnReceipts(t *testing.T) { +// An executor without a receipt store commits state only; the block result still carries receipts. +func TestNoReceiptStoreCommitsStateOnly(t *testing.T) { chainID := big.NewInt(testChainID) key, err := crypto.GenerateKey() require.NoError(t, err) @@ -401,279 +334,28 @@ func TestNoReceiptStoreCommitsStateWithoutWaitingOnReceipts(t *testing.T) { snapshot := newMemoryGigaSnapshot(40) snapshot.setBalance(sender, big.NewInt(testFundedBalanceWei)) - store := &gatedCommitStore{recordingGigaStore: &recordingGigaStore{snapshot: snapshot}, release: make(chan struct{})} - executor := NewExecutor(Config{BlockResultPoolSize: 1}, withTestStores(store, nil, noopChangeSetEncoder)) + store := &recordingGigaStore{snapshot: snapshot} + executor := NewExecutor(Config{}, withTestStores(store, nil, noopChangeSetEncoder)) defer executor.Close() - rawTx := signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil) - result := executePipelinedBlock(t, executor, chainID, 41, rawTx) + result := executePipelinedBlock(t, executor, chainID, 41, + signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) require.Equal(t, uint64(1), result.Receipts[0].Status, "receipts are still produced for the block result") result.Release() - require.NoError(t, executor.AwaitReceipts()) - require.Empty(t, store.commits, "the state commit is still held") - require.Equal(t, BlockResultPoolStats{Capacity: 1, Available: 1}, executor.ResultPoolStats()) - - close(store.release) require.NoError(t, executor.AwaitCommits()) require.Equal(t, []int64{41}, store.commitBlock) } -// signallingReceiptStore reports each block whose receipts it was handed on written. -type signallingReceiptStore struct { - *MemoryReceiptStore - written chan int64 -} - -func (s *signallingReceiptStore) SetReceipts(ctx sdk.Context, records []receipt.ReceiptRecord) error { - if err := s.MemoryReceiptStore.SetReceipts(ctx, records); err != nil { - return err - } - s.written <- ctx.BlockHeight() - return nil -} - -// A block's receipts are written while the loop is still waiting on the previous block's commit, -// rather than after it: the write needs only the result. -func TestReceiptWriteStartsBeforeThePreviousCommitLands(t *testing.T) { - chainID := big.NewInt(testChainID) - key, err := crypto.GenerateKey() - require.NoError(t, err) - sender := crypto.PubkeyToAddress(key.PublicKey) - recipient := testAddress(0xa9) - - snapshot := newMemoryGigaSnapshot(40) - snapshot.setBalance(sender, big.NewInt(testFundedBalanceWei)) - store := &gatedCommitStore{recordingGigaStore: &recordingGigaStore{snapshot: snapshot}, release: make(chan struct{})} - receipts := &signallingReceiptStore{MemoryReceiptStore: NewMemoryReceiptStore(), written: make(chan int64, 2)} - // A store-reading block encoder makes the loop wait for the previous commit before it runs. - blockEncoder := func(BlockContext, *BlockResult) ([]*proto.NamedChangeSet, error) { return nil, nil } - executor := NewExecutor(Config{}, - withTestStores(store, receipts, noopChangeSetEncoder), - WithBlockChangeSetEncoder(blockEncoder)) - defer executor.Close() - - executePipelinedBlock(t, executor, chainID, 41, - signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)).Release() - require.NoError(t, executor.AwaitReceipts()) - require.Equal(t, int64(41), <-receipts.written) - - // Block 42 blocks on the loop until 41's commit is released. - blockCtx := blockContext(chainID) - blockCtx.Number = 42 - prepared, err := executor.PrepareBlock(t.Context(), BlockRequest{Context: blockCtx, - Txs: [][]byte{signLegacyTx(t, key, chainID, 1, &recipient, big.NewInt(9), nil)}}) - require.NoError(t, err) - executed := make(chan error, 1) - go func() { - result, err := executor.ExecutePreparedBlock(t.Context(), prepared) - if err == nil { - result.Release() - } - executed <- err - }() - - require.Equal(t, int64(42), <-receipts.written, "block 42's receipts are written while its loop waits") - require.Empty(t, store.commits, "the previous commit is still held") - select { - case err := <-executed: - t.Fatalf("block 42 returned before the previous commit landed: %v", err) - default: - } - - close(store.release) - require.NoError(t, <-executed) - require.NoError(t, executor.AwaitCommits()) - require.Equal(t, []int64{41, 42}, store.commitBlock) -} - -// A receipt write that fails is a failed block: the state commit is not attempted and both waiters -// report it. -func TestFailedReceiptWriteFailsTheBlockBeforeItsStateCommit(t *testing.T) { - chainID := big.NewInt(testChainID) - key, err := crypto.GenerateKey() - require.NoError(t, err) - sender := crypto.PubkeyToAddress(key.PublicKey) - recipient := testAddress(0xa9) - - snapshot := newMemoryGigaSnapshot(40) - snapshot.setBalance(sender, big.NewInt(testFundedBalanceWei)) - store := &recordingGigaStore{snapshot: snapshot} - executor := NewExecutor(Config{}, withTestStores(store, &failingReceiptStore{MemoryReceiptStore: NewMemoryReceiptStore(), err: errTestReceiptWriteFailed}, noopChangeSetEncoder)) - defer executor.Close() - - executePipelinedBlock(t, executor, chainID, 41, - signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) - - require.ErrorIs(t, executor.AwaitReceipts(), errTestReceiptWriteFailed) - require.ErrorIs(t, executor.AwaitCommits(), errTestReceiptWriteFailed) - require.Empty(t, store.commits, "state must not be committed for a block whose receipts were not") -} - -// The block after a failed receipt write fails too, before its own receipts are attempted: a store -// whose version has moved past a block whose receipts it never got would claim them as written. -func TestReceiptWriteAfterAFailedOneIsNotAttempted(t *testing.T) { - chainID := big.NewInt(testChainID) - key, err := crypto.GenerateKey() - require.NoError(t, err) - sender := crypto.PubkeyToAddress(key.PublicKey) - recipient := testAddress(0xa9) - - snapshot := newMemoryGigaSnapshot(40) - snapshot.setBalance(sender, big.NewInt(testFundedBalanceWei)) - store := &recordingGigaStore{snapshot: snapshot} - receipts := &failingReceiptStore{MemoryReceiptStore: NewMemoryReceiptStore(), err: errTestReceiptWriteFailed} - executor := NewExecutor(Config{}, withTestStores(store, receipts, noopChangeSetEncoder)) - defer executor.Close() - - executePipelinedBlock(t, executor, chainID, 41, - signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) - require.ErrorIs(t, executor.AwaitReceipts(), errTestReceiptWriteFailed) - - // The store would accept block 42's receipts now; the executor must not offer them. - receipts.err = nil - blockCtx := blockContext(chainID) - blockCtx.Number = 42 - rawTx := signLegacyTx(t, key, chainID, 1, &recipient, big.NewInt(9), nil) - prepared, err := executor.PrepareBlock(t.Context(), BlockRequest{Context: blockCtx, Txs: [][]byte{rawTx}}) - require.NoError(t, err) - _, err = executor.ExecutePreparedBlock(t.Context(), prepared) - require.ErrorIs(t, err, errTestReceiptWriteFailed) - - require.ErrorIs(t, executor.AwaitReceipts(), errTestReceiptWriteFailed) - _, err = receipts.GetReceipt(newReceiptContext(t.Context(), 42), decodeTx(t, rawTx).Hash()) - require.ErrorIs(t, err, receipt.ErrNotFound, "block 42's receipts must not land over the hole at 41") - require.Empty(t, store.commits) -} - -// droppingReceiptStore accepts every write and applies none of them, the way a queued store behaves -// once an earlier write has failed: it takes the block, waits out its queue, and its version never -// reaches it. -type droppingReceiptStore struct { - *MemoryReceiptStore -} - -func (s *droppingReceiptStore) SetReceipts(sdk.Context, []receipt.ReceiptRecord) error { return nil } - -func (s *droppingReceiptStore) WaitForPendingWrites() {} - -// A write the store accepted but never applied is this block's failure, not the next one's. -func TestReceiptWriteThatNeverLandsFailsTheBlock(t *testing.T) { - chainID := big.NewInt(testChainID) - key, err := crypto.GenerateKey() - require.NoError(t, err) - sender := crypto.PubkeyToAddress(key.PublicKey) - recipient := testAddress(0xa9) - - snapshot := newMemoryGigaSnapshot(40) - snapshot.setBalance(sender, big.NewInt(testFundedBalanceWei)) - store := &recordingGigaStore{snapshot: snapshot} - receipts := &droppingReceiptStore{MemoryReceiptStore: NewMemoryReceiptStore()} - executor := NewExecutor(Config{}, withTestStores(store, receipts, noopChangeSetEncoder)) - defer executor.Close() - - executePipelinedBlock(t, executor, chainID, 41, - signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) - - err = executor.AwaitReceipts() - require.ErrorContains(t, err, "receipts for block 41 did not land") - require.ErrorIs(t, executor.AwaitCommits(), err) - require.Empty(t, store.commits, "state must not be committed for a block whose receipts were not") -} - -// The state encoder runs after the block result may have been reused, so it must be handed the -// block's own changes and nothing else. -func TestBackgroundEncoderSeesTheBlocksOwnChanges(t *testing.T) { - chainID := big.NewInt(testChainID) - key, err := crypto.GenerateKey() - require.NoError(t, err) - sender := crypto.PubkeyToAddress(key.PublicKey) - recipient := testAddress(0xa9) - - snapshot := newMemoryGigaSnapshot(40) - snapshot.setBalance(sender, big.NewInt(testFundedBalanceWei)) - store := &recordingGigaStore{snapshot: snapshot} - var encoded []StateChangeSet - encoder := func(changes StateChangeSet) ([]*proto.NamedChangeSet, error) { - encoded = append(encoded, changes) - return noopChangeSetEncoder(changes) - } - executor := NewExecutor(Config{BlockResultPoolSize: 1}, withTestStores(store, NewMemoryReceiptStore(), encoder)) - defer executor.Close() - - first := executePipelinedBlock(t, executor, chainID, 41, - signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) - want := *first.ChangeSet.clone() - first.Release() - require.NoError(t, executor.AwaitReceipts()) - // Takes the pooled result back while block 41's encoder may still be running. - second := executePipelinedBlock(t, executor, chainID, 42, - signLegacyTx(t, key, chainID, 1, &recipient, big.NewInt(9), nil)) - second.Release() - - require.NoError(t, executor.AwaitCommits()) - require.Len(t, encoded, 2) - require.Equal(t, want, encoded[0]) - require.Equal(t, []int64{41, 42}, store.commitBlock) -} - -// errTestEncoderFailed is the failure a test block encoder reports. -var errTestEncoderFailed = errors.New("block encoder failed") - -// gatedReceiptStore holds every receipt write until release is closed. -type gatedReceiptStore struct { - *MemoryReceiptStore - release chan struct{} -} - -func (s *gatedReceiptStore) SetReceipts(ctx sdk.Context, records []receipt.ReceiptRecord) error { - <-s.release - return s.MemoryReceiptStore.SetReceipts(ctx, records) -} - -// A block that fails on the loop after its receipts were handed off has no state commit for Close to -// drain, so Close waits for the receipt write itself and the receipts are readable once it returns. -func TestCloseWaitsForTheReceiptsOfABlockThatFailedAfterHandingThemOff(t *testing.T) { - chainID := big.NewInt(testChainID) - key, err := crypto.GenerateKey() - require.NoError(t, err) - sender := crypto.PubkeyToAddress(key.PublicKey) - recipient := testAddress(0xa9) - - snapshot := newMemoryGigaSnapshot(40) - snapshot.setBalance(sender, big.NewInt(testFundedBalanceWei)) - store := &recordingGigaStore{snapshot: snapshot} - receipts := &gatedReceiptStore{MemoryReceiptStore: NewMemoryReceiptStore(), release: make(chan struct{})} - blockEncoder := func(BlockContext, *BlockResult) ([]*proto.NamedChangeSet, error) { - return nil, errTestEncoderFailed - } - executor := NewExecutor(Config{}, - withTestStores(store, receipts, noopChangeSetEncoder), - WithBlockChangeSetEncoder(blockEncoder)) - - rawTx := signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil) +// executePipelinedBlock runs one block through the pipelined path, which returns before the block's +// commit has landed. +func executePipelinedBlock(t *testing.T, executor *Executor, chainID *big.Int, number uint64, txs ...[]byte) *BlockResult { + t.Helper() blockCtx := blockContext(chainID) - blockCtx.Number = 41 - prepared, err := executor.PrepareBlock(t.Context(), BlockRequest{Context: blockCtx, Txs: [][]byte{rawTx}}) + blockCtx.Number = number + prepared, err := executor.PrepareBlock(t.Context(), BlockRequest{Context: blockCtx, Txs: txs}) require.NoError(t, err) - _, err = executor.ExecutePreparedBlock(t.Context(), prepared) - require.ErrorIs(t, err, errTestEncoderFailed) - - closed := make(chan struct{}) - go func() { - executor.Close() - close(closed) - }() - select { - case <-closed: - t.Fatal("Close returned while the receipt write was still held") - case <-time.After(50 * time.Millisecond): - } - - close(receipts.release) - <-closed - stored, err := receipts.GetReceipt(newReceiptContext(t.Context(), 41), decodeTx(t, rawTx).Hash()) + result, err := executor.ExecutePreparedBlock(t.Context(), prepared) require.NoError(t, err) - require.Equal(t, uint64(41), stored.BlockNumber) - require.Empty(t, store.commits, "a block that failed on the loop commits no state") + return result } diff --git a/sei-db/ledger_db/receipt/litt_receipt_store.go b/sei-db/ledger_db/receipt/litt_receipt_store.go index eb471ece02..e9a28e9b59 100644 --- a/sei-db/ledger_db/receipt/litt_receipt_store.go +++ b/sei-db/ledger_db/receipt/litt_receipt_store.go @@ -89,7 +89,7 @@ type littReceiptStore struct { // Receipt writes waiting to be applied, and the meter for time spent waiting on a full queue. A // whole write is queued, so the depth is the receipt write's own. Nil means writes apply inline. - writes chan *receiptWrite + writes chan receiptWrite // Orders admitting a write against shutting the writer down, so none is accepted into a queue // that will not be drained. queueWrite holds it shared; Close takes it exclusively. @@ -99,18 +99,12 @@ type littReceiptStore struct { writeQueue *seidbmetrics.QueueMeter writeErr atomic.Pointer[error] stopSampling context.CancelFunc - // The last write admitted to the queue, for WaitForPendingWrites to wait on. queueMu orders - // publishing it with the send, so the marker is never a write that is still behind another. - queueMu sync.Mutex - lastQueued atomic.Pointer[receiptWrite] } -// receiptWrite is one block's receipts, waiting to be applied. landed is closed once the writer is -// done with it, whether or not the write succeeded. +// receiptWrite is one block's receipts, waiting to be applied. type receiptWrite struct { height int64 receipts []ReceiptRecord - landed chan struct{} } // writeQueueSampleIntervalSeconds is how often the write queue's depth is read. Sampling on a timer @@ -249,7 +243,7 @@ func newLittReceiptStore(cfg dbconfig.ReceiptStoreConfig, storeKey sdk.StoreKey) receiptMeter := otel.Meter("seidb_receipt") s.writePhases = seidbmetrics.NewPhaseTimer(receiptMeter, "receipt_store_write") if cfg.AsyncWriteBuffer > 0 { - s.writes = make(chan *receiptWrite, cfg.AsyncWriteBuffer) + s.writes = make(chan receiptWrite, cfg.AsyncWriteBuffer) s.writeQueue = seidbmetrics.NewQueueMeter(receiptMeter, "receipt_write") s.startWriter() @@ -355,16 +349,7 @@ func (s *littReceiptStore) SetReceipts(ctx sdk.Context, receipts []ReceiptRecord if err := s.writeFailure(); err != nil { return err } - return s.queueWrite(&receiptWrite{height: ctx.BlockHeight(), receipts: receipts, landed: make(chan struct{})}) -} - -// WaitForPendingWrites blocks until every write queued before the call has been applied, so a -// receipt SetReceipts accepted is readable when this returns. It returns at once when writes apply -// inline. -func (s *littReceiptStore) WaitForPendingWrites() { - if write := s.lastQueued.Load(); write != nil { - <-write.landed - } + return s.queueWrite(receiptWrite{height: ctx.BlockHeight(), receipts: receipts}) } // ErrStoreClosed is returned by a write the store can no longer apply, the writer having stopped. @@ -372,7 +357,7 @@ var ErrStoreClosed = errors.New("receipt store is closed") // queueWrite hands a write to the writer, waiting for room when the queue is full and refusing once // the store is closing. -func (s *littReceiptStore) queueWrite(write *receiptWrite) error { +func (s *littReceiptStore) queueWrite(write receiptWrite) error { // Held across the send, not merely to read the flag: Close takes it exclusively before stopping // the writer, so a write admitted here always reaches a writer that is still running. s.admission.RLock() @@ -380,9 +365,6 @@ func (s *littReceiptStore) queueWrite(write *receiptWrite) error { if s.closing { return ErrStoreClosed } - s.queueMu.Lock() - defer s.queueMu.Unlock() - s.lastQueued.Store(write) seidbmetrics.Send(s.writeQueue, s.writes, write) return nil } @@ -597,8 +579,7 @@ func (s *littReceiptStore) startWriter() { // applyWrite performs one queued write, keeping the first failure for its callers to collect. // Nothing is applied after a failure: a later block carries its own version marker and would publish // a head above one whose receipts were never written. -func (s *littReceiptStore) applyWrite(write *receiptWrite) { - defer close(write.landed) +func (s *littReceiptStore) applyWrite(write receiptWrite) { if s.writeFailure() != nil { return } diff --git a/sei-db/ledger_db/receipt/litt_write_failure_internal_test.go b/sei-db/ledger_db/receipt/litt_write_failure_internal_test.go index 4386c941e1..e53242cd06 100644 --- a/sei-db/ledger_db/receipt/litt_write_failure_internal_test.go +++ b/sei-db/ledger_db/receipt/litt_write_failure_internal_test.go @@ -128,7 +128,7 @@ func TestWriteAfterCloseIsRefusedWithAFullQueue(t *testing.T) { // Leftovers with no writer behind them: the send has nowhere to go and nobody to take it. for len(s.writes) < cap(s.writes) { - s.writes <- &receiptWrite{height: 1, landed: make(chan struct{})} + s.writes <- receiptWrite{height: 1} } txHash, rcpt := littCtxTestReceipt(1, 0, common.HexToAddress("0xfa41"), common.HexToHash("0xfa42"), 1) diff --git a/sei-db/ledger_db/receipt/littidx_test.go b/sei-db/ledger_db/receipt/littidx_test.go index 37af9a04f4..4c9119fc53 100644 --- a/sei-db/ledger_db/receipt/littidx_test.go +++ b/sei-db/ledger_db/receipt/littidx_test.go @@ -13,7 +13,6 @@ import ( "github.com/sei-protocol/sei-chain/sei-cosmos/testutil" sdk "github.com/sei-protocol/sei-chain/sei-cosmos/types" dbconfig "github.com/sei-protocol/sei-chain/sei-db/config" - seidbtypes "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" "github.com/sei-protocol/sei-chain/sei-db/ledger_db/receipt" "github.com/sei-protocol/sei-chain/x/evm/types" "github.com/stretchr/testify/require" @@ -130,41 +129,6 @@ func TestLittIdxWriteBufferBoundsLag(t *testing.T) { 5*time.Second, time.Millisecond) } -// TestLittIdxWaitForPendingWritesLandsEveryQueuedBlock pins that once WaitForPendingWrites returns, -// every write SetReceipts accepted before it is readable, whichever way the writer was scheduled. -func TestLittIdxWaitForPendingWritesLandsEveryQueuedBlock(t *testing.T) { - storeKey := storetypes.NewKVStoreKey("evm") - tkey := storetypes.NewTransientStoreKey("evm_transient") - ctx := testutil.DefaultContext(storeKey, tkey).WithBlockHeight(1) - cfg := dbconfig.DefaultReceiptStoreConfig() - cfg.Backend = "littidx" - cfg.DBDirectory = t.TempDir() - cfg.KeepRecent = 0 - cfg.AsyncWriteBuffer = 4 - - store, err := receipt.NewReceiptStore(cfg, storeKey) - require.NoError(t, err) - t.Cleanup(func() { _ = store.Close() }) - waiter, ok := store.(seidbtypes.PendingWriteWaiter) - require.True(t, ok, "the queued receipt store must let a caller wait for its writes") - - // Nothing queued yet: returns at once. - waiter.WaitForPendingWrites() - - addr := common.HexToAddress("0xabc") - const blocks = 6 - for block := uint64(1); block <= blocks; block++ { - record := litReceipt(block, 0, addr, common.HexToHash("0xdead")) - require.NoError(t, store.SetReceipts(ctx.WithBlockHeight(int64(block)), //nolint:gosec // small test heights - []receipt.ReceiptRecord{record})) - waiter.WaitForPendingWrites() - require.Equal(t, int64(block), store.LatestVersion()) //nolint:gosec // small test heights - got, err := store.GetReceipt(ctx.WithBlockHeight(int64(block)), record.TxHash) //nolint:gosec // small test heights - require.NoError(t, err) - require.Equal(t, block, got.BlockNumber) - } -} - func setupLittIdx(t *testing.T, dir string) (receipt.ReceiptStore, sdk.Context) { t.Helper() return setupLittIdxPar(t, dir, dbconfig.DefaultReceiptLogFilterParallelism) diff --git a/sei-db/ledger_db/receipt/receipt_store.go b/sei-db/ledger_db/receipt/receipt_store.go index 965a6e523f..2b5ecb34d1 100644 --- a/sei-db/ledger_db/receipt/receipt_store.go +++ b/sei-db/ledger_db/receipt/receipt_store.go @@ -273,14 +273,6 @@ func (s *receiptStore) LatestVersion() int64 { return s.db.GetLatestVersion() } -// WaitForPendingWrites blocks until every write the state store has queued is applied. It returns -// at once when the store applies writes inline. -func (s *receiptStore) WaitForPendingWrites() { - if waiter, ok := s.db.(seidbtypes.PendingWriteWaiter); ok { - waiter.WaitForPendingWrites() - } -} - func (s *receiptStore) SetLatestVersion(version int64) error { return s.db.SetLatestVersion(version) } diff --git a/sei-tendermint/internal/evmonlyapp/app.go b/sei-tendermint/internal/evmonlyapp/app.go index dc33221834..324b99793e 100644 --- a/sei-tendermint/internal/evmonlyapp/app.go +++ b/sei-tendermint/internal/evmonlyapp/app.go @@ -829,19 +829,12 @@ func (a *evmOnlyApplication) pendingCursor(height int64) (evmOnlyCursor, error) } // Commit acknowledges the finalized block as the one the chain builds on: the -// height it advances is what RPC serves as latest. On a node with a receipt -// store it first waits for the block's receipts to land, so a block RPC serves -// has them; a node without one has nothing to wait for. The block's state -// commit may still be landing in the store: it is not waited for here, since -// that would put the write back on the block loop. A commit that fails halts -// the node here or through the next FinalizeBlock, and a restart resumes from -// the store's own version. +// height it advances is what RPC serves as latest. Neither the block's state +// commit nor its queued receipt write is waited for here, since that would put +// the write back on the block loop, so the newest block's receipts can trail +// latest briefly. A commit that fails halts the node through the next +// FinalizeBlock, and a restart resumes from the store's own version. func (a *evmOnlyApplication) Commit(context.Context) (*abci.ResponseCommit, error) { - if executor, ok := a.settler.Load().Get(); ok { - if err := executor.AwaitReceipts(); err != nil { - return nil, err - } - } for state := range a.cursor.Lock() { pending, ok := state.pending.Get() if !ok { diff --git a/sei-tendermint/internal/evmonlyapp/app_test.go b/sei-tendermint/internal/evmonlyapp/app_test.go index 02d24caef4..4243c9dbbb 100644 --- a/sei-tendermint/internal/evmonlyapp/app_test.go +++ b/sei-tendermint/internal/evmonlyapp/app_test.go @@ -196,30 +196,6 @@ func TestEVMOnlyApplicationExecutesRawEthereumBlock(t *testing.T) { require.Equal(t, uint64(1), receipt.BlockNumber) } -// The height Commit advances is what RPC serves as latest, so the block's receipts are readable -// from the store the moment Commit returns, without waiting for anything else to land. -func TestEVMOnlyApplicationCommitReturnsWithTheBlocksReceiptsReadable(t *testing.T) { - storage := openEVMOnlyTestStorage(t, t.TempDir()) - app, err := NewEVMOnlyApplication(evmOnlyTestChainID, nil, storage, evmonly.NewFlatKVChangeSetEncoder(storage.SC())) - require.NoError(t, err) - t.Cleanup(func() { closeEVMOnlyTestApp(t, app, storage) }) - _, err = app.InitChain(evmOnlyTestInitChain()) - require.NoError(t, err) - - for height := int64(1); height <= 3; height++ { - raw, _ := signedEVMOnlyTestTx(t, evmOnlyTestChainID, 0) - finalizeAndCommitEVMOnlyTestBlock(t, app, evmOnlyTestBlock(height, raw)) - require.Equal(t, height, app.LastBlockHeight()) - - var tx ethtypes.Transaction - require.NoError(t, tx.UnmarshalBinary(raw)) - receiptCtx := sdk.NewContext(nil, tmproto.Header{Height: height}, false).WithContext(t.Context()) - receipt, err := storage.ReceiptDB().GetReceipt(receiptCtx, tx.Hash()) - require.NoError(t, err) - require.Equal(t, uint64(height), receipt.BlockNumber) //nolint:gosec // G115: test heights are positive. - } -} - func TestEVMOnlyApplicationRejectsWrongChain(t *testing.T) { app := newInitializedEVMOnlyTestApp(t) raw, _ := signedEVMOnlyTestTx(t, evmOnlyTestChainID+1, 0)