mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-21 12:16:44 +00:00
engine: add minimal OpenTelemetry tracing for newPayload
This commit is contained in:
parent
a9acb3ff93
commit
3ca77216c3
4 changed files with 91 additions and 30 deletions
|
|
@ -18,6 +18,7 @@
|
|||
package catalyst
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"reflect"
|
||||
|
|
@ -43,6 +44,11 @@ import (
|
|||
"github.com/ethereum/go-ethereum/params/forks"
|
||||
"github.com/ethereum/go-ethereum/rlp"
|
||||
"github.com/ethereum/go-ethereum/rpc"
|
||||
|
||||
"go.opentelemetry.io/otel"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/codes"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
)
|
||||
|
||||
// Register adds the engine API to the full node.
|
||||
|
|
@ -624,16 +630,32 @@ func (api *ConsensusAPI) getBlobs(hashes []common.Hash, v2 bool) ([]*engine.Blob
|
|||
// Helper for NewPayload* methods.
|
||||
var invalidStatus = engine.PayloadStatusV1{Status: engine.INVALID}
|
||||
|
||||
// startNewPayloadSpan starts a tracing span for a new payload.
|
||||
func startNewPayloadSpan(ctx context.Context, name string, params engine.ExecutableData) (context.Context, trace.Span) {
|
||||
tracer := otel.Tracer("")
|
||||
ctx, span := tracer.Start(ctx, name)
|
||||
span.SetAttributes(
|
||||
attribute.Int64("block.number", int64(params.Number)),
|
||||
attribute.String("block.hash", params.BlockHash.Hex()),
|
||||
attribute.Int("tx.count", len(params.Transactions)),
|
||||
)
|
||||
return ctx, span
|
||||
}
|
||||
|
||||
// NewPayloadV1 creates an Eth1 block, inserts it in the chain, and returns the status of the chain.
|
||||
func (api *ConsensusAPI) NewPayloadV1(params engine.ExecutableData) (engine.PayloadStatusV1, error) {
|
||||
func (api *ConsensusAPI) NewPayloadV1(ctx context.Context, params engine.ExecutableData) (engine.PayloadStatusV1, error) {
|
||||
ctx, span := startNewPayloadSpan(ctx, "engine.newPayloadV1", params)
|
||||
defer span.End()
|
||||
if params.Withdrawals != nil {
|
||||
return invalidStatus, paramsErr("withdrawals not supported in V1")
|
||||
}
|
||||
return api.newPayload(params, nil, nil, nil, false)
|
||||
return api.newPayload(ctx, params, nil, nil, nil, false)
|
||||
}
|
||||
|
||||
// NewPayloadV2 creates an Eth1 block, inserts it in the chain, and returns the status of the chain.
|
||||
func (api *ConsensusAPI) NewPayloadV2(params engine.ExecutableData) (engine.PayloadStatusV1, error) {
|
||||
func (api *ConsensusAPI) NewPayloadV2(ctx context.Context, params engine.ExecutableData) (engine.PayloadStatusV1, error) {
|
||||
ctx, span := startNewPayloadSpan(ctx, "engine.newPayloadV2", params)
|
||||
defer span.End()
|
||||
var (
|
||||
cancun = api.config().IsCancun(api.config().LondonBlock, params.Timestamp)
|
||||
shanghai = api.config().IsShanghai(api.config().LondonBlock, params.Timestamp)
|
||||
|
|
@ -650,11 +672,13 @@ func (api *ConsensusAPI) NewPayloadV2(params engine.ExecutableData) (engine.Payl
|
|||
case params.BlobGasUsed != nil:
|
||||
return invalidStatus, paramsErr("non-nil blobGasUsed pre-cancun")
|
||||
}
|
||||
return api.newPayload(params, nil, nil, nil, false)
|
||||
return api.newPayload(ctx, params, nil, nil, nil, false)
|
||||
}
|
||||
|
||||
// NewPayloadV3 creates an Eth1 block, inserts it in the chain, and returns the status of the chain.
|
||||
func (api *ConsensusAPI) NewPayloadV3(params engine.ExecutableData, versionedHashes []common.Hash, beaconRoot *common.Hash) (engine.PayloadStatusV1, error) {
|
||||
func (api *ConsensusAPI) NewPayloadV3(ctx context.Context, params engine.ExecutableData, versionedHashes []common.Hash, beaconRoot *common.Hash) (engine.PayloadStatusV1, error) {
|
||||
ctx, span := startNewPayloadSpan(ctx, "engine.newPayloadV3", params)
|
||||
defer span.End()
|
||||
switch {
|
||||
case params.Withdrawals == nil:
|
||||
return invalidStatus, paramsErr("nil withdrawals post-shanghai")
|
||||
|
|
@ -669,11 +693,13 @@ func (api *ConsensusAPI) NewPayloadV3(params engine.ExecutableData, versionedHas
|
|||
case !api.checkFork(params.Timestamp, forks.Cancun):
|
||||
return invalidStatus, unsupportedForkErr("newPayloadV3 must only be called for cancun payloads")
|
||||
}
|
||||
return api.newPayload(params, versionedHashes, beaconRoot, nil, false)
|
||||
return api.newPayload(ctx, params, versionedHashes, beaconRoot, nil, false)
|
||||
}
|
||||
|
||||
// NewPayloadV4 creates an Eth1 block, inserts it in the chain, and returns the status of the chain.
|
||||
func (api *ConsensusAPI) NewPayloadV4(params engine.ExecutableData, versionedHashes []common.Hash, beaconRoot *common.Hash, executionRequests []hexutil.Bytes) (engine.PayloadStatusV1, error) {
|
||||
func (api *ConsensusAPI) NewPayloadV4(ctx context.Context, params engine.ExecutableData, versionedHashes []common.Hash, beaconRoot *common.Hash, executionRequests []hexutil.Bytes) (engine.PayloadStatusV1, error) {
|
||||
ctx, span := startNewPayloadSpan(ctx, "engine.newPayloadV4", params)
|
||||
defer span.End()
|
||||
switch {
|
||||
case params.Withdrawals == nil:
|
||||
return invalidStatus, paramsErr("nil withdrawals post-shanghai")
|
||||
|
|
@ -694,10 +720,10 @@ func (api *ConsensusAPI) NewPayloadV4(params engine.ExecutableData, versionedHas
|
|||
if err := validateRequests(requests); err != nil {
|
||||
return engine.PayloadStatusV1{Status: engine.INVALID}, engine.InvalidParams.With(err)
|
||||
}
|
||||
return api.newPayload(params, versionedHashes, beaconRoot, requests, false)
|
||||
return api.newPayload(ctx, params, versionedHashes, beaconRoot, requests, false)
|
||||
}
|
||||
|
||||
func (api *ConsensusAPI) newPayload(params engine.ExecutableData, versionedHashes []common.Hash, beaconRoot *common.Hash, requests [][]byte, witness bool) (engine.PayloadStatusV1, error) {
|
||||
func (api *ConsensusAPI) newPayload(ctx context.Context, params engine.ExecutableData, versionedHashes []common.Hash, beaconRoot *common.Hash, requests [][]byte, witness bool) (engine.PayloadStatusV1, error) {
|
||||
// The locking here is, strictly, not required. Without these locks, this can happen:
|
||||
//
|
||||
// 1. NewPayload( execdata-N ) is invoked from the CL. It goes all the way down to
|
||||
|
|
@ -715,7 +741,17 @@ func (api *ConsensusAPI) newPayload(params engine.ExecutableData, versionedHashe
|
|||
defer api.newPayloadLock.Unlock()
|
||||
|
||||
log.Trace("Engine API request received", "method", "NewPayload", "number", params.Number, "hash", params.BlockHash)
|
||||
|
||||
tracer := otel.Tracer("")
|
||||
rootSpan := trace.SpanFromContext(ctx)
|
||||
_, span := tracer.Start(ctx, "engine.newPayload.ExecutableDataToBlock")
|
||||
block, err := engine.ExecutableDataToBlock(params, versionedHashes, beaconRoot, requests)
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
rootSpan.SetStatus(codes.Error, err.Error())
|
||||
}
|
||||
span.End()
|
||||
if err != nil {
|
||||
bgu := "nil"
|
||||
if params.BlobGasUsed != nil {
|
||||
|
|
@ -789,8 +825,15 @@ func (api *ConsensusAPI) newPayload(params engine.ExecutableData, versionedHashe
|
|||
}
|
||||
log.Trace("Inserting block without sethead", "hash", block.Hash(), "number", block.Number())
|
||||
start := time.Now()
|
||||
_, span = tracer.Start(ctx, "engine.newPayload.InsertBlockWithoutSetHead")
|
||||
proofs, err := api.eth.BlockChain().InsertBlockWithoutSetHead(block, witness)
|
||||
processingTime := time.Since(start)
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
rootSpan.SetStatus(codes.Error, err.Error())
|
||||
}
|
||||
span.End()
|
||||
if err != nil {
|
||||
log.Warn("NewPayload: inserting block failed", "error", err)
|
||||
|
||||
|
|
|
|||
|
|
@ -314,7 +314,7 @@ func TestEth2NewBlock(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("Failed to convert executable data to block %v", err)
|
||||
}
|
||||
newResp, err := api.NewPayloadV1(*execData)
|
||||
newResp, err := api.NewPayloadV1(context.Background(), *execData)
|
||||
switch {
|
||||
case err != nil:
|
||||
t.Fatalf("Failed to insert block: %v", err)
|
||||
|
|
@ -356,7 +356,7 @@ func TestEth2NewBlock(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("Failed to convert executable data to block %v", err)
|
||||
}
|
||||
newResp, err := api.NewPayloadV1(*execData)
|
||||
newResp, err := api.NewPayloadV1(context.Background(), *execData)
|
||||
if err != nil || newResp.Status != "VALID" {
|
||||
t.Fatalf("Failed to insert block: %v", err)
|
||||
}
|
||||
|
|
@ -502,7 +502,12 @@ func setupBlocks(t *testing.T, ethservice *eth.Ethereum, n int, parent *types.He
|
|||
}
|
||||
|
||||
envelope := getNewEnvelope(t, api, parent, w, h)
|
||||
execResp, err := api.newPayload(*envelope.ExecutionPayload, []common.Hash{}, h, envelope.Requests, false)
|
||||
|
||||
// NOTE: This span is for the test harness only. Engine oot spans are created
|
||||
// in NewPayloadV* entrypoints. This test calls newPayload() directly.
|
||||
ctx, span := startNewPayloadSpan(context.Background(), "engine.api_test.setupBlocks", *envelope.ExecutionPayload)
|
||||
defer span.End()
|
||||
execResp, err := api.newPayload(ctx, *envelope.ExecutionPayload, []common.Hash{}, h, envelope.Requests, false)
|
||||
if err != nil {
|
||||
t.Fatalf("can't execute payload: %v", err)
|
||||
}
|
||||
|
|
@ -648,7 +653,7 @@ func TestNewPayloadOnInvalidChain(t *testing.T) {
|
|||
t.Fatalf("payload should not be empty")
|
||||
}
|
||||
}
|
||||
execResp, err := api.NewPayloadV1(*payload.ExecutionPayload)
|
||||
execResp, err := api.NewPayloadV1(context.Background(), *payload.ExecutionPayload)
|
||||
if err != nil {
|
||||
t.Fatalf("can't execute payload: %v", err)
|
||||
}
|
||||
|
|
@ -708,7 +713,7 @@ func TestEmptyBlocks(t *testing.T) {
|
|||
// (1) check LatestValidHash by sending a normal payload (P1'')
|
||||
payload := getNewPayload(t, api, commonAncestor, nil, nil)
|
||||
|
||||
status, err := api.NewPayloadV1(*payload)
|
||||
status, err := api.NewPayloadV1(context.Background(), *payload)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -724,7 +729,7 @@ func TestEmptyBlocks(t *testing.T) {
|
|||
payload.GasUsed += 1
|
||||
payload = setBlockhash(payload)
|
||||
// Now latestValidHash should be the common ancestor
|
||||
status, err = api.NewPayloadV1(*payload)
|
||||
status, err = api.NewPayloadV1(context.Background(), *payload)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -742,7 +747,7 @@ func TestEmptyBlocks(t *testing.T) {
|
|||
payload.ParentHash = common.Hash{1}
|
||||
payload = setBlockhash(payload)
|
||||
// Now latestValidHash should be the common ancestor
|
||||
status, err = api.NewPayloadV1(*payload)
|
||||
status, err = api.NewPayloadV1(context.Background(), *payload)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -859,7 +864,7 @@ func TestTrickRemoteBlockCache(t *testing.T) {
|
|||
|
||||
// feed the payloads to node B
|
||||
for _, payload := range invalidChain {
|
||||
status, err := apiB.NewPayloadV1(*payload)
|
||||
status, err := apiB.NewPayloadV1(context.Background(), *payload)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
|
@ -892,7 +897,7 @@ func TestInvalidBloom(t *testing.T) {
|
|||
// (1) check LatestValidHash by sending a normal payload (P1'')
|
||||
payload := getNewPayload(t, api, commonAncestor, nil, nil)
|
||||
payload.LogsBloom = append(payload.LogsBloom, byte(1))
|
||||
status, err := api.NewPayloadV1(*payload)
|
||||
status, err := api.NewPayloadV1(context.Background(), *payload)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -931,7 +936,7 @@ func TestSimultaneousNewBlock(t *testing.T) {
|
|||
for ii := 0; ii < 10; ii++ {
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
if newResp, err := api.NewPayloadV1(*execData); err != nil {
|
||||
if newResp, err := api.NewPayloadV1(context.Background(), *execData); err != nil {
|
||||
errMu.Lock()
|
||||
testErr = fmt.Errorf("failed to insert block: %w", err)
|
||||
errMu.Unlock()
|
||||
|
|
@ -1038,7 +1043,7 @@ func TestWithdrawals(t *testing.T) {
|
|||
}
|
||||
|
||||
// 10: verify locally built block
|
||||
if status, err := api.NewPayloadV2(*execData.ExecutionPayload); err != nil {
|
||||
if status, err := api.NewPayloadV2(context.Background(), *execData.ExecutionPayload); err != nil {
|
||||
t.Fatalf("error validating payload: %v", err)
|
||||
} else if status.Status != engine.VALID {
|
||||
t.Fatalf("invalid payload")
|
||||
|
|
@ -1082,7 +1087,7 @@ func TestWithdrawals(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("error getting payload, err=%v", err)
|
||||
}
|
||||
if status, err := api.NewPayloadV2(*execData.ExecutionPayload); err != nil {
|
||||
if status, err := api.NewPayloadV2(context.Background(), *execData.ExecutionPayload); err != nil {
|
||||
t.Fatalf("error validating payload: %v", err)
|
||||
} else if status.Status != engine.VALID {
|
||||
t.Fatalf("invalid payload")
|
||||
|
|
@ -1225,9 +1230,9 @@ func TestNilWithdrawals(t *testing.T) {
|
|||
}
|
||||
var status engine.PayloadStatusV1
|
||||
if !shanghai {
|
||||
status, err = api.NewPayloadV1(*execData.ExecutionPayload)
|
||||
status, err = api.NewPayloadV1(context.Background(), *execData.ExecutionPayload)
|
||||
} else {
|
||||
status, err = api.NewPayloadV2(*execData.ExecutionPayload)
|
||||
status, err = api.NewPayloadV2(context.Background(), *execData.ExecutionPayload)
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatalf("error validating payload: %v", err.(*engine.EngineAPIError).ErrorData())
|
||||
|
|
@ -1598,7 +1603,7 @@ func TestParentBeaconBlockRoot(t *testing.T) {
|
|||
}
|
||||
|
||||
// 11: verify locally built block
|
||||
if status, err := api.NewPayloadV3(*execData.ExecutionPayload, []common.Hash{}, &common.Hash{42}); err != nil {
|
||||
if status, err := api.NewPayloadV3(context.Background(), *execData.ExecutionPayload, []common.Hash{}, &common.Hash{42}); err != nil {
|
||||
t.Fatalf("error validating payload: %v", err)
|
||||
} else if status.Status != engine.VALID {
|
||||
t.Fatalf("invalid payload")
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@
|
|||
package catalyst
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"crypto/sha256"
|
||||
"errors"
|
||||
|
|
@ -255,8 +256,11 @@ func (c *SimulatedBeacon) sealBlock(withdrawals []*types.Withdrawal, timestamp u
|
|||
requests = envelope.Requests
|
||||
}
|
||||
|
||||
// Mark the payload as canon
|
||||
_, err = c.engineAPI.newPayload(*payload, blobHashes, beaconRoot, requests, false)
|
||||
// NOTE: This span is for the simulated beacon harness only. Normal tracing
|
||||
// of engine_newPayload* is performed at the Engine API entrypoints.
|
||||
ctx, span := startNewPayloadSpan(context.Background(), "engine.simulatedBeacon.sealBlock", *payload)
|
||||
defer span.End()
|
||||
_, err = c.engineAPI.newPayload(ctx, *payload, blobHashes, beaconRoot, requests, false)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@
|
|||
package catalyst
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strconv"
|
||||
"time"
|
||||
|
|
@ -90,7 +91,9 @@ func (api *ConsensusAPI) NewPayloadWithWitnessV1(params engine.ExecutableData) (
|
|||
if params.Withdrawals != nil {
|
||||
return engine.PayloadStatusV1{Status: engine.INVALID}, engine.InvalidParams.With(errors.New("withdrawals not supported in V1"))
|
||||
}
|
||||
return api.newPayload(params, nil, nil, nil, true)
|
||||
ctx, span := startNewPayloadSpan(context.Background(), "engine.newPayloadWithWitnessV1", params)
|
||||
defer span.End()
|
||||
return api.newPayload(ctx, params, nil, nil, nil, true)
|
||||
}
|
||||
|
||||
// NewPayloadWithWitnessV2 is analogous to NewPayloadV2, only it also generates
|
||||
|
|
@ -112,7 +115,9 @@ func (api *ConsensusAPI) NewPayloadWithWitnessV2(params engine.ExecutableData) (
|
|||
case params.BlobGasUsed != nil:
|
||||
return invalidStatus, paramsErr("non-nil blobGasUsed pre-cancun")
|
||||
}
|
||||
return api.newPayload(params, nil, nil, nil, true)
|
||||
ctx, span := startNewPayloadSpan(context.Background(), "engine.newPayloadWithWitnessV2", params)
|
||||
defer span.End()
|
||||
return api.newPayload(ctx, params, nil, nil, nil, true)
|
||||
}
|
||||
|
||||
// NewPayloadWithWitnessV3 is analogous to NewPayloadV3, only it also generates
|
||||
|
|
@ -132,7 +137,9 @@ func (api *ConsensusAPI) NewPayloadWithWitnessV3(params engine.ExecutableData, v
|
|||
case !api.checkFork(params.Timestamp, forks.Cancun):
|
||||
return invalidStatus, unsupportedForkErr("newPayloadV3 must only be called for cancun payloads")
|
||||
}
|
||||
return api.newPayload(params, versionedHashes, beaconRoot, nil, true)
|
||||
ctx, span := startNewPayloadSpan(context.Background(), "engine.newPayloadWithWitnessV3", params)
|
||||
defer span.End()
|
||||
return api.newPayload(ctx, params, versionedHashes, beaconRoot, nil, true)
|
||||
}
|
||||
|
||||
// NewPayloadWithWitnessV4 is analogous to NewPayloadV4, only it also generates
|
||||
|
|
@ -158,7 +165,9 @@ func (api *ConsensusAPI) NewPayloadWithWitnessV4(params engine.ExecutableData, v
|
|||
if err := validateRequests(requests); err != nil {
|
||||
return engine.PayloadStatusV1{Status: engine.INVALID}, engine.InvalidParams.With(err)
|
||||
}
|
||||
return api.newPayload(params, versionedHashes, beaconRoot, requests, true)
|
||||
ctx, span := startNewPayloadSpan(context.Background(), "engine.newPayloadWithWitnessV4", params)
|
||||
defer span.End()
|
||||
return api.newPayload(ctx, params, versionedHashes, beaconRoot, requests, true)
|
||||
}
|
||||
|
||||
// ExecuteStatelessPayloadV1 is analogous to NewPayloadV1, only it operates in
|
||||
|
|
|
|||
Loading…
Reference in a new issue