Skip to content

Commit 453a8a4

Browse files
authored
feat(tracing): add tracing to EngineClient (#2959)
<!-- Please read and fill out this form before submitting your PR. Please make sure you have reviewed our contributors guide before submitting your first PR. NOTE: PR titles should follow semantic commits: https://www.conventionalcommits.org/en/v1.0.0/ --> ## Overview This PR extracts a small interface with methods to call into `engine_forkchoiceUpdatedV3`, `engine_getPayloadV4` and `engine_newPayloadV4`. Creates an implementation with the existing logic, and then a tracing wrapper implementation. In this case, each of these will be child spans as they use the existing context which will have a parent span. We will get traces like ``` Executor.InitChain -> Engine.ForkchoiceUpdated Executor.ExecuteTxs -> Engine.GetPayload ``` <!-- Please provide an explanation of the PR, including the appropriate context, background, goal, and rationale. If there is an issue with this information, please provide a tl;dr and link the issue. Ex: Closes #<issue number> -->
1 parent 41cac58 commit 453a8a4

4 files changed

Lines changed: 208 additions & 21 deletions

File tree

execution/evm/engine_rpc_client.go

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,47 @@
1+
package evm
2+
3+
import (
4+
"context"
5+
6+
"github.com/ethereum/go-ethereum/beacon/engine"
7+
"github.com/ethereum/go-ethereum/rpc"
8+
)
9+
10+
var _ EngineRPCClient = (*engineRPCClient)(nil)
11+
12+
// engineRPCClient is the concrete implementation wrapping *rpc.Client.
13+
type engineRPCClient struct {
14+
client *rpc.Client
15+
}
16+
17+
// NewEngineRPCClient creates a new Engine API client.
18+
func NewEngineRPCClient(client *rpc.Client) EngineRPCClient {
19+
return &engineRPCClient{client: client}
20+
}
21+
22+
func (e *engineRPCClient) ForkchoiceUpdated(ctx context.Context, state engine.ForkchoiceStateV1, args map[string]any) (*engine.ForkChoiceResponse, error) {
23+
var result engine.ForkChoiceResponse
24+
err := e.client.CallContext(ctx, &result, "engine_forkchoiceUpdatedV3", state, args)
25+
if err != nil {
26+
return nil, err
27+
}
28+
return &result, nil
29+
}
30+
31+
func (e *engineRPCClient) GetPayload(ctx context.Context, payloadID engine.PayloadID) (*engine.ExecutionPayloadEnvelope, error) {
32+
var result engine.ExecutionPayloadEnvelope
33+
err := e.client.CallContext(ctx, &result, "engine_getPayloadV4", payloadID)
34+
if err != nil {
35+
return nil, err
36+
}
37+
return &result, nil
38+
}
39+
40+
func (e *engineRPCClient) NewPayload(ctx context.Context, payload *engine.ExecutableData, blobHashes []string, parentBeaconBlockRoot string, executionRequests [][]byte) (*engine.PayloadStatusV1, error) {
41+
var result engine.PayloadStatusV1
42+
err := e.client.CallContext(ctx, &result, "engine_newPayloadV4", payload, blobHashes, parentBeaconBlockRoot, executionRequests)
43+
if err != nil {
44+
return nil, err
45+
}
46+
return &result, nil
47+
}
Lines changed: 122 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,122 @@
1+
package evm
2+
3+
import (
4+
"context"
5+
6+
"github.com/ethereum/go-ethereum/beacon/engine"
7+
"go.opentelemetry.io/otel"
8+
"go.opentelemetry.io/otel/attribute"
9+
"go.opentelemetry.io/otel/codes"
10+
"go.opentelemetry.io/otel/trace"
11+
)
12+
13+
var _ EngineRPCClient = (*tracedEngineRPCClient)(nil)
14+
15+
// tracedEngineRPCClient wraps an EngineRPCClient and records spans.
16+
type tracedEngineRPCClient struct {
17+
inner EngineRPCClient
18+
tracer trace.Tracer
19+
}
20+
21+
// withTracingEngineRPCClient decorates an EngineRPCClient with OpenTelemetry spans.
22+
func withTracingEngineRPCClient(inner EngineRPCClient) EngineRPCClient {
23+
return &tracedEngineRPCClient{
24+
inner: inner,
25+
tracer: otel.Tracer("ev-node/execution/engine-rpc"),
26+
}
27+
}
28+
29+
func (t *tracedEngineRPCClient) ForkchoiceUpdated(ctx context.Context, state engine.ForkchoiceStateV1, args map[string]any) (*engine.ForkChoiceResponse, error) {
30+
ctx, span := t.tracer.Start(ctx, "Engine.ForkchoiceUpdated",
31+
trace.WithAttributes(
32+
attribute.String("method", "engine_forkchoiceUpdatedV3"),
33+
attribute.String("head_block_hash", state.HeadBlockHash.Hex()),
34+
attribute.String("safe_block_hash", state.SafeBlockHash.Hex()),
35+
attribute.String("finalized_block_hash", state.FinalizedBlockHash.Hex()),
36+
),
37+
)
38+
defer span.End()
39+
40+
result, err := t.inner.ForkchoiceUpdated(ctx, state, args)
41+
if err != nil {
42+
span.RecordError(err)
43+
span.SetStatus(codes.Error, err.Error())
44+
return nil, err
45+
}
46+
47+
attributes := []attribute.KeyValue{
48+
attribute.String("payload_status", result.PayloadStatus.Status),
49+
}
50+
51+
if result.PayloadID != nil {
52+
attributes = append(attributes, attribute.String("payload_id", result.PayloadID.String()))
53+
}
54+
55+
if result.PayloadStatus.LatestValidHash != nil {
56+
attributes = append(attributes, attribute.String("latest_valid_hash", result.PayloadStatus.LatestValidHash.Hex()))
57+
}
58+
59+
span.SetAttributes(
60+
attributes...,
61+
)
62+
63+
return result, nil
64+
}
65+
66+
func (t *tracedEngineRPCClient) GetPayload(ctx context.Context, payloadID engine.PayloadID) (*engine.ExecutionPayloadEnvelope, error) {
67+
ctx, span := t.tracer.Start(ctx, "Engine.GetPayload",
68+
trace.WithAttributes(
69+
attribute.String("method", "engine_getPayloadV4"),
70+
attribute.String("payload_id", payloadID.String()),
71+
),
72+
)
73+
defer span.End()
74+
75+
result, err := t.inner.GetPayload(ctx, payloadID)
76+
if err != nil {
77+
span.RecordError(err)
78+
span.SetStatus(codes.Error, err.Error())
79+
return nil, err
80+
}
81+
82+
span.SetAttributes(
83+
attribute.Int64("block_number", int64(result.ExecutionPayload.Number)),
84+
attribute.String("block_hash", result.ExecutionPayload.BlockHash.Hex()),
85+
attribute.String("state_root", result.ExecutionPayload.StateRoot.Hex()),
86+
attribute.Int("tx_count", len(result.ExecutionPayload.Transactions)),
87+
attribute.Int64("gas_used", int64(result.ExecutionPayload.GasUsed)),
88+
)
89+
90+
return result, nil
91+
}
92+
93+
func (t *tracedEngineRPCClient) NewPayload(ctx context.Context, payload *engine.ExecutableData, blobHashes []string, parentBeaconBlockRoot string, executionRequests [][]byte) (*engine.PayloadStatusV1, error) {
94+
ctx, span := t.tracer.Start(ctx, "Engine.NewPayload",
95+
trace.WithAttributes(
96+
attribute.String("method", "engine_newPayloadV4"),
97+
attribute.Int64("block_number", int64(payload.Number)),
98+
attribute.String("block_hash", payload.BlockHash.Hex()),
99+
attribute.String("parent_hash", payload.ParentHash.Hex()),
100+
attribute.Int("tx_count", len(payload.Transactions)),
101+
attribute.Int64("gas_used", int64(payload.GasUsed)),
102+
),
103+
)
104+
defer span.End()
105+
106+
result, err := t.inner.NewPayload(ctx, payload, blobHashes, parentBeaconBlockRoot, executionRequests)
107+
if err != nil {
108+
span.RecordError(err)
109+
span.SetStatus(codes.Error, err.Error())
110+
return nil, err
111+
}
112+
113+
attributes := []attribute.KeyValue{attribute.String("payload_status", result.Status)}
114+
115+
if result.LatestValidHash != nil {
116+
attributes = append(attributes, attribute.String("latest_valid_hash", result.LatestValidHash.Hex()))
117+
}
118+
119+
span.SetAttributes(attributes...)
120+
121+
return result, nil
122+
}

execution/evm/execution.go

Lines changed: 28 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -130,11 +130,23 @@ func retryWithBackoffOnPayloadStatus(ctx context.Context, fn func() error, maxRe
130130
return fmt.Errorf("max retries (%d) exceeded for %s", maxRetries, operation)
131131
}
132132

133+
// EngineRPCClient abstracts Engine API RPC calls for tracing and testing.
134+
type EngineRPCClient interface {
135+
// ForkchoiceUpdated updates the forkchoice state and optionally starts payload building.
136+
ForkchoiceUpdated(ctx context.Context, state engine.ForkchoiceStateV1, args map[string]any) (*engine.ForkChoiceResponse, error)
137+
138+
// GetPayload retrieves a previously requested execution payload.
139+
GetPayload(ctx context.Context, payloadID engine.PayloadID) (*engine.ExecutionPayloadEnvelope, error)
140+
141+
// NewPayload submits a new execution payload for validation.
142+
NewPayload(ctx context.Context, payload *engine.ExecutableData, blobHashes []string, parentBeaconBlockRoot string, executionRequests [][]byte) (*engine.PayloadStatusV1, error)
143+
}
144+
133145
// EngineClient represents a client that interacts with an Ethereum execution engine
134146
// through the Engine API. It manages connections to both the engine and standard Ethereum
135147
// APIs, and maintains state related to block processing.
136148
type EngineClient struct {
137-
engineClient *rpc.Client // Client for Engine API calls
149+
engineClient EngineRPCClient // Client for Engine API calls
138150
ethClient *ethclient.Client // Client for standard Ethereum API calls
139151
genesisHash common.Hash // Hash of the genesis block
140152
initialHeight uint64
@@ -206,11 +218,19 @@ func NewEngineExecutionClient(
206218
}
207219
return nil
208220
}))
209-
engineClient, err := rpc.DialOptions(context.Background(), engineURL, engineOptions...)
221+
rawEngineClient, err := rpc.DialOptions(context.Background(), engineURL, engineOptions...)
210222
if err != nil {
211223
return nil, err
212224
}
213225

226+
// raw engine client
227+
engineClient := NewEngineRPCClient(rawEngineClient)
228+
229+
// if tracing enabled, wrap with traced decorator
230+
if tracingEnabled {
231+
engineClient = withTracingEngineRPCClient(engineClient)
232+
}
233+
214234
return &EngineClient{
215235
engineClient: engineClient,
216236
ethClient: ethClient,
@@ -238,8 +258,7 @@ func (c *EngineClient) InitChain(ctx context.Context, genesisTime time.Time, ini
238258

239259
// Acknowledge the genesis block with retry logic for SYNCING status
240260
err := retryWithBackoffOnPayloadStatus(ctx, func() error {
241-
var forkchoiceResult engine.ForkChoiceResponse
242-
err := c.engineClient.CallContext(ctx, &forkchoiceResult, "engine_forkchoiceUpdatedV3",
261+
forkchoiceResult, err := c.engineClient.ForkchoiceUpdated(ctx,
243262
engine.ForkchoiceStateV1{
244263
HeadBlockHash: c.genesisHash,
245264
SafeBlockHash: c.genesisHash,
@@ -372,8 +391,7 @@ func (c *EngineClient) ExecuteTxs(ctx context.Context, txs [][]byte, blockHeight
372391
// 3. Call forkchoice update to get PayloadID
373392
var newPayloadID *engine.PayloadID
374393
err = retryWithBackoffOnPayloadStatus(ctx, func() error {
375-
var forkchoiceResult engine.ForkChoiceResponse
376-
err := c.engineClient.CallContext(ctx, &forkchoiceResult, "engine_forkchoiceUpdatedV3", args, evPayloadAttrs)
394+
forkchoiceResult, err := c.engineClient.ForkchoiceUpdated(ctx, args, evPayloadAttrs)
377395
if err != nil {
378396
return fmt.Errorf("forkchoice update failed: %w", err)
379397
}
@@ -522,8 +540,7 @@ func (c *EngineClient) setFinalWithHeight(ctx context.Context, blockHash common.
522540
func (c *EngineClient) doForkchoiceUpdate(ctx context.Context, args engine.ForkchoiceStateV1, operation string) error {
523541
// Call forkchoice update with retry logic for SYNCING status
524542
err := retryWithBackoffOnPayloadStatus(ctx, func() error {
525-
var forkchoiceResult engine.ForkChoiceResponse
526-
err := c.engineClient.CallContext(ctx, &forkchoiceResult, "engine_forkchoiceUpdatedV3", args, nil)
543+
forkchoiceResult, err := c.engineClient.ForkchoiceUpdated(ctx, args, nil)
527544
if err != nil {
528545
return fmt.Errorf("forkchoice update failed: %w", err)
529546
}
@@ -774,8 +791,7 @@ func (c *EngineClient) filterTransactions(ctx context.Context, txs [][]byte, blo
774791
// processPayload handles the common logic of getting, submitting, and finalizing a payload.
775792
func (c *EngineClient) processPayload(ctx context.Context, payloadID engine.PayloadID, txs [][]byte) ([]byte, uint64, error) {
776793
// 1. Get Payload
777-
var payloadResult engine.ExecutionPayloadEnvelope
778-
err := c.engineClient.CallContext(ctx, &payloadResult, "engine_getPayloadV4", payloadID)
794+
payloadResult, err := c.engineClient.GetPayload(ctx, payloadID)
779795
if err != nil {
780796
return nil, 0, fmt.Errorf("get payload failed: %w", err)
781797
}
@@ -784,9 +800,8 @@ func (c *EngineClient) processPayload(ctx context.Context, payloadID engine.Payl
784800
blockTimestamp := int64(payloadResult.ExecutionPayload.Timestamp)
785801

786802
// 2. Submit Payload (newPayload)
787-
var newPayloadResult engine.PayloadStatusV1
788803
err = retryWithBackoffOnPayloadStatus(ctx, func() error {
789-
err := c.engineClient.CallContext(ctx, &newPayloadResult, "engine_newPayloadV4",
804+
newPayloadResult, err := c.engineClient.NewPayload(ctx,
790805
payloadResult.ExecutionPayload,
791806
[]string{}, // No blob hashes
792807
common.Hash{}.Hex(), // Use zero hash for parentBeaconBlockRoot
@@ -796,7 +811,7 @@ func (c *EngineClient) processPayload(ctx context.Context, payloadID engine.Payl
796811
return fmt.Errorf("new payload submission failed: %w", err)
797812
}
798813

799-
if err := validatePayloadStatus(newPayloadResult); err != nil {
814+
if err := validatePayloadStatus(*newPayloadResult); err != nil {
800815
c.logger.Warn().
801816
Str("status", newPayloadResult.Status).
802817
Str("latestValidHash", latestValidHashHex(newPayloadResult.LatestValidHash)).

execution/evm/go.mod

Lines changed: 11 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,20 @@ require (
1010
github.com/ipfs/go-datastore v0.9.0
1111
github.com/rs/zerolog v1.34.0
1212
github.com/stretchr/testify v1.11.1
13+
go.opentelemetry.io/otel v1.39.0
14+
go.opentelemetry.io/otel/trace v1.39.0
1315
google.golang.org/protobuf v1.36.10
1416
)
1517

18+
require (
19+
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
20+
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0 // indirect
21+
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0 // indirect
22+
go.opentelemetry.io/otel/metric v1.39.0 // indirect
23+
go.opentelemetry.io/otel/sdk v1.38.0 // indirect
24+
go.opentelemetry.io/proto/otlp v1.7.1 // indirect
25+
)
26+
1627
require (
1728
github.com/DataDog/zstd v1.5.5 // indirect
1829
github.com/Microsoft/go-winio v0.6.2 // indirect
@@ -78,14 +89,6 @@ require (
7889
github.com/syndtr/goleveldb v1.0.1-0.20220721030215-126854af5e6d // indirect
7990
github.com/tklauser/go-sysconf v0.3.12 // indirect
8091
github.com/tklauser/numcpus v0.6.1 // indirect
81-
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
82-
go.opentelemetry.io/otel v1.39.0 // indirect
83-
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0 // indirect
84-
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0 // indirect
85-
go.opentelemetry.io/otel/metric v1.39.0 // indirect
86-
go.opentelemetry.io/otel/sdk v1.38.0 // indirect
87-
go.opentelemetry.io/otel/trace v1.39.0 // indirect
88-
go.opentelemetry.io/proto/otlp v1.7.1 // indirect
8992
go.yaml.in/yaml/v3 v3.0.4 // indirect
9093
golang.org/x/crypto v0.46.0 // indirect
9194
golang.org/x/exp v0.0.0-20251023183803-a4bb9ffd2546 // indirect

0 commit comments

Comments
 (0)