Skip to content

Commit b6baa26

Browse files
committed
feat: handle execution proposer rotation
1 parent e5537b8 commit b6baa26

8 files changed

Lines changed: 922 additions & 11 deletions

‎execution/evm/eth_rpc_client.go‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,8 @@ import (
44
"context"
55
"math/big"
66

7+
"github.com/ethereum/go-ethereum/common"
8+
"github.com/ethereum/go-ethereum/common/hexutil"
79
"github.com/ethereum/go-ethereum/core/types"
810
"github.com/ethereum/go-ethereum/ethclient"
911
)
@@ -30,3 +32,19 @@ func (e *ethRPCClient) GetTxs(ctx context.Context) ([]string, error) {
3032
}
3133
return result, nil
3234
}
35+
36+
func (e *ethRPCClient) GetNextProposer(ctx context.Context, number *big.Int) (common.Hash, error) {
37+
var result common.Hash
38+
err := e.client.Client().CallContext(ctx, &result, "evolve_getNextProposer", blockNumberArg(number))
39+
if err != nil {
40+
return common.Hash{}, err
41+
}
42+
return result, nil
43+
}
44+
45+
func blockNumberArg(number *big.Int) string {
46+
if number == nil {
47+
return "latest"
48+
}
49+
return hexutil.EncodeBig(number)
50+
}

‎execution/evm/eth_rpc_tracing.go‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import (
44
"context"
55
"math/big"
66

7+
"github.com/ethereum/go-ethereum/common"
78
"github.com/ethereum/go-ethereum/core/types"
89
"go.opentelemetry.io/otel"
910
"go.opentelemetry.io/otel/attribute"
@@ -82,3 +83,29 @@ func (t *tracedEthRPCClient) GetTxs(ctx context.Context) ([]string, error) {
8283

8384
return result, nil
8485
}
86+
87+
func (t *tracedEthRPCClient) GetNextProposer(ctx context.Context, number *big.Int) (common.Hash, error) {
88+
blockNumber := "latest"
89+
if number != nil {
90+
blockNumber = number.String()
91+
}
92+
93+
ctx, span := t.tracer.Start(ctx, "Evolve.GetNextProposer",
94+
trace.WithAttributes(
95+
attribute.String("method", "evolve_getNextProposer"),
96+
attribute.String("block_number", blockNumber),
97+
),
98+
)
99+
defer span.End()
100+
101+
result, err := t.inner.GetNextProposer(ctx, number)
102+
if err != nil {
103+
span.RecordError(err)
104+
span.SetStatus(codes.Error, err.Error())
105+
return common.Hash{}, err
106+
}
107+
108+
span.SetAttributes(attribute.String("next_proposer", result.Hex()))
109+
110+
return result, nil
111+
}

‎execution/evm/eth_rpc_tracing_test.go‎

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import (
66
"math/big"
77
"testing"
88

9+
"github.com/ethereum/go-ethereum/common"
910
"github.com/ethereum/go-ethereum/core/types"
1011
"github.com/stretchr/testify/require"
1112
"go.opentelemetry.io/otel"
@@ -40,8 +41,9 @@ func setupTestEthRPCTracing(t *testing.T, mockClient EthRPCClient) (EthRPCClient
4041

4142
// mockEthRPCClient is a simple mock for testing
4243
type mockEthRPCClient struct {
43-
headerByNumberFn func(ctx context.Context, number *big.Int) (*types.Header, error)
44-
getTxsFn func(ctx context.Context) ([]string, error)
44+
headerByNumberFn func(ctx context.Context, number *big.Int) (*types.Header, error)
45+
getTxsFn func(ctx context.Context) ([]string, error)
46+
getNextProposerFn func(ctx context.Context, number *big.Int) (common.Hash, error)
4547
}
4648

4749
func (m *mockEthRPCClient) HeaderByNumber(ctx context.Context, number *big.Int) (*types.Header, error) {
@@ -58,6 +60,13 @@ func (m *mockEthRPCClient) GetTxs(ctx context.Context) ([]string, error) {
5860
return nil, nil
5961
}
6062

63+
func (m *mockEthRPCClient) GetNextProposer(ctx context.Context, number *big.Int) (common.Hash, error) {
64+
if m.getNextProposerFn != nil {
65+
return m.getNextProposerFn(ctx, number)
66+
}
67+
return common.Hash{}, nil
68+
}
69+
6170
func TestTracedEthRPCClient_HeaderByNumber_Success(t *testing.T) {
6271
expectedHeader := &types.Header{
6372
GasLimit: 30000000,

‎execution/evm/execution.go‎

Lines changed: 72 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -160,6 +160,9 @@ type EthRPCClient interface {
160160

161161
// GetTxs retrieves pending transactions from the transaction pool.
162162
GetTxs(ctx context.Context) ([]string, error)
163+
164+
// GetNextProposer retrieves the proposer selected by the execution layer.
165+
GetNextProposer(ctx context.Context, number *big.Int) (common.Hash, error)
163166
}
164167

165168
// EngineClient represents a client that interacts with an Ethereum execution engine
@@ -353,15 +356,15 @@ func (c *EngineClient) ExecuteTxs(ctx context.Context, txs [][]byte, blockHeight
353356
// Continue execution on error, as it might be transient
354357
} else if found {
355358
if stateRoot != nil {
356-
return execution.ExecuteResult{UpdatedStateRoot: stateRoot}, nil
359+
return c.executionResultAfterBlock(ctx, stateRoot, blockHeight)
357360
}
358361
if payloadID != nil {
359362
// Found in-progress execution, attempt to resume
360363
stateRoot, err := c.processPayload(ctx, *payloadID, txs)
361364
if err != nil {
362365
return execution.ExecuteResult{}, err
363366
}
364-
return execution.ExecuteResult{UpdatedStateRoot: stateRoot}, nil
367+
return c.executionResultAfterBlock(ctx, stateRoot, blockHeight)
365368
}
366369
}
367370

@@ -461,7 +464,7 @@ func (c *EngineClient) ExecuteTxs(ctx context.Context, txs [][]byte, blockHeight
461464
if err != nil {
462465
return execution.ExecuteResult{}, err
463466
}
464-
return execution.ExecuteResult{UpdatedStateRoot: stateRoot}, nil
467+
return c.executionResultAfterBlock(ctx, stateRoot, blockHeight)
465468
}
466469

467470
// setHead updates the head block hash without changing safe or finalized.
@@ -858,21 +861,81 @@ func (c *EngineClient) filterTransactions(txs [][]byte) []string {
858861

859862
// GetExecutionInfo returns current execution layer parameters.
860863
func (c *EngineClient) GetExecutionInfo(ctx context.Context) (execution.ExecutionInfo, error) {
864+
var info execution.ExecutionInfo
861865
if cached := c.cachedExecutionInfo.Load(); cached != nil {
862-
return *cached, nil
866+
info.MaxGas = cached.MaxGas
867+
} else {
868+
header, err := c.ethClient.HeaderByNumber(ctx, nil) // nil = latest
869+
if err != nil {
870+
return execution.ExecutionInfo{}, fmt.Errorf("failed to get latest block: %w", err)
871+
}
872+
873+
info.MaxGas = header.GasLimit
874+
c.cachedExecutionInfo.Store(&execution.ExecutionInfo{MaxGas: info.MaxGas})
863875
}
864876

865-
header, err := c.ethClient.HeaderByNumber(ctx, nil) // nil = latest
877+
nextProposer, err := c.ethClient.GetNextProposer(ctx, nil)
866878
if err != nil {
867-
return execution.ExecutionInfo{}, fmt.Errorf("failed to get latest block: %w", err)
879+
if !isRPCMethodNotFound(err) {
880+
return execution.ExecutionInfo{}, fmt.Errorf("failed to get next proposer: %w", err)
881+
}
882+
return info, nil
883+
}
884+
if nextProposer != (common.Hash{}) {
885+
info.NextProposerAddress = nextProposer.Bytes()
868886
}
869-
870-
info := execution.ExecutionInfo{MaxGas: header.GasLimit}
871-
c.cachedExecutionInfo.Store(&info)
872887

873888
return info, nil
874889
}
875890

891+
func isRPCMethodNotFound(err error) bool {
892+
var rpcErr rpc.Error
893+
return errors.As(err, &rpcErr) && rpcErr.ErrorCode() == -32601
894+
}
895+
896+
func (c *EngineClient) executionResultAfterBlock(ctx context.Context, stateRoot []byte, executedBlockHeight uint64) (execution.ExecuteResult, error) {
897+
nextProposer, err := c.nextProposerChangeAfterBlock(ctx, executedBlockHeight)
898+
if err != nil {
899+
return execution.ExecuteResult{}, err
900+
}
901+
return execution.ExecuteResult{
902+
UpdatedStateRoot: stateRoot,
903+
NextProposerAddress: nextProposer,
904+
}, nil
905+
}
906+
907+
func (c *EngineClient) nextProposerChangeAfterBlock(ctx context.Context, executedBlockHeight uint64) ([]byte, error) {
908+
if executedBlockHeight == 0 {
909+
return nil, nil
910+
}
911+
912+
before, supported, err := c.nextProposerAtBlock(ctx, executedBlockHeight-1)
913+
if err != nil || !supported {
914+
return nil, err
915+
}
916+
917+
after, supported, err := c.nextProposerAtBlock(ctx, executedBlockHeight)
918+
if err != nil || !supported {
919+
return nil, err
920+
}
921+
922+
if after == (common.Hash{}) || after == before {
923+
return nil, nil
924+
}
925+
return after.Bytes(), nil
926+
}
927+
928+
func (c *EngineClient) nextProposerAtBlock(ctx context.Context, blockHeight uint64) (common.Hash, bool, error) {
929+
proposer, err := c.ethClient.GetNextProposer(ctx, new(big.Int).SetUint64(blockHeight))
930+
if err != nil {
931+
if isRPCMethodNotFound(err) {
932+
return common.Hash{}, false, nil
933+
}
934+
return common.Hash{}, true, fmt.Errorf("failed to get next proposer at block %d: %w", blockHeight, err)
935+
}
936+
return proposer, true, nil
937+
}
938+
876939
// FilterTxs validates force-included transactions and applies gas and size filtering for all passed txs.
877940
// If hasForceIncludedTransaction is false, skip filtering entirely - mempool batch is already filtered.
878941
// Returns a slice of FilterStatus for each transaction.

‎execution/evm/execution_reconcile_test.go‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import (
88
"time"
99

1010
"github.com/ethereum/go-ethereum/beacon/engine"
11+
"github.com/ethereum/go-ethereum/common"
1112
"github.com/ethereum/go-ethereum/core/types"
1213
ds "github.com/ipfs/go-datastore"
1314
dssync "github.com/ipfs/go-datastore/sync"
@@ -123,3 +124,7 @@ func (mockReconcileEthRPCClient) HeaderByNumber(_ context.Context, _ *big.Int) (
123124
func (mockReconcileEthRPCClient) GetTxs(_ context.Context) ([]string, error) {
124125
return nil, errors.New("unexpected GetTxs call")
125126
}
127+
128+
func (mockReconcileEthRPCClient) GetNextProposer(_ context.Context, _ *big.Int) (common.Hash, error) {
129+
return common.Hash{}, errors.New("unexpected GetNextProposer call")
130+
}

‎execution/evm/proposer_test.go‎

Lines changed: 132 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,132 @@
1+
package evm
2+
3+
import (
4+
"context"
5+
"math/big"
6+
"testing"
7+
"time"
8+
9+
"github.com/ethereum/go-ethereum/beacon/engine"
10+
"github.com/ethereum/go-ethereum/common"
11+
"github.com/ethereum/go-ethereum/core/types"
12+
ds "github.com/ipfs/go-datastore"
13+
dssync "github.com/ipfs/go-datastore/sync"
14+
"github.com/rs/zerolog"
15+
"github.com/stretchr/testify/require"
16+
)
17+
18+
type proposerEthRPCClient struct {
19+
headerByNumberFn func(ctx context.Context, number *big.Int) (*types.Header, error)
20+
getTxsFn func(ctx context.Context) ([]string, error)
21+
getNextProposerFn func(ctx context.Context, number *big.Int) (common.Hash, error)
22+
nextProposerBlocks []*big.Int
23+
}
24+
25+
func (m *proposerEthRPCClient) HeaderByNumber(ctx context.Context, number *big.Int) (*types.Header, error) {
26+
if m.headerByNumberFn != nil {
27+
return m.headerByNumberFn(ctx, number)
28+
}
29+
return &types.Header{GasLimit: 30_000_000}, nil
30+
}
31+
32+
func (m *proposerEthRPCClient) GetTxs(ctx context.Context) ([]string, error) {
33+
if m.getTxsFn != nil {
34+
return m.getTxsFn(ctx)
35+
}
36+
return nil, nil
37+
}
38+
39+
func (m *proposerEthRPCClient) GetNextProposer(ctx context.Context, number *big.Int) (common.Hash, error) {
40+
m.nextProposerBlocks = append(m.nextProposerBlocks, number)
41+
if m.getNextProposerFn != nil {
42+
return m.getNextProposerFn(ctx, number)
43+
}
44+
return common.Hash{}, nil
45+
}
46+
47+
func TestGetExecutionInfoIncludesNextProposer(t *testing.T) {
48+
nextProposer := common.HexToHash("0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa")
49+
ethClient := &proposerEthRPCClient{
50+
getNextProposerFn: func(ctx context.Context, number *big.Int) (common.Hash, error) {
51+
require.Nil(t, number)
52+
return nextProposer, nil
53+
},
54+
}
55+
client := &EngineClient{ethClient: ethClient}
56+
57+
info, err := client.GetExecutionInfo(t.Context())
58+
59+
require.NoError(t, err)
60+
require.Equal(t, uint64(30_000_000), info.MaxGas)
61+
require.Equal(t, nextProposer.Bytes(), info.NextProposerAddress)
62+
require.Len(t, ethClient.nextProposerBlocks, 1)
63+
}
64+
65+
func TestExecuteTxsReturnsNextProposerWhenChanged(t *testing.T) {
66+
timestamp := time.Unix(1_700_000_000, 0)
67+
stateRoot := common.HexToHash("0x01")
68+
prevProposer := common.HexToHash("0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa")
69+
nextProposer := common.HexToHash("0xbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb")
70+
71+
store := NewEVMStore(dssync.MutexWrap(ds.NewMapDatastore()))
72+
require.NoError(t, store.SaveExecMeta(t.Context(), &ExecMeta{
73+
Height: 2,
74+
StateRoot: stateRoot.Bytes(),
75+
Timestamp: timestamp.Unix(),
76+
Stage: ExecStagePromoted,
77+
}))
78+
79+
ethClient := &proposerEthRPCClient{
80+
headerByNumberFn: func(ctx context.Context, number *big.Int) (*types.Header, error) {
81+
require.Equal(t, int64(2), number.Int64())
82+
return &types.Header{
83+
Number: big.NewInt(2),
84+
Time: uint64(timestamp.Unix()),
85+
Root: stateRoot,
86+
GasLimit: 30_000_000,
87+
}, nil
88+
},
89+
getNextProposerFn: func(ctx context.Context, number *big.Int) (common.Hash, error) {
90+
switch number.Uint64() {
91+
case 1:
92+
return prevProposer, nil
93+
case 2:
94+
return nextProposer, nil
95+
default:
96+
t.Fatalf("unexpected proposer block %v", number)
97+
return common.Hash{}, nil
98+
}
99+
},
100+
}
101+
client := &EngineClient{
102+
engineClient: proposerEngineRPCClient{},
103+
ethClient: ethClient,
104+
store: store,
105+
currentSafeBlockHash: common.HexToHash("0x10"),
106+
currentFinalizedBlockHash: common.HexToHash("0x10"),
107+
logger: zerolog.Nop(),
108+
}
109+
110+
result, err := client.ExecuteTxs(t.Context(), nil, 2, timestamp, nil)
111+
112+
require.NoError(t, err)
113+
require.Equal(t, stateRoot.Bytes(), result.UpdatedStateRoot)
114+
require.Equal(t, nextProposer.Bytes(), result.NextProposerAddress)
115+
require.Len(t, ethClient.nextProposerBlocks, 2)
116+
}
117+
118+
type proposerEngineRPCClient struct{}
119+
120+
func (proposerEngineRPCClient) ForkchoiceUpdated(context.Context, engine.ForkchoiceStateV1, map[string]any) (*engine.ForkChoiceResponse, error) {
121+
return &engine.ForkChoiceResponse{
122+
PayloadStatus: engine.PayloadStatusV1{Status: engine.VALID},
123+
}, nil
124+
}
125+
126+
func (proposerEngineRPCClient) GetPayload(context.Context, engine.PayloadID) (*engine.ExecutionPayloadEnvelope, error) {
127+
return nil, nil
128+
}
129+
130+
func (proposerEngineRPCClient) NewPayload(context.Context, *engine.ExecutableData, []string, string, [][]byte) (*engine.PayloadStatusV1, error) {
131+
return nil, nil
132+
}

0 commit comments

Comments
 (0)