diff --git a/eth/catalyst/api.go b/eth/catalyst/api.go index 37bd3172c4..17b67de55c 100644 --- a/eth/catalyst/api.go +++ b/eth/catalyst/api.go @@ -47,6 +47,7 @@ import ( "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" + "go.opentelemetry.io/otel/trace" ) // Register adds the engine API to the full node. @@ -613,25 +614,33 @@ func (api *ConsensusAPI) GetBlobsV2(hashes []common.Hash) ([]*engine.BlobAndProo // Helper for NewPayload* methods. var invalidStatus = engine.PayloadStatusV1{Status: engine.INVALID} -type newPayloadMethod string - -const ( - newPayloadV1 newPayloadMethod = "engine.newPayloadV1" - newPayloadV2 newPayloadMethod = "engine.newPayloadV2" - newPayloadV3 newPayloadMethod = "engine.newPayloadV3" - newPayloadV4 newPayloadMethod = "engine.newPayloadV4" -) +func startNewPayloadSpan(ctx context.Context, name string, params engine.ExecutableData) (context.Context, trace.Span) { + tracer := getEngineTracer() + 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) { + // TODO(jrhea): Decide whether validation errors should be recorded as span errors. + ctx, span := startNewPayloadSpan(context.Background(), "engine.newPayloadV1", params) + defer span.End() if params.Withdrawals != nil { return invalidStatus, paramsErr("withdrawals not supported in V1") } - return api.newPayload(newPayloadV1, 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) { + // TODO(jrhea): Decide whether validation errors should be recorded as span errors. + ctx, span := startNewPayloadSpan(context.Background(), "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) @@ -648,11 +657,14 @@ func (api *ConsensusAPI) NewPayloadV2(params engine.ExecutableData) (engine.Payl case params.BlobGasUsed != nil: return invalidStatus, paramsErr("non-nil blobGasUsed pre-cancun") } - return api.newPayload(newPayloadV2, 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) { + // TODO(jrhea): Decide whether validation errors should be recorded as span errors. + ctx, span := startNewPayloadSpan(context.Background(), "engine.newPayloadV3", params) + defer span.End() switch { case params.Withdrawals == nil: return invalidStatus, paramsErr("nil withdrawals post-shanghai") @@ -667,11 +679,14 @@ 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(newPayloadV3, 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) { + // TODO(jrhea): Decide whether validation errors should be recorded as span errors. + ctx, span := startNewPayloadSpan(context.Background(), "engine.newPayloadV4", params) + defer span.End() switch { case params.Withdrawals == nil: return invalidStatus, paramsErr("nil withdrawals post-shanghai") @@ -692,10 +707,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(newPayloadV4, params, versionedHashes, beaconRoot, requests, false) + return api.newPayload(ctx, params, versionedHashes, beaconRoot, requests, false) } -func (api *ConsensusAPI) newPayload(method newPayloadMethod, 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 @@ -712,24 +727,16 @@ func (api *ConsensusAPI) newPayload(method newPayloadMethod, params engine.Execu api.newPayloadLock.Lock() defer api.newPayloadLock.Unlock() - // TODO(jrhea): Propagate context from RPC. - tracer := getEngineTracer() - ctx, span := tracer.Start(context.Background(), string(method)) - defer span.End() - - span.SetAttributes( - attribute.Int64("block.number", int64(params.Number)), - attribute.String("block.hash", params.BlockHash.Hex()), - attribute.Int("tx.count", len(params.Transactions)), - ) - log.Trace("Engine API request received", "method", "NewPayload", "number", params.Number, "hash", params.BlockHash) + + tracer := getEngineTracer() + rootSpan := trace.SpanFromContext(ctx) _, decodeSpan := tracer.Start(ctx, "payload.decode") block, err := engine.ExecutableDataToBlock(params, versionedHashes, beaconRoot, requests) if err != nil { decodeSpan.RecordError(err) decodeSpan.SetStatus(codes.Error, err.Error()) - span.SetStatus(codes.Error, err.Error()) + rootSpan.SetStatus(codes.Error, err.Error()) } decodeSpan.End() if err != nil { @@ -809,7 +816,7 @@ func (api *ConsensusAPI) newPayload(method newPayloadMethod, params engine.Execu if err != nil { importSpan.RecordError(err) importSpan.SetStatus(codes.Error, err.Error()) - span.SetStatus(codes.Error, err.Error()) + rootSpan.SetStatus(codes.Error, err.Error()) } importSpan.End() if err != nil { diff --git a/eth/catalyst/api_test.go b/eth/catalyst/api_test.go index 4e3b76c993..8a0c30fb56 100644 --- a/eth/catalyst/api_test.go +++ b/eth/catalyst/api_test.go @@ -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(newPayloadV4, *envelope.ExecutionPayload, []common.Hash{}, h, envelope.Requests, false) + + // NOTE: This span is for the test harness only. Engine RPC root 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) } diff --git a/eth/catalyst/simulated_beacon.go b/eth/catalyst/simulated_beacon.go index 4b014cdd5f..d4302a33ab 100644 --- a/eth/catalyst/simulated_beacon.go +++ b/eth/catalyst/simulated_beacon.go @@ -17,6 +17,7 @@ package catalyst import ( + "context" "crypto/rand" "crypto/sha256" "errors" @@ -255,20 +256,13 @@ func (c *SimulatedBeacon) sealBlock(withdrawals []*types.Withdrawal, timestamp u requests = envelope.Requests } - var method newPayloadMethod - switch version { - case engine.PayloadV1: - method = newPayloadV1 - case engine.PayloadV2: - method = newPayloadV2 - case engine.PayloadV3: - method = newPayloadV3 - default: - method = newPayloadV4 - } - - // Mark the payload as canon - _, err = c.engineAPI.newPayload(method, *payload, blobHashes, beaconRoot, requests, false) + // NOTE: This span is for the simulated beacon harness only. It is intentionally + // named differently from engine.newPayload* to avoid double-rooting or + // mislabeling Engine RPC spans. Real 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 } diff --git a/eth/catalyst/witness.go b/eth/catalyst/witness.go index b52be88100..bc71df0863 100644 --- a/eth/catalyst/witness.go +++ b/eth/catalyst/witness.go @@ -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(newPayloadV1, 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(newPayloadV2, 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(newPayloadV3, 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(newPayloadV4, 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