diff --git a/giga/evmonly/executor.go b/giga/evmonly/executor.go index 574e402c1b..ca03f6fb1f 100644 --- a/giga/evmonly/executor.go +++ b/giga/evmonly/executor.go @@ -52,7 +52,9 @@ type Executor struct { pipelineMu sync.Mutex pipelineDone chan struct{} pipelineErr error - pipelineChanges *StateChangeSet + 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. pipelineFailure error } diff --git a/giga/evmonly/giga_store.go b/giga/evmonly/giga_store.go index 1e6538432f..5f6dfe44ea 100644 --- a/giga/evmonly/giga_store.go +++ b/giga/evmonly/giga_store.go @@ -92,7 +92,7 @@ func (e *Executor) executePreparedBlockWithStore(ctx context.Context, req Prepar snapshot: snapshot, missingState: e.missingState, } - source = newPendingOverlay(source, pending) + source = pending.overlay(source) e.blockPhases.SetPhase("execute") result, err := e.executePreparedBlock(ctx, req, source) @@ -201,12 +201,79 @@ func (e *Executor) AwaitCommits() error { // pipelinePending returns the changes of a block whose commit has not been waited on yet, or nil // when the store is caught up. -func (e *Executor) pipelinePending() *StateChangeSet { +func (e *Executor) pipelinePending() *pendingChanges { e.pipelineMu.Lock() defer e.pipelineMu.Unlock() return e.pipelineChanges } +// LatestAccount is the balance and nonce of an account after the last block this executor ran. +type LatestAccount struct { + Balance *big.Int + Nonce uint64 +} + +// ReadLatestAccount returns addr's balance and nonce after the last block this executor ran, +// without waiting for that block's commit to land. It reports the first failed commit instead of +// state that lacks the failed block. +func (e *Executor) ReadLatestAccount(addr common.Address) (LatestAccount, error) { + if e.stateStore == nil { + return LatestAccount{}, errMissingStateStore + } + for { + e.pipelineMu.Lock() + pending, generation, failure := e.pipelineChanges, e.pipelineGeneration, e.pipelineFailureLocked() + e.pipelineMu.Unlock() + if failure != nil { + return LatestAccount{}, failure + } + snapshot := e.stateStore.OpenView() + if snapshot == nil { + return LatestAccount{}, errors.New("giga store returned a nil snapshot") + } + account, ok := e.readLatestAccount(snapshot, pending, generation, addr) + snapshot.Close() + if ok { + return account, nil + } + } +} + +// readLatestAccount reads addr through pending laid over snapshot. It reports false when another +// commit started after generation was read, since the view may then hold a later block's writes and +// pending would replay older values over them; the caller reads again. +func (e *Executor) readLatestAccount(snapshot gigatypes.EVMStateView, pending *pendingChanges, generation uint64, addr common.Address) (LatestAccount, bool) { + e.pipelineMu.Lock() + moved := e.pipelineGeneration != generation + e.pipelineMu.Unlock() + if moved { + return LatestAccount{}, false + } + reader := pending.overlay(gigaSnapshotStateReader{snapshot: snapshot, missingState: e.missingState}) + if rowReader, ok := reader.(accountSnapshotReader); ok { + if row, ok := rowReader.ReadAccount(addr); ok { + balance := row.Balance + if balance == nil { + balance = new(big.Int) + } + return LatestAccount{Balance: balance, Nonce: row.Nonce}, true + } + } + return LatestAccount{Balance: reader.GetBalance(addr), Nonce: reader.GetNonce(addr)}, true +} + +// pipelineFailureLocked returns the first failed commit, whether or not a waiter has retired it yet. +// Callers hold pipelineMu. +func (e *Executor) pipelineFailureLocked() error { + if e.pipelineFailure != nil { + return e.pipelineFailure + } + if e.pipelineErr != nil { + return fmt.Errorf("commit state changes: %w", e.pipelineErr) + } + return nil +} + // awaitPipelineCommit blocks until the in-flight commit has landed, reporting the first commit that // failed. After it returns the store holds every block this executor has run, so the next view // opens on a known height and needs no overlay. @@ -248,7 +315,7 @@ func (e *Executor) awaitPipelineCommit() error { // 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, changesets []*proto.NamedChangeSet, changes *StateChangeSet) error { - pending := changes.clone() + pending := newPendingChanges(changes.clone()) done := make(chan struct{}) e.pipelineMu.Lock() if failure := e.pipelineFailure; failure != nil { @@ -256,6 +323,7 @@ func (e *Executor) startPipelineCommit(blockNumber int64, changesets []*proto.Na return failure } e.pipelineChanges = pending + e.pipelineGeneration++ e.pipelineDone = done e.pipelineErr = nil e.pipelineMu.Unlock() diff --git a/giga/evmonly/pipeline_overlay.go b/giga/evmonly/pipeline_overlay.go index b44fc22e23..947d85ab32 100644 --- a/giga/evmonly/pipeline_overlay.go +++ b/giga/evmonly/pipeline_overlay.go @@ -18,7 +18,12 @@ import ( // a returned balance or code in place would corrupt the pending state for every other reader. type pendingOverlay struct { base StateReader + *pendingChanges +} +// pendingChanges is a StateChangeSet indexed by address and slot, built once per block so every +// reader that lays it over a view shares the index. +type pendingChanges struct { balances map[common.Address]*big.Int nonces map[common.Address]uint64 code map[common.Address][]byte @@ -28,14 +33,26 @@ type pendingOverlay struct { cleared map[common.Address]struct{} } -// newPendingOverlay indexes changes for lookup. It returns base unchanged when there is nothing to +// newPendingOverlay lays changes over base. It returns base unchanged when there is nothing to // overlay, so a caller pays nothing for the first block or after a commit has caught up. func newPendingOverlay(base StateReader, changes *StateChangeSet) StateReader { - if changes == nil || changes.isEmpty() { + return newPendingChanges(changes).overlay(base) +} + +// overlay returns base with the pending changes laid over it, or base itself when there are none. +func (c *pendingChanges) overlay(base StateReader) StateReader { + if c == nil { return base } - o := &pendingOverlay{ - base: base, + return &pendingOverlay{base: base, pendingChanges: c} +} + +// newPendingChanges indexes changes for lookup, or returns nil when they would change nothing. +func newPendingChanges(changes *StateChangeSet) *pendingChanges { + if changes == nil || changes.isEmpty() { + return nil + } + o := &pendingChanges{ balances: make(map[common.Address]*big.Int, len(changes.Balances)), nonces: make(map[common.Address]uint64, len(changes.Nonces)), code: make(map[common.Address][]byte, len(changes.Code)), diff --git a/giga/evmonly/pipeline_store_test.go b/giga/evmonly/pipeline_store_test.go index 733375e573..e0802ab80e 100644 --- a/giga/evmonly/pipeline_store_test.go +++ b/giga/evmonly/pipeline_store_test.go @@ -169,6 +169,161 @@ func TestRetiringACommitWhileOpeningAViewKeepsThePreviousBlocksState(t *testing. require.Contains(t, second.ChangeSet.Balances, BalanceChange{Address: last, Balance: big.NewInt(1_000)}) } +// The store's view never advances, so the account state a block produced is only reachable through +// the pending overlay until AwaitCommits. A latest-account read must report it without settling. +func TestReadLatestAccountSeesTheBlockWhoseCommitIsInFlight(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, NewMemoryReceiptStore(), noopChangeSetEncoder)) + defer executor.Close() + + before, err := executor.ReadLatestAccount(sender) + require.NoError(t, err) + require.Equal(t, LatestAccount{Balance: big.NewInt(testFundedBalanceWei)}, before) + + result := executePipelinedBlock(t, executor, chainID, 41, + signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) + require.Equal(t, uint64(1), result.Txs[0].Status) + + got, err := executor.ReadLatestAccount(sender) + require.NoError(t, err) + require.Equal(t, uint64(1), got.Nonce) + require.Equal(t, -1, got.Balance.Cmp(big.NewInt(testFundedBalanceWei)), "gas and value must be deducted") + paid, err := executor.ReadLatestAccount(recipient) + require.NoError(t, err) + require.Equal(t, LatestAccount{Balance: big.NewInt(7)}, paid) + + require.NoError(t, executor.AwaitCommits()) +} + +// A commit that failed leaves the store behind the run, so a latest-account read reports the +// failure rather than state that omits the failed block. +func TestReadLatestAccountReportsAFailedCommit(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, commitErr: errTestCommitFailed} + executor := NewExecutor(Config{}, withTestStores(store, NewMemoryReceiptStore(), noopChangeSetEncoder)) + defer executor.Close() + + executePipelinedBlock(t, executor, chainID, 41, + signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) + require.ErrorIs(t, executor.AwaitCommits(), errTestCommitFailed) + + _, err = executor.ReadLatestAccount(sender) + require.ErrorIs(t, err, errTestCommitFailed) +} + +// A failed commit is reported as soon as the commit has returned, before any waiter has retired it. +func TestReadLatestAccountReportsAFailedCommitNobodyHasAwaited(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, commitErr: errTestCommitFailed} + executor := NewExecutor(Config{}, withTestStores(store, NewMemoryReceiptStore(), noopChangeSetEncoder)) + defer executor.Close() + + executePipelinedBlock(t, executor, chainID, 41, + signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) + executor.pipelineMu.Lock() + done := executor.pipelineDone + executor.pipelineMu.Unlock() + require.NotNil(t, done) + <-done + + _, err = executor.ReadLatestAccount(sender) + require.ErrorIs(t, err, errTestCommitFailed) +} + +// A block that lands its commit and starts the next one between a reader's pending read and its +// view must not leave the reader replaying the older block over the newer state; the read starts +// over instead. +func TestReadLatestAccountRestartsWhenABlockLandsUnderIt(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 := &retiringOnOpenStore{recordingGigaStore: &recordingGigaStore{snapshot: snapshot}} + executor := NewExecutor(Config{}, withTestStores(store, NewMemoryReceiptStore(), noopChangeSetEncoder)) + defer executor.Close() + + executePipelinedBlock(t, executor, chainID, 41, + signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) + + // The reader has block 41 in hand as pending; block 42 lands underneath while its view opens. + opens := 0 + store.retire = func() { + opens++ + if opens == 1 { + store.retire = nil + executePipelinedBlock(t, executor, chainID, 42, + signLegacyTx(t, key, chainID, 1, &recipient, big.NewInt(7), nil)) + } + } + got, err := executor.ReadLatestAccount(sender) + require.NoError(t, err) + require.Equal(t, uint64(2), got.Nonce, "the read must reflect block 42, not replay block 41 over it") + require.NoError(t, executor.AwaitCommits()) +} + +// A reader whose pending block is retired, and then followed by further blocks that land, between +// its pending read and its view must not replay that retired block over the newer state, even though +// nothing is pending any more by the time it looks again. +func TestReadLatestAccountRestartsWhenItsPendingBlockRetiresUnderIt(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 := &retiringOnOpenStore{recordingGigaStore: &recordingGigaStore{snapshot: snapshot}} + executor := NewExecutor(Config{}, withTestStores(store, NewMemoryReceiptStore(), noopChangeSetEncoder)) + defer executor.Close() + + executePipelinedBlock(t, executor, chainID, 41, + signLegacyTx(t, key, chainID, 0, &recipient, big.NewInt(7), nil)) + + // The reader has block 41 in hand as pending; while its view opens, block 42 runs and both + // blocks land, leaving nothing pending and a store view that already holds them. + opens := 0 + store.retire = func() { + opens++ + if opens == 1 { + store.retire = nil + executePipelinedBlock(t, executor, chainID, 42, + signLegacyTx(t, key, chainID, 1, &recipient, big.NewInt(7), nil)) + require.NoError(t, executor.AwaitCommits()) + snapshot.nonces[sender] = 2 + } + } + got, err := executor.ReadLatestAccount(sender) + require.NoError(t, err) + require.Equal(t, uint64(2), got.Nonce, "the read must not replay retired block 41 over the landed state") +} + // An executor without a receipt store commits state only; the block result still carries receipts. func TestNoReceiptStoreCommitsStateOnly(t *testing.T) { chainID := big.NewInt(testChainID) diff --git a/sei-tendermint/internal/evmonlyapp/app.go b/sei-tendermint/internal/evmonlyapp/app.go index c86bba9161..f89d43d972 100644 --- a/sei-tendermint/internal/evmonlyapp/app.go +++ b/sei-tendermint/internal/evmonlyapp/app.go @@ -470,20 +470,34 @@ func (a *evmOnlyApplication) callBlockContext() (evmonly.BlockContext, error) { panic("unreachable") } -func (a *evmOnlyApplication) EvmNonce(address common.Address) uint64 { +// latestAccount returns address's balance and nonce after the last finalized block, read through +// the executor's in-flight commit rather than waiting for it. Before InitChain, or once a commit +// has failed, it reads the settled store instead. +func (a *evmOnlyApplication) latestAccount(address common.Address) evmonly.LatestAccount { + if executor, ok := a.settler.Load().Get(); ok { + if account, err := executor.ReadLatestAccount(address); err == nil { + return account + } + } snapshot := a.openSettledView() defer snapshot.Close() - return snapshot.GetNonce(evmOnlyStoreAddress(address)) + storeAddress := evmOnlyStoreAddress(address) + if !snapshot.AccountExists(storeAddress) { + return evmonly.LatestAccount{Balance: new(big.Int).Set(evmOnlyBaseBalance)} + } + balance := snapshot.GetBalance(storeAddress) + return evmonly.LatestAccount{ + Balance: new(big.Int).SetBytes(balance[:]), + Nonce: snapshot.GetNonce(storeAddress), + } +} + +func (a *evmOnlyApplication) EvmNonce(address common.Address) uint64 { + return a.latestAccount(address).Nonce } func (a *evmOnlyApplication) EvmBalance(address common.Address, _ []byte) uint256.Int { - snapshot := a.openSettledView() - defer snapshot.Close() - if !snapshot.AccountExists(evmOnlyStoreAddress(address)) { - return *uint256.MustFromBig(evmOnlyBaseBalance) - } - balance := snapshot.GetBalance(evmOnlyStoreAddress(address)) - return *new(uint256.Int).SetBytes(balance[:]) + return *uint256.MustFromBig(a.latestAccount(address).Balance) } func (a *evmOnlyApplication) EvmChainID() uint64 { diff --git a/sei-tendermint/internal/evmonlyapp/app_test.go b/sei-tendermint/internal/evmonlyapp/app_test.go index 2f35d5895c..74b9aeef2c 100644 --- a/sei-tendermint/internal/evmonlyapp/app_test.go +++ b/sei-tendermint/internal/evmonlyapp/app_test.go @@ -588,6 +588,57 @@ func TestEVMOnlyApplicationReadsSettleBehindFinalizeBlock(t *testing.T) { require.Equal(t, int64(4), latest) } +// Nonce and balance reads answer for the block FinalizeBlock just ran, before its commit lands, +// for every account it touched; accounts it did not touch keep the funded default. +func TestEVMOnlyApplicationNonceAndBalanceReflectTheFinalizedBlock(t *testing.T) { + app := newInitializedEVMOnlyTestApp(t) + key, err := crypto.GenerateKey() + require.NoError(t, err) + sender := crypto.PubkeyToAddress(key.PublicKey) + recipient := common.HexToAddress("0x1000000000000000000000000000000000000001") + untouched := common.HexToAddress("0x2000000000000000000000000000000000000002") + + response, err := app.FinalizeBlock(t.Context(), evmOnlyTestBlock(1, signedEVMOnlyTestTxFrom(t, key, evmOnlyTestChainID, 0))) + require.NoError(t, err) + require.Len(t, response.TxResults, 1) + + require.Equal(t, uint64(1), app.EvmNonce(sender)) + gasPaid := new(big.Int).Mul(big.NewInt(evmOnlyMinGasPrice), big.NewInt(response.TxResults[0].GasUsed)) + wantSender := new(big.Int).Sub(new(big.Int).Sub(new(big.Int).Set(evmOnlyBaseBalance), big.NewInt(1)), gasPaid) + senderBalance := app.EvmBalance(sender, nil) + require.Equal(t, wantSender, senderBalance.ToBig()) + recipientBalance := app.EvmBalance(recipient, nil) + require.Equal(t, new(big.Int).Add(new(big.Int).Set(evmOnlyBaseBalance), big.NewInt(1)), recipientBalance.ToBig()) + require.Equal(t, uint64(0), app.EvmNonce(untouched)) + untouchedBalance := app.EvmBalance(untouched, nil) + require.Equal(t, evmOnlyBaseBalance, untouchedBalance.ToBig()) + + _, err = app.Commit(t.Context()) + require.NoError(t, err) + require.Equal(t, uint64(1), app.EvmNonce(sender)) +} + +// Once a commit has failed, nonce and balance reads fall back to the store, which stays at the +// last version that landed. +func TestEVMOnlyApplicationNonceAndBalanceFallBackToTheStoreAfterAFailedCommit(t *testing.T) { + storage := openEVMOnlyTestStorage(t, t.TempDir()) + app, err := NewEVMOnlyApplication(evmOnlyTestChainID, nil, storage, unwritableEVMChangeSetEncoder) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, storage.Close()) }) + _, err = app.InitChain(evmOnlyTestInitChain()) + require.NoError(t, err) + settler, ok := app.(*evmOnlyApplication) + require.True(t, ok) + + raw, sender := signedEVMOnlyTestTx(t, evmOnlyTestChainID, 0) + finalizeAndCommitEVMOnlyTestBlock(t, app, evmOnlyTestBlock(1, raw)) + require.Error(t, settler.AwaitCommits()) + + require.Equal(t, uint64(0), app.EvmNonce(sender)) + balance := app.EvmBalance(sender, nil) + require.Equal(t, evmOnlyBaseBalance, balance.ToBig()) +} + // unwritableEVMChangeSetEncoder encodes every block with an EVM pair the store // refuses to apply, so the block executes and encodes cleanly and its commit is // the first thing that fails. diff --git a/sei-tendermint/internal/mempool/tx.go b/sei-tendermint/internal/mempool/tx.go index 7b53d436ba..2a325995ab 100644 --- a/sei-tendermint/internal/mempool/tx.go +++ b/sei-tendermint/internal/mempool/tx.go @@ -126,6 +126,9 @@ type txStoreInner struct { byEvmHash map[common.Hash]*WrappedTx byNonce map[evmAddrNonce]*WrappedTx accounts map[common.Address]*evmAccount + // Incremented whenever accounts is reset, so account state fetched outside the lock is only + // installed against the height it was fetched for. + accountsEpoch uint64 softLimit txCounter hardLimit txCounter @@ -203,6 +206,7 @@ func (s *txStore) Clear() { inner.byEvmHash = map[common.Hash]*WrappedTx{} inner.byNonce = map[evmAddrNonce]*WrappedTx{} inner.accounts = map[common.Address]*evmAccount{} + inner.accountsEpoch++ inner.state.Store(txStoreState{}) s.readyTxs.Clear() } @@ -362,7 +366,8 @@ func (s *txStore) insert(inner *txStoreInner, wtx *WrappedTx, recordAdded bool) // Fetch the evm account state. account, ok := inner.accounts[evm.address] if !ok { - // TODO(gprusak): consider whether we should move these queries out of the mutex. + // Insert prefetches the account outside the mutex; this only runs when that prefetch + // was discarded, or for txs reinserted by compact. b := s.app.EvmBalance(evm.address, evm.seiAddress) n := s.app.EvmNonce(evm.address) account = &evmAccount{b, n, n} @@ -506,10 +511,53 @@ func (inner *txStoreInner) inInclusionOrder() []*WrappedTx { return res } +// prefetchedAccount is evm account state fetched from the app outside the store lock, tagged with +// the accounts epoch it is valid for. +type prefetchedAccount struct { + address common.Address + epoch uint64 + account *evmAccount +} + +// prefetchAccount fetches the evm account state of wtx's sender without holding the store lock, when +// the store does not know the account yet. Reading the app can wait on the state store, and doing +// so under the lock would stall every other mempool operation. +func (s *txStore) prefetchAccount(wtx *WrappedTx) utils.Option[prefetchedAccount] { + evm, ok := wtx.evm.Get() + if !ok { + return utils.None[prefetchedAccount]() + } + var epoch uint64 + for inner := range s.inner.RLock() { + if _, ok := inner.accounts[evm.address]; ok { + return utils.None[prefetchedAccount]() + } + epoch = inner.accountsEpoch + } + b := s.app.EvmBalance(evm.address, evm.seiAddress) + n := s.app.EvmNonce(evm.address) + return utils.Some(prefetchedAccount{address: evm.address, epoch: epoch, account: &evmAccount{b, n, n}}) +} + +// installPrefetched adds a prefetched account to inner unless the accounts were reset since it was +// fetched or another inserter got there first. +func installPrefetched(inner *txStoreInner, prefetched utils.Option[prefetchedAccount]) { + p, ok := prefetched.Get() + if !ok || p.epoch != inner.accountsEpoch { + return + } + if _, ok := inner.accounts[p.address]; ok { + return + } + inner.accounts[p.address] = p.account +} + // Inserts a new transaction to txStore. // txStore takes ownership of wtx. func (s *txStore) Insert(wtx *WrappedTx) error { + prefetched := s.prefetchAccount(wtx) for inner := range s.inner.Lock() { + installPrefetched(inner, prefetched) if err := s.insert(inner, wtx, true); err != nil { return err } @@ -538,6 +586,7 @@ func (s *txStore) compact(inner *txStoreInner, clearAccounts bool) { inner.byNonce = map[evmAddrNonce]*WrappedTx{} if clearAccounts { inner.accounts = map[common.Address]*evmAccount{} + inner.accountsEpoch++ } for _, account := range inner.accounts { account.nextNonce = account.firstNonce diff --git a/sei-tendermint/internal/mempool/tx_test.go b/sei-tendermint/internal/mempool/tx_test.go index 79d7c87bd6..8b8a95b770 100644 --- a/sei-tendermint/internal/mempool/tx_test.go +++ b/sei-tendermint/internal/mempool/tx_test.go @@ -1,6 +1,7 @@ package mempool import ( + "context" "errors" "fmt" "testing" @@ -12,6 +13,7 @@ import ( "github.com/sei-protocol/sei-chain/sei-tendermint/internal/proxy" "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils" "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils/require" + "github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils/scope" "github.com/sei-protocol/sei-chain/sei-tendermint/types" ) @@ -789,3 +791,111 @@ func TestTxStore_InsertCompactionKeepsReadyListInSync(t *testing.T) { require.ElementsMatch(t, toTxs(expected), listed) } } + +// gatedNonceApp is an evmNonceApp whose nonce reads block until released, so a test can hold an +// account fetch in flight and observe what the store does meanwhile. +type gatedNonceApp struct { + *evmNonceApp + inFlight utils.AtomicSend[int] + release utils.AtomicRecv[bool] +} + +// EvmNonce reads the nonce first and only then blocks, so the value it returns is the one that +// was current when the fetch began. +func (a *gatedNonceApp) EvmNonce(addr common.Address) uint64 { + nonce := a.evmNonceApp.EvmNonce(addr) + a.inFlight.Store(a.inFlight.Load() + 1) + if _, err := a.release.Wait(context.Background(), func(v bool) bool { return v }); err != nil { + panic(err) + } + return nonce +} + +func evmTxForTest(rng utils.Rng, address common.Address, nonce uint64) *WrappedTx { + return &WrappedTx{ + hashedTx: newHashedTx(utils.GenBytes(rng, 32)), + timestamp: time.Now(), + priority: 1, + gasWanted: 1, + estimatedGas: 1, + evm: utils.Some(evmTx{ + address: address, + seiAddress: address.Bytes(), + hash: genEvmHash(rng), + nonce: nonce, + }), + } +} + +// Fetching a first-seen account from the app must not hold the store lock: other inserts and +// reads proceed while the fetch is in flight. +func TestTxStore_InsertFetchesFirstSeenAccountOutsideTheLock(t *testing.T) { + rng := utils.TestRng() + release := utils.NewAtomicSend(false) + app := &gatedNonceApp{evmNonceApp: newEVMNonceApp(), inFlight: utils.NewAtomicSend(0), release: release.Subscribe()} + txStore := NewTxStore(TestConfig(), proxy.New(app)) + cold := common.BytesToAddress(utils.GenBytes(rng, 20)) + warm := common.BytesToAddress(utils.GenBytes(rng, 20)) + + // warm is known to the store before any fetch is gated; its fetch does not count. + release.Store(true) + require.NoError(t, txStore.Insert(evmTxForTest(rng, warm, 0))) + release.Store(false) + app.inFlight.Store(0) + + coldTx := evmTxForTest(rng, cold, 0) + warmTx := evmTxForTest(rng, warm, 1) + require.NoError(t, scope.Run(t.Context(), func(ctx context.Context, s scope.Scope) error { + s.Spawn(func() error { return txStore.Insert(coldTx) }) + inFlight := app.inFlight.Subscribe() + if _, err := inFlight.Wait(ctx, func(n int) bool { return n == 1 }); err != nil { + return err + } + // Both a read and a write of the store complete while the cold fetch is blocked. + if got := txStore.NextNonce(warm); got != 1 { + return fmt.Errorf("NextNonce(warm) = %d, want 1", got) + } + if err := txStore.Insert(warmTx); err != nil { + return err + } + release.Store(true) + return nil + })) + _, ok := txStore.ByHash(coldTx.Hash()) + require.True(t, ok) + _, ok = txStore.ByHash(warmTx.Hash()) + require.True(t, ok) +} + +// Account state fetched before an Update belongs to the previous height and must not be installed +// after it: the insert is judged against the account nonce of the new height. +func TestTxStore_InsertDiscardsAccountFetchedBeforeUpdate(t *testing.T) { + rng := utils.TestRng() + release := utils.NewAtomicSend(false) + app := &gatedNonceApp{evmNonceApp: newEVMNonceApp(), inFlight: utils.NewAtomicSend(0), release: release.Subscribe()} + txStore := NewTxStore(TestConfig(), proxy.New(app)) + sender := common.BytesToAddress(utils.GenBytes(rng, 20)) + + staleTx := evmTxForTest(rng, sender, 0) + err := scope.Run(t.Context(), func(ctx context.Context, s scope.Scope) error { + s.Spawn(func() error { return txStore.Insert(staleTx) }) + inFlight := app.inFlight.Subscribe() + if _, err := inFlight.Wait(ctx, func(n int) bool { return n == 1 }); err != nil { + return err + } + // Nonce 0 is mined while the fetch that read it as pending is still in flight. + app.markMined(sender) + txStore.Update(updateSpec{ + Now: time.Now(), + Height: 1, + TxResults: map[types.TxHash]bool{}, + Constraints: NopTxConstraints(), + NewPriorities: map[types.TxHash]int64{}, + }) + release.Store(true) + return nil + }) + require.ErrorIs(t, err, errOldNonce) + _, ok := txStore.ByHash(staleTx.Hash()) + require.False(t, ok) +}