Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 9 additions & 7 deletions giga/evmonly/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -94,13 +94,15 @@ 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; `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.
`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.

`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
Expand Down
17 changes: 13 additions & 4 deletions giga/evmonly/executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ 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.
Expand Down Expand Up @@ -48,9 +49,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 into its stages. One block is persisted at a
// time, so one timer serves it.
// 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.
Expand Down Expand Up @@ -106,11 +111,13 @@ 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: seidbmetrics.NewPhaseTimer(otel.Meter(executorMeterName), "evmonly_pipeline"),
pipelinePhases: pipelineTimers.Build(attribute.String("stage", "state")),
receiptPhases: pipelineTimers.Build(attribute.String("stage", "receipts")),
}
e.parseSizer = newParseSizer(e.cfg.ParseWorkers)
if e.cfg.OCCWorkers > 1 {
Expand All @@ -128,8 +135,10 @@ 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.
// 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.
_ = e.awaitPipelineCommit()
_ = e.AwaitReceipts()
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
if e.occPool != nil {
e.occPool.Close()
}
Expand Down
70 changes: 52 additions & 18 deletions giga/evmonly/giga_store.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,10 @@ 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
Expand Down Expand Up @@ -150,7 +154,7 @@ func (e *Executor) executePreparedBlockWithStore(ctx context.Context, req Prepar
}
}
e.blockPhases.SetPhase("start_commit")
if err := e.startPipelineCommit(ctx, blockNumber, result, extra); err != nil {
if err := e.startPipelineCommit(blockNumber, result, receipts, extra); err != nil {
return nil, fmt.Errorf("commit state changes for block %d: %w", req.Context.Number, err)
}
ok = true
Expand Down Expand Up @@ -310,17 +314,51 @@ func (e *Executor) awaitPipelineCommit() error {
return e.pipelineFailure
}

// 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.
// 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.
//
// 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{})}
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)
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
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.
//
// Commits stay ordered because only one is ever in flight: awaitPipelineCommit lands the previous
// one before this is called.
func (e *Executor) startPipelineCommit(ctx context.Context, blockNumber int64, result *BlockResult, extra []*proto.NamedChangeSet) error {
func (e *Executor) startPipelineCommit(blockNumber int64, result *BlockResult, receipts *receiptWrite, 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 {
Expand All @@ -330,20 +368,16 @@ func (e *Executor) startPipelineCommit(ctx context.Context, blockNumber int64, r
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() {
defer close(done)
defer e.pipelinePhases.Reset()
receipts.err = e.persistReceipts(bgCtx, blockNumber, result)
releaseResult()
close(receipts.done)
// 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
Expand All @@ -361,19 +395,19 @@ func (e *Executor) startPipelineCommit(ctx context.Context, blockNumber int64, r
// 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")
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.pipelinePhases.SetPhase("write_receipts")
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.pipelinePhases.SetPhase("await_receipts")
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)
Expand Down
166 changes: 166 additions & 0 deletions giga/evmonly/pipeline_store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"errors"
"math/big"
"testing"
"time"

"github.com/ethereum/go-ethereum/crypto"
"github.com/stretchr/testify/require"
Expand Down Expand Up @@ -389,6 +390,74 @@ func TestAwaitReceiptsReturnsBeforeTheStateCommitLands(t *testing.T) {
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) {
Expand All @@ -412,6 +481,42 @@ func TestFailedReceiptWriteFailsTheBlockBeforeItsStateCommit(t *testing.T) {
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.
Expand Down Expand Up @@ -482,3 +587,64 @@ func TestBackgroundEncoderSeesTheBlocksOwnChanges(t *testing.T) {
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)
blockCtx := blockContext(chainID)
blockCtx.Number = 41
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, 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())
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")
}
Loading