From 30dd821a987b1222e41fff02c1d700ab63218978 Mon Sep 17 00:00:00 2001 From: Brandon Chatham Date: Sat, 19 Sep 2026 21:05:42 +0000 Subject: [PATCH 1/3] evmonly: persist receipts and encode state changes behind the block; Commit waits for the receipts to land Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- giga/evmonly/executor.go | 22 ++- giga/evmonly/executor_test.go | 12 +- giga/evmonly/giga_store.go | 167 ++++++++++++------ giga/evmonly/giga_store_test.go | 9 +- giga/evmonly/pipeline_store_test.go | 109 ++++++++++++ .../ledger_db/receipt/litt_receipt_store.go | 27 ++- .../litt_write_failure_internal_test.go | 2 +- sei-db/ledger_db/receipt/littidx_test.go | 36 ++++ sei-tendermint/internal/evmonlyapp/app.go | 17 +- .../internal/evmonlyapp/app_test.go | 24 +++ 10 files changed, 343 insertions(+), 82 deletions(-) diff --git a/giga/evmonly/executor.go b/giga/evmonly/executor.go index ca03f6fb1f..0f34d97677 100644 --- a/giga/evmonly/executor.go +++ b/giga/evmonly/executor.go @@ -46,13 +46,20 @@ 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 into its stages. One block is persisted at a + // time, so one timer serves it. + pipelinePhases *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 - pipelineDone chan struct{} - pipelineErr error - pipelineChanges *pendingChanges + 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 // 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. @@ -98,9 +105,10 @@ func WithBlockChangeSetEncoder(encoder BlockChangeSetEncoder) Option { // execution on this executor. func NewExecutor(cfg Config, opts ...Option) *Executor { e := &Executor{ - cfg: cfg.WithDefaults(), - resultPool: newBlockResultPool(cfg.BlockResultPoolSize), - blockPhases: seidbmetrics.NewPhaseTimer(otel.Meter(executorMeterName), "evmonly_block"), + cfg: cfg.WithDefaults(), + resultPool: newBlockResultPool(cfg.BlockResultPoolSize), + blockPhases: seidbmetrics.NewPhaseTimer(otel.Meter(executorMeterName), "evmonly_block"), + pipelinePhases: seidbmetrics.NewPhaseTimer(otel.Meter(executorMeterName), "evmonly_pipeline"), } if e.cfg.OCCWorkers > 1 { e.occPool = newOCCWorkerPool(e.cfg.OCCWorkers) diff --git a/giga/evmonly/executor_test.go b/giga/evmonly/executor_test.go index 87aaac9178..4d682dcba6 100644 --- a/giga/evmonly/executor_test.go +++ b/giga/evmonly/executor_test.go @@ -174,19 +174,21 @@ func TestExecutorReturnsReceiptStoreError(t *testing.T) { require.ErrorIs(t, err, storeErr) require.Nil(t, result) - require.Empty(t, sink.results) + // Receipts are written behind the block, so the sink has already seen the result. + require.Len(t, sink.results, 1) + sink.releases[0]() 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.NoError(t, err) - require.NotNil(t, result) + require.ErrorIs(t, err, storeErr) + require.Nil(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 085077fa04..1c6709910c 100644 --- a/giga/evmonly/giga_store.go +++ b/giga/evmonly/giga_store.go @@ -10,6 +10,7 @@ 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,17 +27,12 @@ var ( var _ StateReader = gigaSnapshotStateReader{} // NamedChangeSetEncoder converts an executor-native state result into the -// 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. +// 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. // -// 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. +// It may read the store, which then holds every earlier block and none of this one; expanding a +// storage clear does so. type NamedChangeSetEncoder func(StateChangeSet) ([]*proto.NamedChangeSet, error) // BlockChangeSetEncoder contributes named changesets that are committed in the @@ -54,8 +50,7 @@ func (e *Executor) executePreparedBlockWithStore(ctx context.Context, req Prepar if stateStore == nil { return nil, errMissingStateStore } - receiptStore := e.receiptStore - if receiptStore == nil { + if e.receiptStore == nil { return nil, errMissingReceiptStore } if e.changeSetEncoder == nil { @@ -115,23 +110,24 @@ func (e *Executor) executePreparedBlockWithStore(ctx context.Context, req Prepar return nil, err } gigametrics.SetPhase(gigametrics.PhaseStorage) - // 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 { + // 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 { e.blockPhases.SetPhase("await_commit") if err := e.awaitPipelineCommit(); err != nil { return nil, err } } - 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) - } + // 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 if e.blockChangeSetEncoder != nil { - extra, err := e.blockChangeSetEncoder(req.Context, result) + e.blockPhases.SetPhase("encode_block_changesets") + 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) } @@ -143,45 +139,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 } - 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 := 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 { + if !settleBeforeBlockEncoder { e.blockPhases.SetPhase("await_commit") if err := e.awaitPipelineCommit(); err != nil { return nil, err } } - e.blockPhases.SetPhase("commit_state") - if err := e.startPipelineCommit(blockNumber, changesets, &result.ChangeSet); err != nil { + e.blockPhases.SetPhase("start_commit") + if err := e.startPipelineCommit(ctx, blockNumber, result, extra); err != nil { return nil, fmt.Errorf("commit state changes for block %d: %w", req.Context.Number, err) } ok = true return result, nil } -// 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. +// 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. // -// 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) +// 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 } // AwaitCommits blocks until every block this executor has run is committed, and reports the first @@ -271,7 +270,7 @@ func (e *Executor) pipelineFailureLocked() error { return e.pipelineFailure } if e.pipelineErr != nil { - return fmt.Errorf("commit state changes: %w", e.pipelineErr) + return fmt.Errorf("persist block: %w", e.pipelineErr) } return nil } @@ -298,7 +297,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("commit state changes: %w", e.pipelineErr) + e.pipelineFailure = fmt.Errorf("persist block: %w", e.pipelineErr) } e.pipelineDone = nil e.pipelineChanges = nil @@ -311,13 +310,17 @@ func (e *Executor) awaitPipelineCommit() error { return e.pipelineFailure } -// 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. +// startPipelineCommit persists the block in the background, receipts first and then the encoded +// state changes, and records what it changed, so the next block reads those changes through an +// overlay rather than waiting for the write. The block result is held until its receipts are +// encoded; the state changes are encoded from a copy that outlives it. // // 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 := newPendingChanges(changes.clone()) +func (e *Executor) startPipelineCommit(ctx context.Context, blockNumber int64, result *BlockResult, extra []*proto.NamedChangeSet) error { + changes := result.ChangeSet.clone() + pending := newPendingChanges(changes) + receipts := &receiptWrite{done: make(chan struct{})} done := make(chan struct{}) e.pipelineMu.Lock() if failure := e.pipelineFailure; failure != nil { @@ -327,19 +330,75 @@ func (e *Executor) startPipelineCommit(blockNumber int64, changesets []*proto.Na e.pipelineChanges = pending e.pipelineGeneration++ e.pipelineDone = done + e.pipelineReceipts = receipts e.pipelineErr = nil 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() { - err := e.stateStore.CommitStateChanges(blockNumber, changesets) + defer close(done) + defer e.pipelinePhases.Reset() + receipts.err = e.persistReceipts(bgCtx, blockNumber, result) + releaseResult() + close(receipts.done) + if receipts.err != nil { + e.pipelineMu.Lock() + e.pipelineErr = receipts.err + e.pipelineMu.Unlock() + return + } + err := e.commitStateChanges(blockNumber, changes, extra) 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.pipelinePhases.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.pipelinePhases.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.pipelinePhases.SetPhase("await_receipts") + 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 a29dd20075..8b0c27f1a9 100644 --- a/giga/evmonly/giga_store_test.go +++ b/giga/evmonly/giga_store_test.go @@ -338,7 +338,7 @@ func TestExecutorGigaStoreFailuresDoNotCommitPartialState(t *testing.T) { require.Equal(t, 1, snapshot.closeCount) }) - t.Run("context canceled during encoding", func(t *testing.T) { + t.Run("context canceled during encoding still persists the executed block", func(t *testing.T) { snapshot := newMemoryGigaSnapshot(0) store := &recordingGigaStore{snapshot: snapshot} ctx, cancel := context.WithCancel(t.Context()) @@ -349,9 +349,10 @@ func TestExecutorGigaStoreFailuresDoNotCommitPartialState(t *testing.T) { result, err := executor.ExecuteBlock(ctx, BlockRequest{Context: blockContext(big.NewInt(testChainID))}) - require.ErrorIs(t, err, context.Canceled) - require.Nil(t, result) - require.Empty(t, store.commits) + require.NoError(t, err) + require.NotNil(t, result) + result.Release() + require.Len(t, store.commits, 1) 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 c622affb71..43388d75f6 100644 --- a/giga/evmonly/pipeline_store_test.go +++ b/giga/evmonly/pipeline_store_test.go @@ -336,3 +336,112 @@ func executePipelinedBlock(t *testing.T, executor *Executor, chainID *big.Int, n 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) +} + +// 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 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) +} diff --git a/sei-db/ledger_db/receipt/litt_receipt_store.go b/sei-db/ledger_db/receipt/litt_receipt_store.go index e9a28e9b59..fb033b2c90 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,12 +99,16 @@ 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. + lastQueued atomic.Pointer[receiptWrite] } -// receiptWrite is one block's receipts, waiting to be applied. +// 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. 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 @@ -243,7 +247,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() @@ -349,7 +353,16 @@ 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}) + 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 + } } // ErrStoreClosed is returned by a write the store can no longer apply, the writer having stopped. @@ -357,7 +370,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() @@ -365,6 +378,7 @@ func (s *littReceiptStore) queueWrite(write receiptWrite) error { if s.closing { return ErrStoreClosed } + s.lastQueued.Store(write) seidbmetrics.Send(s.writeQueue, s.writes, write) return nil } @@ -579,7 +593,8 @@ 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) { +func (s *littReceiptStore) applyWrite(write *receiptWrite) { + defer close(write.landed) 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 e53242cd06..4386c941e1 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} + s.writes <- &receiptWrite{height: 1, landed: make(chan struct{})} } 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 4c9119fc53..37af9a04f4 100644 --- a/sei-db/ledger_db/receipt/littidx_test.go +++ b/sei-db/ledger_db/receipt/littidx_test.go @@ -13,6 +13,7 @@ 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" @@ -129,6 +130,41 @@ 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-tendermint/internal/evmonlyapp/app.go b/sei-tendermint/internal/evmonlyapp/app.go index 9de73b706e..9d5a87849b 100644 --- a/sei-tendermint/internal/evmonlyapp/app.go +++ b/sei-tendermint/internal/evmonlyapp/app.go @@ -822,12 +822,19 @@ func (a *evmOnlyApplication) pendingCursor(height int64) (evmOnlyCursor, error) panic("unreachable") } -// Commit acknowledges the finalized block as the one the chain builds on. 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 through the next FinalizeBlock, and a restart resumes -// from the store's own version. +// Commit acknowledges the finalized block as the one the chain builds on, once +// its receipts are readable: the height this advances is what RPC serves as +// latest, and a block it serves has its receipts. 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. 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 7eff0fa9ed..2dc7856ba4 100644 --- a/sei-tendermint/internal/evmonlyapp/app_test.go +++ b/sei-tendermint/internal/evmonlyapp/app_test.go @@ -196,6 +196,30 @@ 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) From ef8a7547593b6c5c73c50836e2b6b0b533b6c192 Mon Sep 17 00:00:00 2001 From: Brandon Chatham Date: Sat, 19 Sep 2026 21:17:43 +0000 Subject: [PATCH 2/3] receipt: publish the wait marker under the send; pebble store forwards WaitForPendingWrites; test a write that never lands Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- giga/evmonly/pipeline_store_test.go | 37 +++++++++++++++++++ .../ledger_db/receipt/litt_receipt_store.go | 6 ++- sei-db/ledger_db/receipt/receipt_store.go | 8 ++++ 3 files changed, 50 insertions(+), 1 deletion(-) diff --git a/giga/evmonly/pipeline_store_test.go b/giga/evmonly/pipeline_store_test.go index 43388d75f6..3abe38438f 100644 --- a/giga/evmonly/pipeline_store_test.go +++ b/giga/evmonly/pipeline_store_test.go @@ -8,6 +8,8 @@ import ( "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" ) @@ -410,6 +412,41 @@ func TestFailedReceiptWriteFailsTheBlockBeforeItsStateCommit(t *testing.T) { require.Empty(t, store.commits, "state must not be committed for a block whose receipts were not") } +// 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) { diff --git a/sei-db/ledger_db/receipt/litt_receipt_store.go b/sei-db/ledger_db/receipt/litt_receipt_store.go index fb033b2c90..eb471ece02 100644 --- a/sei-db/ledger_db/receipt/litt_receipt_store.go +++ b/sei-db/ledger_db/receipt/litt_receipt_store.go @@ -99,7 +99,9 @@ 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. + // 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] } @@ -378,6 +380,8 @@ 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 diff --git a/sei-db/ledger_db/receipt/receipt_store.go b/sei-db/ledger_db/receipt/receipt_store.go index 2b5ecb34d1..965a6e523f 100644 --- a/sei-db/ledger_db/receipt/receipt_store.go +++ b/sei-db/ledger_db/receipt/receipt_store.go @@ -273,6 +273,14 @@ 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) } From c6d9b4eed8f509bfb3c955e93c62e227bbc2fec4 Mon Sep 17 00:00:00 2001 From: Brandon Chatham Date: Sun, 20 Sep 2026 04:32:46 +0000 Subject: [PATCH 3/3] evmonly: describe the pipelined receipt/state writes in the README Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- giga/evmonly/README.md | 24 ++++++++++++++---------- 1 file changed, 14 insertions(+), 10 deletions(-) diff --git a/giga/evmonly/README.md b/giga/evmonly/README.md index 7a656d5299..64ce24e987 100644 --- a/giga/evmonly/README.md +++ b/giga/evmonly/README.md @@ -81,22 +81,26 @@ 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 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. +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. 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, state commit, or -receipt-store failures release the block result and return an error without -invoking `ResultSink`. Ethereum receipts are converted into +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 `receipt.ReceiptRecord` values and persisted through the shared `receipt.ReceiptStore` interface before the height-advancing state commit, -including for empty blocks. 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 succeed. +including for empty blocks; `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 can leave receipts behind, which re-executing the +block after a restart overwrites. `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