diff --git a/packages/beacon-node/src/chain/chain.ts b/packages/beacon-node/src/chain/chain.ts index fc6d683c197c..f30642f9a621 100644 --- a/packages/beacon-node/src/chain/chain.ts +++ b/packages/beacon-node/src/chain/chain.ts @@ -858,6 +858,24 @@ export class BeaconChain implements IBeaconChain { return null; } + async getSerializedExecutionPayloadEnvelope(blockSlot: Slot, blockRootHex: string): Promise { + const payloadInput = this.seenPayloadEnvelopeInputCache.get(blockRootHex); + if (payloadInput?.hasPayloadEnvelope()) { + const envelope = payloadInput.getPayloadEnvelope(); + const serialized = this.serializedCache.get(envelope); + if (serialized) { + return serialized; + } + return ssz.gloas.SignedExecutionPayloadEnvelope.serialize(envelope); + } + + return ( + (await this.db.executionPayloadEnvelope.getBinary(fromHex(blockRootHex))) ?? + (await this.db.executionPayloadEnvelopeArchive.getBinary(blockSlot)) ?? + null + ); + } + async getDataColumnSidecars(blockSlot: Slot, blockRootHex: string): Promise { const blockInput = this.seenBlockInputCache.get(blockRootHex); if (blockInput) { diff --git a/packages/beacon-node/src/chain/emitter.ts b/packages/beacon-node/src/chain/emitter.ts index 17818b122021..544d32cf14c7 100644 --- a/packages/beacon-node/src/chain/emitter.ts +++ b/packages/beacon-node/src/chain/emitter.ts @@ -66,6 +66,15 @@ export enum ChainEvent { * cut-off window passes for waiting on gossip */ incompleteBlockInput = "incompleteBlockInput", + /** + * Trigger BlockInputSync to fetch a missing execution payload envelope for a known beacon block root + */ + unknownPayloadEnvelope = "unknownPayloadEnvelope", + /** + * Trigger BlockInputSync when a gossip block's parent block is known but parent payload is missing. + * Tracks the child block in pendingBlocks and triggers parent payload fetch. + */ + unknownParentPayload = "unknownParentPayload", } export type HeadEventData = routes.events.EventData[routes.events.EventType.head]; @@ -78,6 +87,8 @@ export type ChainEventData = { [ChainEvent.unknownParent]: {blockInput: IBlockInput; peer: PeerIdStr; source: BlockInputSource}; [ChainEvent.unknownBlockRoot]: {rootHex: RootHex; peer?: PeerIdStr; source: BlockInputSource}; [ChainEvent.incompleteBlockInput]: {blockInput: IBlockInput; peer: PeerIdStr; source: BlockInputSource}; + [ChainEvent.unknownPayloadEnvelope]: {blockRootHex: RootHex; peer?: PeerIdStr; source: BlockInputSource}; + [ChainEvent.unknownParentPayload]: {blockInput: IBlockInput; peer: PeerIdStr; source: BlockInputSource}; }; export type IChainEvents = ApiEvents & { @@ -99,6 +110,8 @@ export type IChainEvents = ApiEvents & { [ChainEvent.unknownParent]: (data: ChainEventData[ChainEvent.unknownParent]) => void; [ChainEvent.unknownBlockRoot]: (data: ChainEventData[ChainEvent.unknownBlockRoot]) => void; [ChainEvent.incompleteBlockInput]: (data: ChainEventData[ChainEvent.incompleteBlockInput]) => void; + [ChainEvent.unknownPayloadEnvelope]: (data: ChainEventData[ChainEvent.unknownPayloadEnvelope]) => void; + [ChainEvent.unknownParentPayload]: (data: ChainEventData[ChainEvent.unknownParentPayload]) => void; }; /** diff --git a/packages/beacon-node/src/chain/errors/blockError.ts b/packages/beacon-node/src/chain/errors/blockError.ts index 7eb02a60bb40..b3ed21870e17 100644 --- a/packages/beacon-node/src/chain/errors/blockError.ts +++ b/packages/beacon-node/src/chain/errors/blockError.ts @@ -70,6 +70,8 @@ export enum BlockErrorCode { TOO_MANY_KZG_COMMITMENTS = "BLOCK_ERROR_TOO_MANY_KZG_COMMITMENTS", /** Bid parent block root does not match block parent root */ BID_PARENT_ROOT_MISMATCH = "BLOCK_ERROR_BID_PARENT_ROOT_MISMATCH", + /** Parent block is known but the matching FULL payload variant is missing */ + PARENT_PAYLOAD_UNKNOWN = "BLOCK_ERROR_PARENT_PAYLOAD_UNKNOWN", } type ExecutionErrorStatus = Exclude< @@ -114,7 +116,8 @@ export type BlockErrorType = | {code: BlockErrorCode.EXECUTION_ENGINE_ERROR; execStatus: ExecutionErrorStatus; errorMessage: string} | {code: BlockErrorCode.DATA_UNAVAILABLE} | {code: BlockErrorCode.TOO_MANY_KZG_COMMITMENTS; blobKzgCommitmentsLen: number; commitmentLimit: number} - | {code: BlockErrorCode.BID_PARENT_ROOT_MISMATCH; bidParentRoot: RootHex; blockParentRoot: RootHex}; + | {code: BlockErrorCode.BID_PARENT_ROOT_MISMATCH; bidParentRoot: RootHex; blockParentRoot: RootHex} + | {code: BlockErrorCode.PARENT_PAYLOAD_UNKNOWN; parentRoot: RootHex; parentBlockHash: RootHex}; export class BlockGossipError extends GossipActionError {} diff --git a/packages/beacon-node/src/chain/interface.ts b/packages/beacon-node/src/chain/interface.ts index f9d78cb8d51a..cf37e336a9b0 100644 --- a/packages/beacon-node/src/chain/interface.ts +++ b/packages/beacon-node/src/chain/interface.ts @@ -223,6 +223,7 @@ export interface IBeaconChain { blockRootHex: string, indices: number[] ): Promise<(Uint8Array | undefined)[]>; + getSerializedExecutionPayloadEnvelope(blockSlot: Slot, blockRootHex: string): Promise; produceCommonBlockBody(blockAttributes: BlockAttributes): Promise; produceBlock(blockAttributes: BlockAttributes & {commonBlockBodyPromise: Promise}): Promise<{ diff --git a/packages/beacon-node/src/chain/validation/block.ts b/packages/beacon-node/src/chain/validation/block.ts index f7addf9a5794..9e97982733e2 100644 --- a/packages/beacon-node/src/chain/validation/block.ts +++ b/packages/beacon-node/src/chain/validation/block.ts @@ -79,6 +79,15 @@ export async function validateGossipBlock( ) : chain.forkChoice.getBlockHexDefaultStatus(parentRoot); if (parentBlock === null) { + // For Gloas blocks: if the parent block root is known but the matching FULL variant + // (by parentBlockHash) is missing, this is a payload availability issue, not a missing block. + if (isGloasBeaconBlock(block) && chain.forkChoice.hasBlockHexUnsafe(parentRoot)) { + throw new BlockGossipError(GossipAction.IGNORE, { + code: BlockErrorCode.PARENT_PAYLOAD_UNKNOWN, + parentRoot, + parentBlockHash: toRootHex(block.body.signedExecutionPayloadBid.message.parentBlockHash), + }); + } // If fork choice does *not* consider the parent to be a descendant of the finalized block, // then there are two more cases: // diff --git a/packages/beacon-node/src/metrics/metrics/lodestar.ts b/packages/beacon-node/src/metrics/metrics/lodestar.ts index 00087979fc9a..4fad820e3e0e 100644 --- a/packages/beacon-node/src/metrics/metrics/lodestar.ts +++ b/packages/beacon-node/src/metrics/metrics/lodestar.ts @@ -680,6 +680,23 @@ export function createLodestarMetrics( labelNames: ["code", "client"], }), }, + pendingPayloads: register.gauge({ + name: "lodestar_sync_block_input_pending_payloads_size", + help: "Current size of pending payloads cache in BlockInputSync", + }), + payloadFetchSuccess: register.gauge({ + name: "lodestar_sync_block_input_payload_fetch_success_total", + help: "Total successful payload envelope fetches", + }), + payloadFetchError: register.gauge({ + name: "lodestar_sync_block_input_payload_fetch_error_total", + help: "Total failed payload envelope fetches", + }), + payloadRequests: register.gauge<{source: string}>({ + name: "lodestar_sync_block_input_payload_requests_total", + help: "Total unknown payload events by source", + labelNames: ["source"], + }), peerBalancer: { peersMetaCount: register.gauge({ name: "lodestar_sync_unknown_block_peer_balancer_peers_meta_count", @@ -1746,6 +1763,29 @@ export function createLodestarMetrics( }), }, + awaitingEnvelopeGossipMessages: { + resolve: register.gauge<{topic: GossipType}>({ + name: "lodestar_awaiting_envelope_gossip_messages_resolve_total", + help: "Total number of gossip messages reprocessed after payload envelope import", + labelNames: ["topic"], + }), + waitSecBeforeResolve: register.gauge<{topic: GossipType}>({ + name: "lodestar_awaiting_envelope_gossip_messages_wait_time_resolve_seconds", + help: "Time to wait for unknown payload envelope in seconds", + labelNames: ["topic"], + }), + reject: register.gauge<{reason: ReprocessRejectReason; topic: GossipType}>({ + name: "lodestar_awaiting_envelope_gossip_messages_reject_total", + help: "Total number of gossip messages rejected while waiting for payload envelope", + labelNames: ["reason", "topic"], + }), + waitSecBeforeReject: register.gauge<{reason: ReprocessRejectReason; topic: GossipType}>({ + name: "lodestar_awaiting_envelope_gossip_messages_wait_time_reject_seconds", + help: "Time to wait for unknown payload envelope before being rejected", + labelNames: ["reason", "topic"], + }), + }, + lightclientServer: { onSyncAggregate: register.gauge<{event: string}>({ name: "lodestar_lightclient_server_on_sync_aggregate_event_total", diff --git a/packages/beacon-node/src/network/interface.ts b/packages/beacon-node/src/network/interface.ts index 25abc19d6c51..0c748e3eef70 100644 --- a/packages/beacon-node/src/network/interface.ts +++ b/packages/beacon-node/src/network/interface.ts @@ -38,7 +38,12 @@ import { import {BlockInputSource} from "../chain/blocks/blockInput/types.js"; import {CustodyConfig} from "../util/dataColumns.js"; import {PeerIdStr} from "../util/peerId.js"; -import {BeaconBlocksByRootRequest, BlobSidecarsByRootRequest, DataColumnSidecarsByRootRequest} from "../util/types.js"; +import { + BeaconBlocksByRootRequest, + BlobSidecarsByRootRequest, + DataColumnSidecarsByRootRequest, + ExecutionPayloadEnvelopesByRootRequest, +} from "../util/types.js"; import {INetworkCorePublic} from "./core/types.js"; import {INetworkEventBus} from "./events.js"; import {GossipType} from "./gossip/interface.js"; @@ -82,6 +87,14 @@ export interface INetwork extends INetworkCorePublic { peerId: PeerIdStr, request: DataColumnSidecarsByRootRequest ): Promise; + sendExecutionPayloadEnvelopesByRange( + peerId: PeerIdStr, + request: gloas.ExecutionPayloadEnvelopesByRangeRequest + ): Promise; + sendExecutionPayloadEnvelopesByRoot( + peerId: PeerIdStr, + request: ExecutionPayloadEnvelopesByRootRequest + ): Promise; // Gossip publishBeaconBlock(signedBlock: SignedBeaconBlock): Promise; diff --git a/packages/beacon-node/src/network/network.ts b/packages/beacon-node/src/network/network.ts index 141295db4b7b..ebc84cda0a33 100644 --- a/packages/beacon-node/src/network/network.ts +++ b/packages/beacon-node/src/network/network.ts @@ -40,7 +40,12 @@ import {IClock} from "../util/clock.js"; import {CustodyConfig} from "../util/dataColumns.js"; import {PeerIdStr, peerIdToString} from "../util/peerId.js"; import {promiseAllMaybeAsync} from "../util/promises.js"; -import {BeaconBlocksByRootRequest, BlobSidecarsByRootRequest, DataColumnSidecarsByRootRequest} from "../util/types.js"; +import { + BeaconBlocksByRootRequest, + BlobSidecarsByRootRequest, + DataColumnSidecarsByRootRequest, + ExecutionPayloadEnvelopesByRootRequest, +} from "../util/types.js"; import {INetworkCore, NetworkCore, WorkerNetworkCore} from "./core/index.js"; import {INetworkEventBus, NetworkEvent, NetworkEventBus, NetworkEventData} from "./events.js"; import {getActiveForkBoundaries} from "./forks.js"; @@ -636,6 +641,29 @@ export class Network implements INetwork { ); } + async sendExecutionPayloadEnvelopesByRange( + peerId: PeerIdStr, + request: gloas.ExecutionPayloadEnvelopesByRangeRequest + ): Promise { + return collectMaxResponseTyped( + this.sendReqRespRequest(peerId, ReqRespMethod.ExecutionPayloadEnvelopesByRange, [Version.V1], request), + request.count, + responseSszTypeByMethod[ReqRespMethod.ExecutionPayloadEnvelopesByRange] + ); + } + + async sendExecutionPayloadEnvelopesByRoot( + peerId: PeerIdStr, + request: ExecutionPayloadEnvelopesByRootRequest + ): Promise { + return collectMaxResponseTyped( + this.sendReqRespRequest(peerId, ReqRespMethod.ExecutionPayloadEnvelopesByRoot, [Version.V1], request), + request.length, + responseSszTypeByMethod[ReqRespMethod.ExecutionPayloadEnvelopesByRoot], + this.chain.serializedCache + ); + } + private sendReqRespRequest( peerId: PeerIdStr, method: ReqRespMethod, diff --git a/packages/beacon-node/src/network/processor/gossipHandlers.ts b/packages/beacon-node/src/network/processor/gossipHandlers.ts index 9c2a9a031c9b..b927bee5ca72 100644 --- a/packages/beacon-node/src/network/processor/gossipHandlers.ts +++ b/packages/beacon-node/src/network/processor/gossipHandlers.ts @@ -197,6 +197,16 @@ function getSequentialHandlers(modules: ValidatorFnsModules, options: GossipHand throw e; } + if (e.type.code === BlockErrorCode.PARENT_PAYLOAD_UNKNOWN && blockInput) { + logger.debug("Gossip block has parent payload unknown", {slot, root: blockShortHex, code: e.type.code}); + chain.emitter.emit(ChainEvent.unknownParentPayload, { + blockInput, + peer: peerIdStr, + source: BlockInputSource.gossip, + }); + throw e; + } + if (e.action === GossipAction.REJECT) { chain.persistInvalidSszValue(forkTypes.SignedBeaconBlock, signedBlock, `gossip_reject_slot_${slot}`); } @@ -842,7 +852,7 @@ function getSequentialHandlers(modules: ValidatorFnsModules, options: GossipHand }: GossipHandlerParamGeneric) => { const {serializedData} = gossipData; const executionPayloadEnvelope = sszDeserialize(topic, serializedData); - // TODO GLOAS: handle BLOCK_ROOT_UNKNOWN error to trigger sync + // BLOCK_ROOT_UNKNWON should not happen here. It would be caught early on and be queued in awaitingMessagesByBlockRoot await validateGossipExecutionPayloadEnvelope(chain, executionPayloadEnvelope); const slot = executionPayloadEnvelope.message.slot; diff --git a/packages/beacon-node/src/network/processor/index.ts b/packages/beacon-node/src/network/processor/index.ts index 73e52c382487..ee77c6aa0455 100644 --- a/packages/beacon-node/src/network/processor/index.ts +++ b/packages/beacon-node/src/network/processor/index.ts @@ -1,4 +1,5 @@ import {routes} from "@lodestar/api"; +import {PayloadStatus} from "@lodestar/fork-choice"; import {ForkSeq} from "@lodestar/params"; import {computeStartSlotAtEpoch} from "@lodestar/state-transition"; import {RootHex, Slot, SlotRootHex} from "@lodestar/types"; @@ -159,7 +160,10 @@ export class NetworkProcessor { // we may not receive the block for messages like Attestation and SignedAggregateAndProof messages, in that case PendingGossipsubMessage needs // to be stored in this Map and reprocessed once the block comes private readonly awaitingMessagesByBlockRoot: MapDef>; + // Messages waiting for a payload envelope to be imported for a known beacon block root + private readonly awaitingMessagesByEnvelopeBlockRoot: MapDef>; private unknownBlocksBySlot = new MapDef>(() => new Set()); + private unknownEnvelopesBySlot = new MapDef>(() => new Set()); constructor( modules: NetworkProcessorModules, @@ -181,9 +185,11 @@ export class NetworkProcessor { events.on(NetworkEvent.pendingGossipsubMessage, this.onPendingGossipsubMessage.bind(this)); this.chain.emitter.on(routes.events.EventType.block, this.onBlockProcessed.bind(this)); + this.chain.emitter.on(routes.events.EventType.executionPayloadAvailable, this.onPayloadProcessed.bind(this)); this.chain.clock.on(ClockEvent.slot, this.onClockSlot.bind(this)); this.awaitingMessagesByBlockRoot = new MapDef>(() => new Set()); + this.awaitingMessagesByEnvelopeBlockRoot = new MapDef>(() => new Set()); // TODO: Implement queues and priorization for ReqResp incoming requests // Listens to NetworkEvent.reqRespIncomingRequest event @@ -212,6 +218,7 @@ export class NetworkProcessor { async stop(): Promise { this.events.off(NetworkEvent.pendingGossipsubMessage, this.onPendingGossipsubMessage); this.chain.emitter.off(routes.events.EventType.block, this.onBlockProcessed); + this.chain.emitter.off(routes.events.EventType.executionPayloadAvailable, this.onPayloadProcessed); this.chain.emitter.off(ClockEvent.slot, this.onClockSlot); } @@ -248,6 +255,22 @@ export class NetworkProcessor { this.chain.emitter.emit(ChainEvent.unknownBlockRoot, {rootHex: root, peer, source}); } + /** + * Search for a missing execution payload envelope for a known beacon block root. + * Only emits if the block is known but FULL variant is missing. + */ + searchUnknownEnvelope({slot, root: blockRoot}: SlotRootHex, source: BlockInputSource, peer?: PeerIdStr): void { + if (this.awaitingMessagesByEnvelopeBlockRoot.has(blockRoot)) return; + if (this.unknownEnvelopesBySlot.getOrDefault(slot).has(blockRoot)) return; + // Only search if block is known — if block is unknown, searchUnknownBlock handles it + if (!this.chain.forkChoice.hasBlockHexUnsafe(blockRoot)) return; + // If FULL variant already exists, no need to search + if (this.chain.forkChoice.getBlockHex(blockRoot, PayloadStatus.FULL)) return; + + this.unknownEnvelopesBySlot.getOrDefault(slot).add(blockRoot); + this.chain.emitter.emit(ChainEvent.unknownPayloadEnvelope, {blockRootHex: blockRoot, peer, source}); + } + private onPendingGossipsubMessage(message: PendingGossipsubMessage): void { const topicType = message.topic.type; const extractBlockSlotRootFn = this.extractBlockSlotRootFns[topicType]; @@ -301,6 +324,14 @@ export class NetworkProcessor { return; } + // Block is known — check if message requires FULL payload state + // TODO GLOAS (PR #9025): Add evidence routing for specific gossip types: + // - attestation index === 1: needs FULL payload + // - PTC payloadPresent === true: needs FULL payload + // - data_column_sidecar in Gloas: needs FULL payload + // For now, these messages proceed to validation where they will fail with + // appropriate errors if FULL variant is missing + this.pushPendingGossipsubMessageToQueue(message); } @@ -345,6 +376,32 @@ export class NetworkProcessor { this.awaitingMessagesByBlockRoot.delete(rootHex); } + private async onPayloadProcessed({blockRoot}: {slot: Slot; blockRoot: string}): Promise { + const waitingMessages = this.awaitingMessagesByEnvelopeBlockRoot.get(blockRoot); + if (!waitingMessages || waitingMessages.size === 0) { + return; + } + + const nowSec = Date.now() / 1000; + let count = 0; + for (const message of waitingMessages) { + const topicType = message.topic.type; + this.metrics?.awaitingEnvelopeGossipMessages.waitSecBeforeResolve.set( + {topic: topicType}, + nowSec - message.seenTimestampSec + ); + this.metrics?.awaitingEnvelopeGossipMessages.resolve.inc({topic: topicType}); + this.pushPendingGossipsubMessageToQueue(message); + count++; + if (count === MAX_AWAITING_GOSSIP_OBJECTS_PER_TICK) { + count = 0; + await sleep(AWAITING_GOSSIP_OBJECTS_YIELD_EVERY_MS); + } + } + + this.awaitingMessagesByEnvelopeBlockRoot.delete(blockRoot); + } + private onClockSlot(clockSlot: Slot): void { const nowSec = Date.now() / 1000; const minSlot = clockSlot - MAX_UNKNOWN_ROOTS_SLOT_CACHE_SIZE; @@ -371,6 +428,29 @@ export class NetworkProcessor { } this.unknownBlocksBySlot.delete(slot); } + + // Prune envelope waiting maps in parallel + for (const [slot, roots] of this.unknownEnvelopesBySlot) { + if (slot > minSlot) continue; + for (const rootHex of roots) { + const gossipMessages = this.awaitingMessagesByEnvelopeBlockRoot.get(rootHex); + if (gossipMessages !== undefined) { + for (const message of gossipMessages) { + const topicType = message.topic.type; + this.metrics?.awaitingEnvelopeGossipMessages.reject.inc({ + topic: topicType, + reason: ReprocessRejectReason.expired, + }); + this.metrics?.awaitingEnvelopeGossipMessages.waitSecBeforeReject.set( + {topic: topicType, reason: ReprocessRejectReason.expired}, + nowSec - message.seenTimestampSec + ); + } + this.awaitingMessagesByEnvelopeBlockRoot.delete(rootHex); + } + } + this.unknownEnvelopesBySlot.delete(slot); + } } private executeWork(): void { diff --git a/packages/beacon-node/src/network/reqresp/ReqRespBeaconNode.ts b/packages/beacon-node/src/network/reqresp/ReqRespBeaconNode.ts index 9a3d5dc83204..5aba7277cbac 100644 --- a/packages/beacon-node/src/network/reqresp/ReqRespBeaconNode.ts +++ b/packages/beacon-node/src/network/reqresp/ReqRespBeaconNode.ts @@ -297,6 +297,19 @@ export class ReqRespBeaconNode extends ReqResp { ); } + if (ForkSeq[fork] >= ForkSeq.gloas) { + protocolsAtFork.push( + [ + protocols.ExecutionPayloadEnvelopesByRoot(fork, this.config), + this.getHandler(ReqRespMethod.ExecutionPayloadEnvelopesByRoot), + ], + [ + protocols.ExecutionPayloadEnvelopesByRange(fork, this.config), + this.getHandler(ReqRespMethod.ExecutionPayloadEnvelopesByRange), + ] + ); + } + return protocolsAtFork; } diff --git a/packages/beacon-node/src/network/reqresp/handlers/executionPayloadEnvelopesByRange.ts b/packages/beacon-node/src/network/reqresp/handlers/executionPayloadEnvelopesByRange.ts new file mode 100644 index 000000000000..4ecf28b6dae6 --- /dev/null +++ b/packages/beacon-node/src/network/reqresp/handlers/executionPayloadEnvelopesByRange.ts @@ -0,0 +1,79 @@ +import {ChainConfig} from "@lodestar/config"; +import {GENESIS_SLOT} from "@lodestar/params"; +import {RespStatus, ResponseError, ResponseOutgoing} from "@lodestar/reqresp"; +import {computeEpochAtSlot} from "@lodestar/state-transition"; +import {gloas} from "@lodestar/types"; +import {IBeaconChain} from "../../../chain/index.js"; +import {IBeaconDb} from "../../../db/index.js"; + +export async function* onExecutionPayloadEnvelopesByRange( + request: gloas.ExecutionPayloadEnvelopesByRangeRequest, + chain: IBeaconChain, + db: IBeaconDb +): AsyncIterable { + const {startSlot, count} = validateExecutionPayloadEnvelopesByRangeRequest(chain.config, request); + const endSlot = startSlot + count; + + const finalized = db.executionPayloadEnvelopeArchive; + const finalizedSlot = chain.forkChoice.getFinalizedBlock().slot; + + // Finalized range of envelopes + if (startSlot <= finalizedSlot) { + for await (const {key, value: envelopeBytes} of finalized.binaryEntriesStream({ + gte: startSlot, + lt: endSlot, + })) { + const slot = finalized.decodeKey(key); + yield { + data: envelopeBytes, + boundary: chain.config.getForkBoundaryAtEpoch(computeEpochAtSlot(slot)), + }; + } + } + + // Non-finalized range of envelopes + if (endSlot > finalizedSlot) { + const headBlock = chain.forkChoice.getHead(); + const headRoot = headBlock.blockRoot; + const headChain = chain.forkChoice.getAllAncestorBlocks(headRoot, headBlock.payloadStatus); + + // Iterate head chain with ascending block numbers + for (let i = headChain.length - 1; i >= 0; i--) { + const block = headChain[i]; + + if (block.slot >= startSlot && block.slot < endSlot) { + const envelopeBytes = await chain.getSerializedExecutionPayloadEnvelope(block.slot, block.blockRoot); + if (envelopeBytes) { + yield { + data: envelopeBytes, + boundary: chain.config.getForkBoundaryAtEpoch(computeEpochAtSlot(block.slot)), + }; + } + // In ePBS, missing envelopes are valid (payload withholding) — skip silently + } else if (block.slot >= endSlot) { + break; + } + } + } +} + +export function validateExecutionPayloadEnvelopesByRangeRequest( + config: ChainConfig, + request: gloas.ExecutionPayloadEnvelopesByRangeRequest +): gloas.ExecutionPayloadEnvelopesByRangeRequest { + const {startSlot} = request; + let {count} = request; + + if (count < 1) { + throw new ResponseError(RespStatus.INVALID_REQUEST, "count < 1"); + } + if (startSlot < GENESIS_SLOT) { + throw new ResponseError(RespStatus.INVALID_REQUEST, "startSlot < genesis"); + } + + if (count > config.MAX_REQUEST_BLOCKS_DENEB) { + count = config.MAX_REQUEST_BLOCKS_DENEB; + } + + return {startSlot, count}; +} diff --git a/packages/beacon-node/src/network/reqresp/handlers/executionPayloadEnvelopesByRoot.ts b/packages/beacon-node/src/network/reqresp/handlers/executionPayloadEnvelopesByRoot.ts new file mode 100644 index 000000000000..87decc1a1946 --- /dev/null +++ b/packages/beacon-node/src/network/reqresp/handlers/executionPayloadEnvelopesByRoot.ts @@ -0,0 +1,43 @@ +import {ResponseOutgoing} from "@lodestar/reqresp"; +import {computeEpochAtSlot} from "@lodestar/state-transition"; +import {toRootHex} from "@lodestar/utils"; +import {IBeaconChain} from "../../../chain/index.js"; +import {IBeaconDb} from "../../../db/index.js"; +import {ExecutionPayloadEnvelopesByRootRequest} from "../../../util/types.js"; + +export async function* onExecutionPayloadEnvelopesByRoot( + requestBody: ExecutionPayloadEnvelopesByRootRequest, + chain: IBeaconChain, + db: IBeaconDb +): AsyncIterable { + // Spec: [max(GLOAS_FORK_EPOCH, current_epoch - MIN_EPOCHS_FOR_BLOCK_REQUESTS), current_epoch] + const currentEpoch = chain.clock.currentEpoch; + const minimumRequestEpoch = Math.max( + currentEpoch - chain.config.MIN_EPOCHS_FOR_BLOCK_REQUESTS, + chain.config.GLOAS_FORK_EPOCH + ); + + for (const root of requestBody) { + const rootHex = toRootHex(root); + const block = chain.forkChoice.getBlockHexDefaultStatus(rootHex); + // If the block is not in fork choice, it may be finalized. Attempt to find its slot in block archive + const slot = block ? block.slot : await db.blockArchive.getSlotByRoot(root); + + if (slot === null) { + continue; + } + + const requestedEpoch = computeEpochAtSlot(slot); + if (requestedEpoch < minimumRequestEpoch) { + continue; + } + + const envelopeBytes = await chain.getSerializedExecutionPayloadEnvelope(slot, rootHex); + if (envelopeBytes) { + yield { + data: envelopeBytes, + boundary: chain.config.getForkBoundaryAtEpoch(requestedEpoch), + }; + } + } +} diff --git a/packages/beacon-node/src/network/reqresp/handlers/index.ts b/packages/beacon-node/src/network/reqresp/handlers/index.ts index 9777bbf23964..c5c6109741ea 100644 --- a/packages/beacon-node/src/network/reqresp/handlers/index.ts +++ b/packages/beacon-node/src/network/reqresp/handlers/index.ts @@ -6,6 +6,7 @@ import { BeaconBlocksByRootRequestType, BlobSidecarsByRootRequestType, DataColumnSidecarsByRootRequestType, + ExecutionPayloadEnvelopesByRootRequestType, } from "../../../util/types.js"; import {GetReqRespHandlerFn, ReqRespMethod} from "../types.js"; import {onBeaconBlocksByRange} from "./beaconBlocksByRange.js"; @@ -14,6 +15,8 @@ import {onBlobSidecarsByRange} from "./blobSidecarsByRange.js"; import {onBlobSidecarsByRoot} from "./blobSidecarsByRoot.js"; import {onDataColumnSidecarsByRange} from "./dataColumnSidecarsByRange.js"; import {onDataColumnSidecarsByRoot} from "./dataColumnSidecarsByRoot.js"; +import {onExecutionPayloadEnvelopesByRange} from "./executionPayloadEnvelopesByRange.js"; +import {onExecutionPayloadEnvelopesByRoot} from "./executionPayloadEnvelopesByRoot.js"; import {onLightClientBootstrap} from "./lightClientBootstrap.js"; import {onLightClientFinalityUpdate} from "./lightClientFinalityUpdate.js"; import {onLightClientOptimisticUpdate} from "./lightClientOptimisticUpdate.js"; @@ -62,6 +65,15 @@ export function getReqRespHandlers({db, chain}: {db: IBeaconDb; chain: IBeaconCh return onDataColumnSidecarsByRoot(body, chain, db, peerId, peerClient); }, + [ReqRespMethod.ExecutionPayloadEnvelopesByRoot]: (req) => { + const body = ExecutionPayloadEnvelopesByRootRequestType(chain.config).deserialize(req.data); + return onExecutionPayloadEnvelopesByRoot(body, chain, db); + }, + [ReqRespMethod.ExecutionPayloadEnvelopesByRange]: (req) => { + const body = ssz.gloas.ExecutionPayloadEnvelopesByRangeRequest.deserialize(req.data); + return onExecutionPayloadEnvelopesByRange(body, chain, db); + }, + [ReqRespMethod.LightClientBootstrap]: (req) => { const body = ssz.Root.deserialize(req.data); return onLightClientBootstrap(body, chain); diff --git a/packages/beacon-node/src/network/reqresp/protocols.ts b/packages/beacon-node/src/network/reqresp/protocols.ts index 1bc83984bedb..17dfc84b4910 100644 --- a/packages/beacon-node/src/network/reqresp/protocols.ts +++ b/packages/beacon-node/src/network/reqresp/protocols.ts @@ -94,6 +94,18 @@ export const DataColumnSidecarsByRoot = toProtocol({ contextBytesType: ContextBytesType.ForkDigest, }); +export const ExecutionPayloadEnvelopesByRoot = toProtocol({ + method: ReqRespMethod.ExecutionPayloadEnvelopesByRoot, + version: Version.V1, + contextBytesType: ContextBytesType.ForkDigest, +}); + +export const ExecutionPayloadEnvelopesByRange = toProtocol({ + method: ReqRespMethod.ExecutionPayloadEnvelopesByRange, + version: Version.V1, + contextBytesType: ContextBytesType.ForkDigest, +}); + export const LightClientBootstrap = toProtocol({ method: ReqRespMethod.LightClientBootstrap, version: Version.V1, diff --git a/packages/beacon-node/src/network/reqresp/rateLimit.ts b/packages/beacon-node/src/network/reqresp/rateLimit.ts index 0bacda6e931d..924ece0d6902 100644 --- a/packages/beacon-node/src/network/reqresp/rateLimit.ts +++ b/packages/beacon-node/src/network/reqresp/rateLimit.ts @@ -73,6 +73,24 @@ export const rateLimitQuotas: (fork: ForkName, config: BeaconConfig) => Record total + item.columns.length, 0) ), }, + [ReqRespMethod.ExecutionPayloadEnvelopesByRoot]: { + byPeer: {quota: config.MAX_REQUEST_PAYLOADS, quotaTimeMs: 10_000}, + getRequestCount: getRequestCountFn( + fork, + config, + ReqRespMethod.ExecutionPayloadEnvelopesByRoot, + (req) => req.length + ), + }, + [ReqRespMethod.ExecutionPayloadEnvelopesByRange]: { + byPeer: {quota: config.MAX_REQUEST_BLOCKS_DENEB, quotaTimeMs: 10_000}, + getRequestCount: getRequestCountFn( + fork, + config, + ReqRespMethod.ExecutionPayloadEnvelopesByRange, + (req) => req.count + ), + }, [ReqRespMethod.LightClientBootstrap]: { // As similar in the nature of `Status` protocol so we use the same rate limits. byPeer: {quota: 5, quotaTimeMs: 15_000}, diff --git a/packages/beacon-node/src/network/reqresp/score.ts b/packages/beacon-node/src/network/reqresp/score.ts index e9f7cee3607d..7ed5e7dd1bec 100644 --- a/packages/beacon-node/src/network/reqresp/score.ts +++ b/packages/beacon-node/src/network/reqresp/score.ts @@ -46,6 +46,8 @@ export function onOutgoingReqRespError(e: RequestError, method: ReqRespMethod): return PeerAction.LowToleranceError; case ReqRespMethod.BeaconBlocksByRange: case ReqRespMethod.BeaconBlocksByRoot: + case ReqRespMethod.ExecutionPayloadEnvelopesByRoot: + case ReqRespMethod.ExecutionPayloadEnvelopesByRange: return PeerAction.MidToleranceError; default: return null; diff --git a/packages/beacon-node/src/network/reqresp/types.ts b/packages/beacon-node/src/network/reqresp/types.ts index b9bcec66f618..722324fbed50 100644 --- a/packages/beacon-node/src/network/reqresp/types.ts +++ b/packages/beacon-node/src/network/reqresp/types.ts @@ -14,6 +14,7 @@ import { altair, deneb, fulu, + gloas, phase0, ssz, sszTypesFor, @@ -25,6 +26,8 @@ import { BlobSidecarsByRootRequestType, DataColumnSidecarsByRootRequest, DataColumnSidecarsByRootRequestType, + ExecutionPayloadEnvelopesByRootRequest, + ExecutionPayloadEnvelopesByRootRequestType, } from "../../util/types.js"; export type ProtocolNoHandler = Omit; @@ -42,6 +45,8 @@ export enum ReqRespMethod { BlobSidecarsByRoot = "blob_sidecars_by_root", DataColumnSidecarsByRange = "data_column_sidecars_by_range", DataColumnSidecarsByRoot = "data_column_sidecars_by_root", + ExecutionPayloadEnvelopesByRoot = "execution_payload_envelopes_by_root", + ExecutionPayloadEnvelopesByRange = "execution_payload_envelopes_by_range", LightClientBootstrap = "light_client_bootstrap", LightClientUpdatesByRange = "light_client_updates_by_range", LightClientFinalityUpdate = "light_client_finality_update", @@ -60,6 +65,8 @@ export type RequestBodyByMethod = { [ReqRespMethod.BlobSidecarsByRoot]: BlobSidecarsByRootRequest; [ReqRespMethod.DataColumnSidecarsByRange]: fulu.DataColumnSidecarsByRangeRequest; [ReqRespMethod.DataColumnSidecarsByRoot]: DataColumnSidecarsByRootRequest; + [ReqRespMethod.ExecutionPayloadEnvelopesByRoot]: ExecutionPayloadEnvelopesByRootRequest; + [ReqRespMethod.ExecutionPayloadEnvelopesByRange]: gloas.ExecutionPayloadEnvelopesByRangeRequest; [ReqRespMethod.LightClientBootstrap]: Root; [ReqRespMethod.LightClientUpdatesByRange]: altair.LightClientUpdatesByRange; [ReqRespMethod.LightClientFinalityUpdate]: null; @@ -78,6 +85,8 @@ type ResponseBodyByMethod = { [ReqRespMethod.BlobSidecarsByRoot]: deneb.BlobSidecar; [ReqRespMethod.DataColumnSidecarsByRange]: fulu.DataColumnSidecar; [ReqRespMethod.DataColumnSidecarsByRoot]: fulu.DataColumnSidecar; + [ReqRespMethod.ExecutionPayloadEnvelopesByRoot]: gloas.SignedExecutionPayloadEnvelope; + [ReqRespMethod.ExecutionPayloadEnvelopesByRange]: gloas.SignedExecutionPayloadEnvelope; [ReqRespMethod.LightClientBootstrap]: LightClientBootstrap; [ReqRespMethod.LightClientUpdatesByRange]: LightClientUpdate; @@ -105,6 +114,8 @@ export const requestSszTypeByMethod: ( [ReqRespMethod.BlobSidecarsByRoot]: BlobSidecarsByRootRequestType(fork, config), [ReqRespMethod.DataColumnSidecarsByRange]: ssz.fulu.DataColumnSidecarsByRangeRequest, [ReqRespMethod.DataColumnSidecarsByRoot]: DataColumnSidecarsByRootRequestType(config), + [ReqRespMethod.ExecutionPayloadEnvelopesByRoot]: ExecutionPayloadEnvelopesByRootRequestType(config), + [ReqRespMethod.ExecutionPayloadEnvelopesByRange]: ssz.gloas.ExecutionPayloadEnvelopesByRangeRequest, [ReqRespMethod.LightClientBootstrap]: ssz.Root, [ReqRespMethod.LightClientUpdatesByRange]: ssz.altair.LightClientUpdatesByRange, @@ -137,6 +148,8 @@ export const responseSszTypeByMethod: {[K in ReqRespMethod]: ResponseTypeGetter< [ReqRespMethod.LightClientFinalityUpdate]: (fork) => sszTypesFor(onlyPostAltairFork(fork)).LightClientFinalityUpdate, [ReqRespMethod.DataColumnSidecarsByRange]: () => ssz.fulu.DataColumnSidecar, [ReqRespMethod.DataColumnSidecarsByRoot]: () => ssz.fulu.DataColumnSidecar, + [ReqRespMethod.ExecutionPayloadEnvelopesByRoot]: () => ssz.gloas.SignedExecutionPayloadEnvelope, + [ReqRespMethod.ExecutionPayloadEnvelopesByRange]: () => ssz.gloas.SignedExecutionPayloadEnvelope, [ReqRespMethod.LightClientOptimisticUpdate]: (fork) => sszTypesFor(onlyPostAltairFork(fork)).LightClientOptimisticUpdate, }; diff --git a/packages/beacon-node/src/sync/types.ts b/packages/beacon-node/src/sync/types.ts index ca36095c5ca8..4330ec8e7c24 100644 --- a/packages/beacon-node/src/sync/types.ts +++ b/packages/beacon-node/src/sync/types.ts @@ -1,5 +1,6 @@ import {RootHex, Slot} from "@lodestar/types"; import {IBlockInput} from "../chain/blocks/blockInput/index.js"; +import {PeerIdStr} from "../util/peerId.js"; export enum PendingBlockType { /** @@ -55,3 +56,12 @@ export function getBlockInputSyncCacheItemRootHex(block: BlockInputSyncCacheItem export function getBlockInputSyncCacheItemSlot(block: BlockInputSyncCacheItem): Slot | string { return isPendingBlockInput(block) ? block.blockInput.slot : "unknown"; } + +export type PendingPayloadEnvelope = { + status: "pending" | "fetching"; + blockRootHex: RootHex; + slot: Slot; + attempts: number; + peerIdStrings: Set; + timeAddedSec: number; +}; diff --git a/packages/beacon-node/src/sync/unknownBlock.ts b/packages/beacon-node/src/sync/unknownBlock.ts index 11db2d8cae74..0f4819348f6b 100644 --- a/packages/beacon-node/src/sync/unknownBlock.ts +++ b/packages/beacon-node/src/sync/unknownBlock.ts @@ -1,11 +1,14 @@ +import {routes} from "@lodestar/api"; import {ChainForkConfig} from "@lodestar/config"; -import {ForkSeq} from "@lodestar/params"; +import {PayloadStatus} from "@lodestar/fork-choice"; +import {ForkSeq, isForkPostGloas} from "@lodestar/params"; import {RequestError, RequestErrorCode} from "@lodestar/reqresp"; import {computeTimeAtSlot} from "@lodestar/state-transition"; -import {RootHex} from "@lodestar/types"; -import {Logger, prettyPrintIndices, pruneSetToMax, sleep} from "@lodestar/utils"; +import {RootHex, isGloasBeaconBlock} from "@lodestar/types"; +import {Logger, fromHex, prettyPrintIndices, pruneSetToMax, sleep, toRootHex} from "@lodestar/utils"; import {isBlockInputBlobs, isBlockInputColumns} from "../chain/blocks/blockInput/blockInput.js"; import {BlockInputSource, IBlockInput} from "../chain/blocks/blockInput/types.js"; +import {PayloadEnvelopeInputSource} from "../chain/blocks/payloadEnvelopeInput/types.js"; import {BlockError, BlockErrorCode} from "../chain/errors/index.js"; import {ChainEvent, ChainEventData, IBeaconChain} from "../chain/index.js"; import {Metrics} from "../metrics/index.js"; @@ -22,6 +25,7 @@ import { PendingBlockInput, PendingBlockInputStatus, PendingBlockType, + PendingPayloadEnvelope, getBlockInputSyncCacheItemRootHex, getBlockInputSyncCacheItemSlot, isPendingBlockInput, @@ -32,6 +36,8 @@ import {getAllDescendantBlocks, getDescendantBlocks, getUnknownAndAncestorBlocks const MAX_ATTEMPTS_PER_BLOCK = 5; const MAX_KNOWN_BAD_BLOCKS = 500; const MAX_PENDING_BLOCKS = 100; +const MAX_PENDING_PAYLOADS = 100; +const MAX_ATTEMPTS_PER_PAYLOAD = 5; enum FetchResult { SuccessResolved = "success_resolved", @@ -78,6 +84,7 @@ export class BlockInputSync { * block RootHex -> PendingBlock. To avoid finding same root at the same time */ private readonly pendingBlocks = new Map(); + private readonly pendingPayloads = new Map(); private readonly knownBadBlocks = new Set(); private readonly maxPendingBlocks; private subscribedToNetworkEvents = false; @@ -101,6 +108,9 @@ export class BlockInputSync { metrics.blockInputSync.knownBadBlocks.addCollect(() => metrics.blockInputSync.knownBadBlocks.set(this.knownBadBlocks.size) ); + metrics.blockInputSync.pendingPayloads.addCollect(() => + metrics.blockInputSync.pendingPayloads.set(this.pendingPayloads.size) + ); } } @@ -116,6 +126,9 @@ export class BlockInputSync { this.chain.emitter.on(ChainEvent.unknownBlockRoot, this.onUnknownBlockRoot); this.chain.emitter.on(ChainEvent.incompleteBlockInput, this.onIncompleteBlockInput); this.chain.emitter.on(ChainEvent.unknownParent, this.onUnknownParent); + this.chain.emitter.on(ChainEvent.unknownPayloadEnvelope, this.onUnknownPayloadEnvelope); + this.chain.emitter.on(ChainEvent.unknownParentPayload, this.onUnknownParentPayload); + this.chain.emitter.on(routes.events.EventType.executionPayloadAvailable, this.onPayloadImported); this.network.events.on(NetworkEvent.peerConnected, this.onPeerConnected); this.network.events.on(NetworkEvent.peerDisconnected, this.onPeerDisconnected); this.subscribedToNetworkEvents = true; @@ -127,6 +140,9 @@ export class BlockInputSync { this.chain.emitter.off(ChainEvent.unknownBlockRoot, this.onUnknownBlockRoot); this.chain.emitter.off(ChainEvent.incompleteBlockInput, this.onIncompleteBlockInput); this.chain.emitter.off(ChainEvent.unknownParent, this.onUnknownParent); + this.chain.emitter.off(ChainEvent.unknownPayloadEnvelope, this.onUnknownPayloadEnvelope); + this.chain.emitter.off(ChainEvent.unknownParentPayload, this.onUnknownParentPayload); + this.chain.emitter.off(routes.events.EventType.executionPayloadAvailable, this.onPayloadImported); this.network.events.off(NetworkEvent.peerConnected, this.onPeerConnected); this.network.events.off(NetworkEvent.peerDisconnected, this.onPeerDisconnected); this.subscribedToNetworkEvents = false; @@ -183,6 +199,66 @@ export class BlockInputSync { } }; + private onUnknownPayloadEnvelope = (data: ChainEventData[ChainEvent.unknownPayloadEnvelope]): void => { + try { + const {blockRootHex, peer} = data; + const block = this.chain.forkChoice.getBlockHexDefaultStatus(blockRootHex); + if (!block) return; + this.addPendingPayload(blockRootHex, block.slot, peer); + this.metrics?.blockInputSync.payloadRequests.inc({source: data.source}); + } catch (e) { + this.logger.debug("Error handling unknownPayloadBlockRoot event", {}, e as Error); + } + }; + + private onUnknownParentPayload = (data: ChainEventData[ChainEvent.unknownParentPayload]): void => { + try { + const {blockInput, peer} = data; + // Track the child block (don't add parent to pendingBlocks — it's already in fork choice) + this.addByBlockInput(blockInput, peer); + // Trigger parent payload fetch + const parentRootHex = blockInput.parentRootHex; + const parentBlock = this.chain.forkChoice.getBlockHexDefaultStatus(parentRootHex); + if (parentBlock) { + this.addPendingPayload(parentRootHex, parentBlock.slot, peer); + } + this.triggerUnknownBlockSearch(); + this.metrics?.blockInputSync.requests.inc({type: PendingBlockType.UNKNOWN_PARENT}); + this.metrics?.blockInputSync.source.inc({source: data.source}); + } catch (e) { + this.logger.debug("Error handling unknownParentPayload event", {}, e as Error); + } + }; + + private onPayloadImported = ({blockRoot}: {slot: number; blockRoot: string}): void => { + this.pendingPayloads.delete(blockRoot); + this.triggerUnknownBlockSearch(); + // TODO GLOAS (PR #9025): call networkProcessor.onPayloadProcessed({root: blockRoot}) + // to flush awaitingMessagesByEnvelopeBlockRoot gossip queue + }; + + addPendingPayload(rootHex: RootHex, slot: number, peer?: PeerIdStr): void { + const payloadInput = this.chain.seenPayloadEnvelopeInputCache.get(rootHex); + if (payloadInput?.hasPayloadEnvelope()) return; + if (this.chain.forkChoice.getBlockHex(rootHex, PayloadStatus.FULL)) return; + if (this.pendingPayloads.size >= MAX_PENDING_PAYLOADS) return; + + let pending = this.pendingPayloads.get(rootHex); + if (!pending) { + pending = { + status: "pending", + blockRootHex: rootHex, + slot, + attempts: 0, + peerIdStrings: new Set(), + timeAddedSec: Date.now() / 1000, + }; + this.pendingPayloads.set(rootHex, pending); + } + if (peer) pending.peerIdStrings.add(peer); + this.triggerPayloadSearch(); + } + private addByRootHex = (rootHex: RootHex, peerIdStr?: PeerIdStr): void => { let pendingBlock = this.pendingBlocks.get(rootHex); if (!pendingBlock) { @@ -248,6 +324,7 @@ export class BlockInputSync { const peerSyncMeta = this.network.getConnectedPeerSyncMeta(peerId); this.peerBalancer.onPeerConnected(data.peer, peerSyncMeta); this.triggerUnknownBlockSearch(); + this.triggerPayloadSearch(); } catch (e) { this.logger.debug("Error handling peerConnected event", {}, e as Error); } @@ -282,7 +359,10 @@ export class BlockInputSync { for (const block of ancestors) { // when this happens, it's likely the block and parent block are processed by head sync - if (this.chain.forkChoice.hasBlockHex(block.blockInput.parentRootHex)) { + if ( + this.chain.forkChoice.hasBlockHex(block.blockInput.parentRootHex) && + this.isParentPayloadResolved(block.blockInput) + ) { processedBlocks++; this.processBlock(block).catch((e) => { this.logger.debug("Unexpected error - process old downloaded block", {}, e); @@ -341,7 +421,7 @@ export class BlockInputSync { }; this.logger.verbose("Downloaded unknown block", logCtx2); - if (parentInForkChoice) { + if (parentInForkChoice && this.isParentPayloadResolved(pending.blockInput)) { // Bingo! Process block. Add to pending blocks anyway for recycle the cache that prevents duplicate processing this.processBlock(pending).catch((e) => { this.logger.debug("Unexpected error - process newly downloaded block", logCtx2, e); @@ -441,6 +521,9 @@ export class BlockInputSync { }); } } + + // Re-check pending payloads — block import may have revealed new payload evidence + this.triggerPayloadSearch(); } else { const errorData = {slot: pendingBlock.blockInput.slot, root: pendingBlock.blockInput.blockRootHex}; if (res.err instanceof BlockError) { @@ -456,6 +539,12 @@ export class BlockInputSync { pendingBlock.status = PendingBlockInputStatus.downloaded; break; + case BlockErrorCode.PARENT_PAYLOAD_UNKNOWN: + this.logger.debug("Block parent payload unknown", errorData, res.err); + this.addPendingPayload(res.err.type.parentRoot, pendingBlock.blockInput.slot - 1); + pendingBlock.status = PendingBlockInputStatus.downloaded; + break; + case BlockErrorCode.EXECUTION_ENGINE_ERROR: // Removing the block(s) without penalizing the peers, hoping for EL to // recover on a latter download + verify attempt @@ -477,6 +566,104 @@ export class BlockInputSync { } } + /** + * Check whether the parent block's FULL payload variant is available for a Gloas block. + * Returns true for pre-Gloas blocks or if the parent's FULL variant is already in fork choice. + * If the parent payload is missing but the parent block is known, triggers a pending payload fetch. + */ + private isParentPayloadResolved(blockInput: IBlockInput): boolean { + const fork = this.config.getForkName(blockInput.slot); + if (!isForkPostGloas(fork)) return true; + + const block = blockInput.getBlock().message; + if (!isGloasBeaconBlock(block)) return true; + + const parentRootHex = blockInput.parentRootHex; + const parentBlockHash = toRootHex(block.body.signedExecutionPayloadBid.message.parentBlockHash); + + if (this.chain.forkChoice.getBlockHexAndBlockHash(parentRootHex, parentBlockHash) !== null) { + return true; + } + + // Parent block known but FULL variant missing — trigger payload fetch + const parentBlock = this.chain.forkChoice.getBlockHexDefaultStatus(parentRootHex); + if (parentBlock) { + this.addPendingPayload(parentRootHex, parentBlock.slot); + } + return false; + } + + private triggerPayloadSearch = (): void => { + if (this.pendingPayloads.size === 0) return; + if (this.network.getConnectedPeers().length === 0) return; + + const finalizedSlot = this.chain.forkChoice.getFinalizedBlock().slot; + + for (const [rootHex, pending] of this.pendingPayloads) { + // Prune payloads for finalized slots + if (pending.slot <= finalizedSlot) { + this.pendingPayloads.delete(rootHex); + continue; + } + if ( + this.chain.seenPayloadEnvelopeInputCache.get(rootHex)?.hasPayloadEnvelope() || + this.chain.forkChoice.getBlockHex(rootHex, PayloadStatus.FULL) + ) { + this.pendingPayloads.delete(rootHex); + continue; + } + if (!this.chain.forkChoice.hasBlockHexUnsafe(rootHex)) continue; + if (pending.status !== "pending") continue; + if (pending.attempts >= MAX_ATTEMPTS_PER_PAYLOAD) { + this.pendingPayloads.delete(rootHex); + continue; + } + this.fetchPayloadEnvelope(pending).catch((e) => { + this.logger.debug("Unexpected error - fetchPayloadEnvelope", {root: pending.blockRootHex}, e); + }); + } + }; + + private async fetchPayloadEnvelope(pending: PendingPayloadEnvelope): Promise { + pending.status = "fetching"; + pending.attempts++; + try { + const peer = this.peerBalancer.bestPeerForPendingColumns(new Set(), new Set())?.peerId; + if (!peer) { + pending.status = "pending"; + return; + } + + const envelopes = await this.network.sendExecutionPayloadEnvelopesByRoot(peer, [fromHex(pending.blockRootHex)]); + if (envelopes.length === 0) { + pending.status = "pending"; + return; + } + + const envelope = envelopes[0]; + const payloadInput = this.chain.seenPayloadEnvelopeInputCache.get(pending.blockRootHex); + if (!payloadInput) { + throw Error(`PayloadEnvelopeInput missing for known block root ${pending.blockRootHex}`); + } + + payloadInput.addPayloadEnvelope({ + envelope, + source: PayloadEnvelopeInputSource.byRoot, + seenTimestampSec: Date.now() / 1000, + peerIdStr: peer, + }); + + if (payloadInput.isComplete()) { + await this.chain.processExecutionPayload(payloadInput); + } + this.metrics?.blockInputSync.payloadFetchSuccess.inc(); + } catch (e) { + this.logger.debug("Error fetching payload envelope", {root: pending.blockRootHex}, e as Error); + pending.status = "pending"; + this.metrics?.blockInputSync.payloadFetchError.inc(); + } + } + /** * From a set of shuffled peers: * - fetch the block @@ -671,6 +858,7 @@ export class BlockInputSync { for (const block of badPendingBlocks) { const rootHex = getBlockInputSyncCacheItemRootHex(block); this.pendingBlocks.delete(rootHex); + this.pendingPayloads.delete(rootHex); this.chain.seenBlockInputCache.prune(rootHex); this.logger.debug("Removing bad/unknown/incomplete BlockInputSyncCacheItem", { slot, diff --git a/packages/beacon-node/src/util/types.ts b/packages/beacon-node/src/util/types.ts index fc3880c616e9..10a26de4838f 100644 --- a/packages/beacon-node/src/util/types.ts +++ b/packages/beacon-node/src/util/types.ts @@ -29,3 +29,9 @@ export type BlobSidecarsByRootRequest = ValueOf new ListCompositeType(ssz.fulu.DataColumnsByRootIdentifier, config.MAX_REQUEST_BLOCKS_DENEB); export type DataColumnSidecarsByRootRequest = ValueOf>; + +export const ExecutionPayloadEnvelopesByRootRequestType = (config: BeaconConfig) => + new ListCompositeType(ssz.Root, config.MAX_REQUEST_PAYLOADS); +export type ExecutionPayloadEnvelopesByRootRequest = ValueOf< + ReturnType +>; diff --git a/packages/types/src/gloas/sszTypes.ts b/packages/types/src/gloas/sszTypes.ts index 38f1af75f369..b9bf89ad583f 100644 --- a/packages/types/src/gloas/sszTypes.ts +++ b/packages/types/src/gloas/sszTypes.ts @@ -282,3 +282,8 @@ export const DataColumnSidecar = new ContainerType( ); export const DataColumnSidecars = new ListCompositeType(DataColumnSidecar, NUMBER_OF_COLUMNS); + +export const ExecutionPayloadEnvelopesByRangeRequest = new ContainerType( + {startSlot: Slot, count: UintNum64}, + {typeName: "ExecutionPayloadEnvelopesByRangeRequest", jsonCase: "eth2"} +); diff --git a/packages/types/src/gloas/types.ts b/packages/types/src/gloas/types.ts index 6ef793c2e8cf..f90e74bbf416 100644 --- a/packages/types/src/gloas/types.ts +++ b/packages/types/src/gloas/types.ts @@ -21,3 +21,5 @@ export type BeaconState = ValueOf; export type DataColumnSidecar = ValueOf; export type DataColumnSidecars = ValueOf; + +export type ExecutionPayloadEnvelopesByRangeRequest = ValueOf;