From db5de4741b60118c61860447480a0c30237d0627 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Tue, 1 Sep 2026 01:48:48 +0100 Subject: [PATCH 01/11] feat(builder): add Engine payload source --- .../builder/src/services/payloadSource.ts | 157 +++++++++++++++++ .../test/unit/services/payloadSource.test.ts | 159 ++++++++++++++++++ 2 files changed, 316 insertions(+) create mode 100644 packages/builder/src/services/payloadSource.ts create mode 100644 packages/builder/test/unit/services/payloadSource.test.ts diff --git a/packages/builder/src/services/payloadSource.ts b/packages/builder/src/services/payloadSource.ts new file mode 100644 index 000000000000..edd4e09e9190 --- /dev/null +++ b/packages/builder/src/services/payloadSource.ts @@ -0,0 +1,157 @@ +import type {ForkPostGloas} from "@lodestar/params"; +import type {BlobsBundle, ExecutionPayload, ExecutionRequests, RootHex, gloas} from "@lodestar/types"; +import {LodestarError} from "@lodestar/utils"; + +export type PayloadId = string; + +export type ForkchoiceState = { + headBlockHash: RootHex; + safeBlockHash: RootHex; + finalizedBlockHash: RootHex; +}; + +export type BuildRequest = { + fork: ForkPostGloas; + forkchoiceState: ForkchoiceState; + payloadAttributes: gloas.PayloadAttributes; +}; + +export type BuildHandle = { + sourceId: string; + fork: ForkPostGloas; + payloadId: PayloadId; +}; + +export type BuiltPayload = { + sourceId: string; + fork: ForkPostGloas; + executionPayload: ExecutionPayload; + executionRequests: ExecutionRequests; + blobsBundle: BlobsBundle; + executionPayloadValue: bigint; +}; + +export type EnginePayloadResult = { + executionPayload: ExecutionPayload; + executionPayloadValue: bigint; + blobsBundle?: BlobsBundle; + executionRequests?: ExecutionRequests; +}; + +/** Narrow Engine API boundary required by an Engine-backed payload source. */ +export interface PayloadSourceEngine { + notifyForkchoiceUpdate( + fork: ForkPostGloas, + headBlockHash: RootHex, + safeBlockHash: RootHex, + finalizedBlockHash: RootHex, + payloadAttributes: gloas.PayloadAttributes + ): Promise; + getPayload(fork: ForkPostGloas, payloadId: PayloadId): Promise; +} + +/** Source that prepares and retrieves complete execution payloads. */ +export interface PayloadSource { + readonly id: string; + prepare(request: BuildRequest): Promise; + getPayload(handle: BuildHandle): Promise; +} + +export enum PayloadSourceErrorCode { + NO_PAYLOAD_ID = "PAYLOAD_SOURCE_ERROR_NO_PAYLOAD_ID", + SOURCE_MISMATCH = "PAYLOAD_SOURCE_ERROR_SOURCE_MISMATCH", + MISSING_BLOBS_BUNDLE = "PAYLOAD_SOURCE_ERROR_MISSING_BLOBS_BUNDLE", + MISSING_EXECUTION_REQUESTS = "PAYLOAD_SOURCE_ERROR_MISSING_EXECUTION_REQUESTS", +} + +export type PayloadSourceErrorType = + | {code: PayloadSourceErrorCode.NO_PAYLOAD_ID; sourceId: string} + | { + code: PayloadSourceErrorCode.SOURCE_MISMATCH; + sourceId: string; + handleSourceId: string; + } + | { + code: PayloadSourceErrorCode.MISSING_BLOBS_BUNDLE | PayloadSourceErrorCode.MISSING_EXECUTION_REQUESTS; + sourceId: string; + payloadId: PayloadId; + }; + +export class PayloadSourceError extends LodestarError {} + +/** Payload source backed by an execution client's Engine API. */ +export class EnginePayloadSource implements PayloadSource { + constructor( + readonly id: string, + private readonly engine: PayloadSourceEngine + ) {} + + async prepare(request: BuildRequest): Promise { + const {headBlockHash, safeBlockHash, finalizedBlockHash} = request.forkchoiceState; + const payloadId = await this.engine.notifyForkchoiceUpdate( + request.fork, + headBlockHash, + safeBlockHash, + finalizedBlockHash, + request.payloadAttributes + ); + + if (payloadId === null) { + throw new PayloadSourceError( + {code: PayloadSourceErrorCode.NO_PAYLOAD_ID, sourceId: this.id}, + `Execution client did not return a payload ID sourceId=${this.id}` + ); + } + + return {sourceId: this.id, fork: request.fork, payloadId}; + } + + async getPayload(handle: BuildHandle): Promise { + if (handle.sourceId !== this.id) { + throw new PayloadSourceError( + { + code: PayloadSourceErrorCode.SOURCE_MISMATCH, + sourceId: this.id, + handleSourceId: handle.sourceId, + }, + `Payload handle belongs to another source sourceId=${this.id} handleSourceId=${handle.sourceId}` + ); + } + + const {executionPayload, executionPayloadValue, blobsBundle, executionRequests} = await this.engine.getPayload( + handle.fork, + handle.payloadId + ); + + if (blobsBundle === undefined) { + throw new PayloadSourceError( + { + code: PayloadSourceErrorCode.MISSING_BLOBS_BUNDLE, + sourceId: this.id, + payloadId: handle.payloadId, + }, + `Execution client did not return a blobs bundle sourceId=${this.id} payloadId=${handle.payloadId}` + ); + } + + if (executionRequests === undefined) { + throw new PayloadSourceError( + { + code: PayloadSourceErrorCode.MISSING_EXECUTION_REQUESTS, + sourceId: this.id, + payloadId: handle.payloadId, + }, + `Execution client did not return execution requests sourceId=${this.id} payloadId=${handle.payloadId}` + ); + } + + return { + sourceId: this.id, + fork: handle.fork, + executionPayload: executionPayload as ExecutionPayload, + executionRequests: executionRequests as ExecutionRequests, + blobsBundle: blobsBundle as BlobsBundle, + executionPayloadValue, + }; + } +} diff --git a/packages/builder/test/unit/services/payloadSource.test.ts b/packages/builder/test/unit/services/payloadSource.test.ts new file mode 100644 index 000000000000..0dc48c7cbb89 --- /dev/null +++ b/packages/builder/test/unit/services/payloadSource.test.ts @@ -0,0 +1,159 @@ +import {Mocked, beforeEach, describe, expect, it, vi} from "vitest"; +import {ForkName} from "@lodestar/params"; +import {ssz} from "@lodestar/types"; +import {ErrorAborted, TimeoutError, toRootHex} from "@lodestar/utils"; +import { + BuildHandle, + BuildRequest, + EnginePayloadResult, + EnginePayloadSource, + PayloadSourceEngine, + PayloadSourceError, + PayloadSourceErrorCode, +} from "../../../src/services/payloadSource.js"; + +describe("EnginePayloadSource", () => { + const sourceId = "engine-0"; + const payloadId = "0x0102030405060708"; + const forkchoiceState = { + headBlockHash: toRootHex(Uint8Array.from({length: 32}, () => 1)), + safeBlockHash: toRootHex(Uint8Array.from({length: 32}, () => 2)), + finalizedBlockHash: toRootHex(Uint8Array.from({length: 32}, () => 3)), + }; + const payloadAttributes = ssz.gloas.PayloadAttributes.defaultValue(); + const request: BuildRequest = {fork: ForkName.gloas, forkchoiceState, payloadAttributes}; + const handle: BuildHandle = {sourceId, fork: ForkName.gloas, payloadId}; + + let engine: Mocked; + let source: EnginePayloadSource; + + beforeEach(() => { + engine = { + notifyForkchoiceUpdate: vi.fn(), + getPayload: vi.fn(), + }; + source = new EnginePayloadSource(sourceId, engine); + }); + + it("prepares a payload and returns a source-bound handle", async () => { + engine.notifyForkchoiceUpdate.mockResolvedValue(payloadId); + + const result = await source.prepare(request); + + expect(engine.notifyForkchoiceUpdate).toHaveBeenCalledWith( + ForkName.gloas, + forkchoiceState.headBlockHash, + forkchoiceState.safeBlockHash, + forkchoiceState.finalizedBlockHash, + payloadAttributes + ); + expect(result).toEqual(handle); + }); + + it("supports post-Gloas forks without narrowing the fork", async () => { + engine.notifyForkchoiceUpdate.mockResolvedValue(payloadId); + engine.getPayload.mockResolvedValue(getEnginePayloadResult()); + + const result = await source.prepare({...request, fork: ForkName.heze}); + const builtPayload = await source.getPayload(result); + + expect(result.fork).toBe(ForkName.heze); + expect(engine.notifyForkchoiceUpdate).toHaveBeenCalledWith( + ForkName.heze, + forkchoiceState.headBlockHash, + forkchoiceState.safeBlockHash, + forkchoiceState.finalizedBlockHash, + payloadAttributes + ); + expect(engine.getPayload).toHaveBeenCalledWith(ForkName.heze, payloadId); + expect(builtPayload.fork).toBe(ForkName.heze); + }); + + it("rejects a missing payload ID with a structured error", async () => { + engine.notifyForkchoiceUpdate.mockResolvedValue(null); + + const error = await getPayloadSourceError(source.prepare(request)); + + expect(error.type).toEqual({code: PayloadSourceErrorCode.NO_PAYLOAD_ID, sourceId}); + }); + + it.each([ + ["transport", new Error("connection reset")], + ["timeout", new TimeoutError("engine_forkchoiceUpdatedV4")], + ["cancellation", new ErrorAborted("engine_forkchoiceUpdatedV4")], + ])("propagates %s errors from payload preparation", async (_name, error) => { + engine.notifyForkchoiceUpdate.mockRejectedValue(error); + + await expect(source.prepare(request)).rejects.toBe(error); + }); + + it("retrieves a complete payload without rebuilding exact-width values", async () => { + const result = getEnginePayloadResult(); + engine.getPayload.mockResolvedValue(result); + + const builtPayload = await source.getPayload(handle); + + expect(engine.getPayload).toHaveBeenCalledWith(ForkName.gloas, payloadId); + expect(builtPayload.sourceId).toBe(sourceId); + expect(builtPayload.fork).toBe(ForkName.gloas); + expect(builtPayload.executionPayload).toBe(result.executionPayload); + expect(builtPayload.blobsBundle).toBe(result.blobsBundle); + expect(builtPayload.executionRequests).toBe(result.executionRequests); + expect(builtPayload.executionPayloadValue).toBe(result.executionPayloadValue); + }); + + it("rejects a handle belonging to another source before calling the Engine API", async () => { + const error = await getPayloadSourceError(source.getPayload({...handle, sourceId: "engine-1"})); + + expect(error.type).toEqual({ + code: PayloadSourceErrorCode.SOURCE_MISMATCH, + sourceId, + handleSourceId: "engine-1", + }); + expect(engine.getPayload).not.toHaveBeenCalled(); + }); + + it("rejects a response without a blobs bundle", async () => { + engine.getPayload.mockResolvedValue({...getEnginePayloadResult(), blobsBundle: undefined}); + + const error = await getPayloadSourceError(source.getPayload(handle)); + + expect(error.type).toEqual({code: PayloadSourceErrorCode.MISSING_BLOBS_BUNDLE, sourceId, payloadId}); + }); + + it("rejects a response without execution requests", async () => { + engine.getPayload.mockResolvedValue({...getEnginePayloadResult(), executionRequests: undefined}); + + const error = await getPayloadSourceError(source.getPayload(handle)); + + expect(error.type).toEqual({code: PayloadSourceErrorCode.MISSING_EXECUTION_REQUESTS, sourceId, payloadId}); + }); + + it("propagates retrieval errors without replacing their type", async () => { + const error = new TimeoutError("engine_getPayloadV6"); + engine.getPayload.mockRejectedValue(error); + + await expect(source.getPayload(handle)).rejects.toBe(error); + }); +}); + +function getEnginePayloadResult(): EnginePayloadResult { + return { + executionPayload: ssz.gloas.ExecutionPayload.defaultValue(), + blobsBundle: ssz.gloas.BlobsBundle.defaultValue(), + executionRequests: ssz.gloas.ExecutionRequests.defaultValue(), + executionPayloadValue: 12_345_678_901_234_567_890n, + }; +} + +async function getPayloadSourceError(promise: Promise): Promise { + try { + await promise; + throw Error("Expected PayloadSourceError"); + } catch (error) { + if (!(error instanceof PayloadSourceError)) { + throw error; + } + return error; + } +} From dcb81cb2b4c9b6594cecd5a7c712066c1ba3922b Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Tue, 1 Sep 2026 01:53:31 +0100 Subject: [PATCH 02/11] test(builder): cover unsupported payload source errors --- packages/builder/test/unit/services/payloadSource.test.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/packages/builder/test/unit/services/payloadSource.test.ts b/packages/builder/test/unit/services/payloadSource.test.ts index 0dc48c7cbb89..8f274b9145e5 100644 --- a/packages/builder/test/unit/services/payloadSource.test.ts +++ b/packages/builder/test/unit/services/payloadSource.test.ts @@ -79,6 +79,7 @@ describe("EnginePayloadSource", () => { it.each([ ["transport", new Error("connection reset")], + ["unsupported Engine response", new Error("Method not found")], ["timeout", new TimeoutError("engine_forkchoiceUpdatedV4")], ["cancellation", new ErrorAborted("engine_forkchoiceUpdatedV4")], ])("propagates %s errors from payload preparation", async (_name, error) => { From b084a2e5beb75534adee53b6ffe84c427546946f Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Tue, 1 Sep 2026 02:41:56 +0100 Subject: [PATCH 03/11] refactor(builder): preserve fork-specific payload attributes --- .../builder/src/services/payloadSource.ts | 50 ++++++++++--------- .../test/unit/services/payloadSource.test.ts | 22 ++++++-- 2 files changed, 43 insertions(+), 29 deletions(-) diff --git a/packages/builder/src/services/payloadSource.ts b/packages/builder/src/services/payloadSource.ts index edd4e09e9190..62fb8ee9a838 100644 --- a/packages/builder/src/services/payloadSource.ts +++ b/packages/builder/src/services/payloadSource.ts @@ -1,5 +1,5 @@ import type {ForkPostGloas} from "@lodestar/params"; -import type {BlobsBundle, ExecutionPayload, ExecutionRequests, RootHex, gloas} from "@lodestar/types"; +import type {BlobsBundle, ExecutionPayload, ExecutionRequests, RootHex, SSEPayloadAttributes} from "@lodestar/types"; import {LodestarError} from "@lodestar/utils"; export type PayloadId = string; @@ -10,51 +10,53 @@ export type ForkchoiceState = { finalizedBlockHash: RootHex; }; -export type BuildRequest = { - fork: ForkPostGloas; +export type PayloadAttributes = SSEPayloadAttributes["payloadAttributes"]; + +export type BuildRequest = { + fork: F; forkchoiceState: ForkchoiceState; - payloadAttributes: gloas.PayloadAttributes; + payloadAttributes: PayloadAttributes; }; -export type BuildHandle = { +export type BuildHandle = { sourceId: string; - fork: ForkPostGloas; + fork: F; payloadId: PayloadId; }; -export type BuiltPayload = { +export type BuiltPayload = { sourceId: string; - fork: ForkPostGloas; - executionPayload: ExecutionPayload; - executionRequests: ExecutionRequests; - blobsBundle: BlobsBundle; + fork: F; + executionPayload: ExecutionPayload; + executionRequests: ExecutionRequests; + blobsBundle: BlobsBundle; executionPayloadValue: bigint; }; export type EnginePayloadResult = { - executionPayload: ExecutionPayload; + executionPayload: ExecutionPayload; executionPayloadValue: bigint; - blobsBundle?: BlobsBundle; - executionRequests?: ExecutionRequests; + blobsBundle?: BlobsBundle; + executionRequests?: ExecutionRequests; }; -/** Narrow Engine API boundary required by an Engine-backed payload source. */ +/** Narrow Engine API boundary whose transport owns request retries, timeouts, and Builder-lifetime cancellation. */ export interface PayloadSourceEngine { notifyForkchoiceUpdate( fork: ForkPostGloas, headBlockHash: RootHex, safeBlockHash: RootHex, finalizedBlockHash: RootHex, - payloadAttributes: gloas.PayloadAttributes + payloadAttributes: PayloadAttributes ): Promise; getPayload(fork: ForkPostGloas, payloadId: PayloadId): Promise; } -/** Source that prepares and retrieves complete execution payloads. */ +/** Source that prepares and retrieves complete execution payloads without owning build scheduling policy. */ export interface PayloadSource { readonly id: string; - prepare(request: BuildRequest): Promise; - getPayload(handle: BuildHandle): Promise; + prepare(request: BuildRequest): Promise>; + getPayload(handle: BuildHandle): Promise>; } export enum PayloadSourceErrorCode { @@ -86,7 +88,7 @@ export class EnginePayloadSource implements PayloadSource { private readonly engine: PayloadSourceEngine ) {} - async prepare(request: BuildRequest): Promise { + async prepare(request: BuildRequest): Promise> { const {headBlockHash, safeBlockHash, finalizedBlockHash} = request.forkchoiceState; const payloadId = await this.engine.notifyForkchoiceUpdate( request.fork, @@ -106,7 +108,7 @@ export class EnginePayloadSource implements PayloadSource { return {sourceId: this.id, fork: request.fork, payloadId}; } - async getPayload(handle: BuildHandle): Promise { + async getPayload(handle: BuildHandle): Promise> { if (handle.sourceId !== this.id) { throw new PayloadSourceError( { @@ -148,9 +150,9 @@ export class EnginePayloadSource implements PayloadSource { return { sourceId: this.id, fork: handle.fork, - executionPayload: executionPayload as ExecutionPayload, - executionRequests: executionRequests as ExecutionRequests, - blobsBundle: blobsBundle as BlobsBundle, + executionPayload: executionPayload as ExecutionPayload, + executionRequests: executionRequests as ExecutionRequests, + blobsBundle: blobsBundle as BlobsBundle, executionPayloadValue, }; } diff --git a/packages/builder/test/unit/services/payloadSource.test.ts b/packages/builder/test/unit/services/payloadSource.test.ts index 8f274b9145e5..3be79772e610 100644 --- a/packages/builder/test/unit/services/payloadSource.test.ts +++ b/packages/builder/test/unit/services/payloadSource.test.ts @@ -21,8 +21,8 @@ describe("EnginePayloadSource", () => { finalizedBlockHash: toRootHex(Uint8Array.from({length: 32}, () => 3)), }; const payloadAttributes = ssz.gloas.PayloadAttributes.defaultValue(); - const request: BuildRequest = {fork: ForkName.gloas, forkchoiceState, payloadAttributes}; - const handle: BuildHandle = {sourceId, fork: ForkName.gloas, payloadId}; + const request: BuildRequest = {fork: ForkName.gloas, forkchoiceState, payloadAttributes}; + const handle: BuildHandle = {sourceId, fork: ForkName.gloas, payloadId}; let engine: Mocked; let source: EnginePayloadSource; @@ -51,10 +51,21 @@ describe("EnginePayloadSource", () => { }); it("supports post-Gloas forks without narrowing the fork", async () => { + const hezePayloadAttributes = ssz.heze.PayloadAttributes.defaultValue(); + hezePayloadAttributes.inclusionListTransactions = [Uint8Array.from([1, 2, 3])]; engine.notifyForkchoiceUpdate.mockResolvedValue(payloadId); - engine.getPayload.mockResolvedValue(getEnginePayloadResult()); + engine.getPayload.mockResolvedValue({ + executionPayload: ssz.heze.ExecutionPayload.defaultValue(), + blobsBundle: ssz.heze.BlobsBundle.defaultValue(), + executionRequests: ssz.heze.ExecutionRequests.defaultValue(), + executionPayloadValue: 12_345_678_901_234_567_890n, + }); - const result = await source.prepare({...request, fork: ForkName.heze}); + const result = await source.prepare({ + fork: ForkName.heze, + forkchoiceState, + payloadAttributes: hezePayloadAttributes, + }); const builtPayload = await source.getPayload(result); expect(result.fork).toBe(ForkName.heze); @@ -63,10 +74,11 @@ describe("EnginePayloadSource", () => { forkchoiceState.headBlockHash, forkchoiceState.safeBlockHash, forkchoiceState.finalizedBlockHash, - payloadAttributes + hezePayloadAttributes ); expect(engine.getPayload).toHaveBeenCalledWith(ForkName.heze, payloadId); expect(builtPayload.fork).toBe(ForkName.heze); + expect(hezePayloadAttributes.inclusionListTransactions).toHaveLength(1); }); it("rejects a missing payload ID with a structured error", async () => { From 9513820f322162bf94b8d8fc95fa464b2e3b24a2 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Tue, 1 Sep 2026 02:44:09 +0100 Subject: [PATCH 04/11] feat(builder): carry custody columns in build requests --- packages/builder/src/services/payloadSource.ts | 16 +++++++++++++--- .../test/unit/services/payloadSource.test.ts | 15 ++++++++++++--- 2 files changed, 25 insertions(+), 6 deletions(-) diff --git a/packages/builder/src/services/payloadSource.ts b/packages/builder/src/services/payloadSource.ts index 62fb8ee9a838..68cba99fe03c 100644 --- a/packages/builder/src/services/payloadSource.ts +++ b/packages/builder/src/services/payloadSource.ts @@ -1,5 +1,12 @@ import type {ForkPostGloas} from "@lodestar/params"; -import type {BlobsBundle, ExecutionPayload, ExecutionRequests, RootHex, SSEPayloadAttributes} from "@lodestar/types"; +import type { + BlobsBundle, + ColumnIndex, + ExecutionPayload, + ExecutionRequests, + RootHex, + SSEPayloadAttributes, +} from "@lodestar/types"; import {LodestarError} from "@lodestar/utils"; export type PayloadId = string; @@ -16,6 +23,7 @@ export type BuildRequest = { fork: F; forkchoiceState: ForkchoiceState; payloadAttributes: PayloadAttributes; + custodyColumns: ColumnIndex[]; }; export type BuildHandle = { @@ -47,7 +55,8 @@ export interface PayloadSourceEngine { headBlockHash: RootHex, safeBlockHash: RootHex, finalizedBlockHash: RootHex, - payloadAttributes: PayloadAttributes + payloadAttributes: PayloadAttributes, + custodyColumns: ColumnIndex[] ): Promise; getPayload(fork: ForkPostGloas, payloadId: PayloadId): Promise; } @@ -95,7 +104,8 @@ export class EnginePayloadSource implements PayloadSource { headBlockHash, safeBlockHash, finalizedBlockHash, - request.payloadAttributes + request.payloadAttributes, + request.custodyColumns ); if (payloadId === null) { diff --git a/packages/builder/test/unit/services/payloadSource.test.ts b/packages/builder/test/unit/services/payloadSource.test.ts index 3be79772e610..2a13e85b09fa 100644 --- a/packages/builder/test/unit/services/payloadSource.test.ts +++ b/packages/builder/test/unit/services/payloadSource.test.ts @@ -21,7 +21,13 @@ describe("EnginePayloadSource", () => { finalizedBlockHash: toRootHex(Uint8Array.from({length: 32}, () => 3)), }; const payloadAttributes = ssz.gloas.PayloadAttributes.defaultValue(); - const request: BuildRequest = {fork: ForkName.gloas, forkchoiceState, payloadAttributes}; + const custodyColumns = [0, 3, 127]; + const request: BuildRequest = { + fork: ForkName.gloas, + forkchoiceState, + payloadAttributes, + custodyColumns, + }; const handle: BuildHandle = {sourceId, fork: ForkName.gloas, payloadId}; let engine: Mocked; @@ -45,7 +51,8 @@ describe("EnginePayloadSource", () => { forkchoiceState.headBlockHash, forkchoiceState.safeBlockHash, forkchoiceState.finalizedBlockHash, - payloadAttributes + payloadAttributes, + custodyColumns ); expect(result).toEqual(handle); }); @@ -65,6 +72,7 @@ describe("EnginePayloadSource", () => { fork: ForkName.heze, forkchoiceState, payloadAttributes: hezePayloadAttributes, + custodyColumns, }); const builtPayload = await source.getPayload(result); @@ -74,7 +82,8 @@ describe("EnginePayloadSource", () => { forkchoiceState.headBlockHash, forkchoiceState.safeBlockHash, forkchoiceState.finalizedBlockHash, - hezePayloadAttributes + hezePayloadAttributes, + custodyColumns ); expect(engine.getPayload).toHaveBeenCalledWith(ForkName.heze, payloadId); expect(builtPayload.fork).toBe(ForkName.heze); From f6b60531623a48d677f30c300908c4375f3c3d22 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Tue, 1 Sep 2026 03:00:39 +0100 Subject: [PATCH 05/11] refactor(builder): preserve Engine fork types --- .../builder/src/services/payloadSource.ts | 22 +++---- .../test/unit/services/payloadSource.test.ts | 57 ++++++++++++------- 2 files changed, 46 insertions(+), 33 deletions(-) diff --git a/packages/builder/src/services/payloadSource.ts b/packages/builder/src/services/payloadSource.ts index 68cba99fe03c..4b14fae4ed5e 100644 --- a/packages/builder/src/services/payloadSource.ts +++ b/packages/builder/src/services/payloadSource.ts @@ -41,24 +41,24 @@ export type BuiltPayload = { executionPayloadValue: bigint; }; -export type EnginePayloadResult = { - executionPayload: ExecutionPayload; +export type EnginePayloadResult = { + executionPayload: ExecutionPayload; executionPayloadValue: bigint; - blobsBundle?: BlobsBundle; - executionRequests?: ExecutionRequests; + blobsBundle?: BlobsBundle; + executionRequests?: ExecutionRequests; }; /** Narrow Engine API boundary whose transport owns request retries, timeouts, and Builder-lifetime cancellation. */ export interface PayloadSourceEngine { - notifyForkchoiceUpdate( - fork: ForkPostGloas, + notifyForkchoiceUpdate( + fork: F, headBlockHash: RootHex, safeBlockHash: RootHex, finalizedBlockHash: RootHex, - payloadAttributes: PayloadAttributes, + payloadAttributes: PayloadAttributes, custodyColumns: ColumnIndex[] ): Promise; - getPayload(fork: ForkPostGloas, payloadId: PayloadId): Promise; + getPayload(fork: F, payloadId: PayloadId): Promise>; } /** Source that prepares and retrieves complete execution payloads without owning build scheduling policy. */ @@ -160,9 +160,9 @@ export class EnginePayloadSource implements PayloadSource { return { sourceId: this.id, fork: handle.fork, - executionPayload: executionPayload as ExecutionPayload, - executionRequests: executionRequests as ExecutionRequests, - blobsBundle: blobsBundle as BlobsBundle, + executionPayload, + executionRequests, + blobsBundle, executionPayloadValue, }; } diff --git a/packages/builder/test/unit/services/payloadSource.test.ts b/packages/builder/test/unit/services/payloadSource.test.ts index 2a13e85b09fa..4e91150c1318 100644 --- a/packages/builder/test/unit/services/payloadSource.test.ts +++ b/packages/builder/test/unit/services/payloadSource.test.ts @@ -1,12 +1,14 @@ -import {Mocked, beforeEach, describe, expect, it, vi} from "vitest"; -import {ForkName} from "@lodestar/params"; -import {ssz} from "@lodestar/types"; +import {type Mock, beforeEach, describe, expect, it, vi} from "vitest"; +import {ForkName, type ForkPostGloas} from "@lodestar/params"; +import {type ColumnIndex, type RootHex, ssz} from "@lodestar/types"; import {ErrorAborted, TimeoutError, toRootHex} from "@lodestar/utils"; import { BuildHandle, BuildRequest, EnginePayloadResult, EnginePayloadSource, + PayloadAttributes, + PayloadId, PayloadSourceEngine, PayloadSourceError, PayloadSourceErrorCode, @@ -30,23 +32,23 @@ describe("EnginePayloadSource", () => { }; const handle: BuildHandle = {sourceId, fork: ForkName.gloas, payloadId}; - let engine: Mocked; + let notifyForkchoiceUpdate: Mock; + let getPayload: Mock; let source: EnginePayloadSource; beforeEach(() => { - engine = { - notifyForkchoiceUpdate: vi.fn(), - getPayload: vi.fn(), - }; + notifyForkchoiceUpdate = vi.fn(); + getPayload = vi.fn(); + const engine = {notifyForkchoiceUpdate, getPayload} as unknown as PayloadSourceEngine; source = new EnginePayloadSource(sourceId, engine); }); it("prepares a payload and returns a source-bound handle", async () => { - engine.notifyForkchoiceUpdate.mockResolvedValue(payloadId); + notifyForkchoiceUpdate.mockResolvedValue(payloadId); const result = await source.prepare(request); - expect(engine.notifyForkchoiceUpdate).toHaveBeenCalledWith( + expect(notifyForkchoiceUpdate).toHaveBeenCalledWith( ForkName.gloas, forkchoiceState.headBlockHash, forkchoiceState.safeBlockHash, @@ -60,8 +62,8 @@ describe("EnginePayloadSource", () => { it("supports post-Gloas forks without narrowing the fork", async () => { const hezePayloadAttributes = ssz.heze.PayloadAttributes.defaultValue(); hezePayloadAttributes.inclusionListTransactions = [Uint8Array.from([1, 2, 3])]; - engine.notifyForkchoiceUpdate.mockResolvedValue(payloadId); - engine.getPayload.mockResolvedValue({ + notifyForkchoiceUpdate.mockResolvedValue(payloadId); + getPayload.mockResolvedValue({ executionPayload: ssz.heze.ExecutionPayload.defaultValue(), blobsBundle: ssz.heze.BlobsBundle.defaultValue(), executionRequests: ssz.heze.ExecutionRequests.defaultValue(), @@ -77,7 +79,7 @@ describe("EnginePayloadSource", () => { const builtPayload = await source.getPayload(result); expect(result.fork).toBe(ForkName.heze); - expect(engine.notifyForkchoiceUpdate).toHaveBeenCalledWith( + expect(notifyForkchoiceUpdate).toHaveBeenCalledWith( ForkName.heze, forkchoiceState.headBlockHash, forkchoiceState.safeBlockHash, @@ -85,13 +87,13 @@ describe("EnginePayloadSource", () => { hezePayloadAttributes, custodyColumns ); - expect(engine.getPayload).toHaveBeenCalledWith(ForkName.heze, payloadId); + expect(getPayload).toHaveBeenCalledWith(ForkName.heze, payloadId); expect(builtPayload.fork).toBe(ForkName.heze); expect(hezePayloadAttributes.inclusionListTransactions).toHaveLength(1); }); it("rejects a missing payload ID with a structured error", async () => { - engine.notifyForkchoiceUpdate.mockResolvedValue(null); + notifyForkchoiceUpdate.mockResolvedValue(null); const error = await getPayloadSourceError(source.prepare(request)); @@ -104,18 +106,18 @@ describe("EnginePayloadSource", () => { ["timeout", new TimeoutError("engine_forkchoiceUpdatedV4")], ["cancellation", new ErrorAborted("engine_forkchoiceUpdatedV4")], ])("propagates %s errors from payload preparation", async (_name, error) => { - engine.notifyForkchoiceUpdate.mockRejectedValue(error); + notifyForkchoiceUpdate.mockRejectedValue(error); await expect(source.prepare(request)).rejects.toBe(error); }); it("retrieves a complete payload without rebuilding exact-width values", async () => { const result = getEnginePayloadResult(); - engine.getPayload.mockResolvedValue(result); + getPayload.mockResolvedValue(result); const builtPayload = await source.getPayload(handle); - expect(engine.getPayload).toHaveBeenCalledWith(ForkName.gloas, payloadId); + expect(getPayload).toHaveBeenCalledWith(ForkName.gloas, payloadId); expect(builtPayload.sourceId).toBe(sourceId); expect(builtPayload.fork).toBe(ForkName.gloas); expect(builtPayload.executionPayload).toBe(result.executionPayload); @@ -132,11 +134,11 @@ describe("EnginePayloadSource", () => { sourceId, handleSourceId: "engine-1", }); - expect(engine.getPayload).not.toHaveBeenCalled(); + expect(getPayload).not.toHaveBeenCalled(); }); it("rejects a response without a blobs bundle", async () => { - engine.getPayload.mockResolvedValue({...getEnginePayloadResult(), blobsBundle: undefined}); + getPayload.mockResolvedValue({...getEnginePayloadResult(), blobsBundle: undefined}); const error = await getPayloadSourceError(source.getPayload(handle)); @@ -144,7 +146,7 @@ describe("EnginePayloadSource", () => { }); it("rejects a response without execution requests", async () => { - engine.getPayload.mockResolvedValue({...getEnginePayloadResult(), executionRequests: undefined}); + getPayload.mockResolvedValue({...getEnginePayloadResult(), executionRequests: undefined}); const error = await getPayloadSourceError(source.getPayload(handle)); @@ -153,7 +155,7 @@ describe("EnginePayloadSource", () => { it("propagates retrieval errors without replacing their type", async () => { const error = new TimeoutError("engine_getPayloadV6"); - engine.getPayload.mockRejectedValue(error); + getPayload.mockRejectedValue(error); await expect(source.getPayload(handle)).rejects.toBe(error); }); @@ -168,6 +170,17 @@ function getEnginePayloadResult(): EnginePayloadResult { }; } +type NotifyForkchoiceUpdate = ( + fork: ForkPostGloas, + headBlockHash: RootHex, + safeBlockHash: RootHex, + finalizedBlockHash: RootHex, + payloadAttributes: PayloadAttributes, + custodyColumns: ColumnIndex[] +) => Promise; + +type GetPayload = (fork: ForkPostGloas, payloadId: PayloadId) => Promise; + async function getPayloadSourceError(promise: Promise): Promise { try { await promise; From b45355d2e7cef066f3b7f45fd428e675372d4fd0 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Tue, 1 Sep 2026 23:52:44 +0100 Subject: [PATCH 06/11] refactor(builder): keep Engine ownership provisional --- packages/builder/src/services/payloadSource.ts | 9 +++++---- .../test/unit/services/payloadSource.test.ts | 17 ++++++++++++++++- 2 files changed, 21 insertions(+), 5 deletions(-) diff --git a/packages/builder/src/services/payloadSource.ts b/packages/builder/src/services/payloadSource.ts index 4b14fae4ed5e..2a84bfa2b438 100644 --- a/packages/builder/src/services/payloadSource.ts +++ b/packages/builder/src/services/payloadSource.ts @@ -23,7 +23,8 @@ export type BuildRequest = { fork: F; forkchoiceState: ForkchoiceState; payloadAttributes: PayloadAttributes; - custodyColumns: ColumnIndex[]; + /** Logical custody set. The transport serializes it for Engine API; null means no custody service. */ + custodyColumns: ColumnIndex[] | null; }; export type BuildHandle = { @@ -48,7 +49,7 @@ export type EnginePayloadResult = { executionRequests?: ExecutionRequests; }; -/** Narrow Engine API boundary whose transport owns request retries, timeouts, and Builder-lifetime cancellation. */ +/** Narrow Engine boundary whose transport owns serialization, retries, timeouts, and Builder-lifetime cancellation. */ export interface PayloadSourceEngine { notifyForkchoiceUpdate( fork: F, @@ -56,7 +57,7 @@ export interface PayloadSourceEngine { safeBlockHash: RootHex, finalizedBlockHash: RootHex, payloadAttributes: PayloadAttributes, - custodyColumns: ColumnIndex[] + custodyColumns: ColumnIndex[] | null ): Promise; getPayload(fork: F, payloadId: PayloadId): Promise>; } @@ -90,7 +91,7 @@ export type PayloadSourceErrorType = export class PayloadSourceError extends LodestarError {} -/** Payload source backed by an execution client's Engine API. */ +/** Payload source backed by an injected Engine boundary. Engine ownership and lifecycle remain caller policy. */ export class EnginePayloadSource implements PayloadSource { constructor( readonly id: string, diff --git a/packages/builder/test/unit/services/payloadSource.test.ts b/packages/builder/test/unit/services/payloadSource.test.ts index 4e91150c1318..563b28132235 100644 --- a/packages/builder/test/unit/services/payloadSource.test.ts +++ b/packages/builder/test/unit/services/payloadSource.test.ts @@ -59,6 +59,21 @@ describe("EnginePayloadSource", () => { expect(result).toEqual(handle); }); + it("preserves a null custody set", async () => { + notifyForkchoiceUpdate.mockResolvedValue(payloadId); + + await source.prepare({...request, custodyColumns: null}); + + expect(notifyForkchoiceUpdate).toHaveBeenCalledWith( + ForkName.gloas, + forkchoiceState.headBlockHash, + forkchoiceState.safeBlockHash, + forkchoiceState.finalizedBlockHash, + payloadAttributes, + null + ); + }); + it("supports post-Gloas forks without narrowing the fork", async () => { const hezePayloadAttributes = ssz.heze.PayloadAttributes.defaultValue(); hezePayloadAttributes.inclusionListTransactions = [Uint8Array.from([1, 2, 3])]; @@ -176,7 +191,7 @@ type NotifyForkchoiceUpdate = ( safeBlockHash: RootHex, finalizedBlockHash: RootHex, payloadAttributes: PayloadAttributes, - custodyColumns: ColumnIndex[] + custodyColumns: ColumnIndex[] | null ) => Promise; type GetPayload = (fork: ForkPostGloas, payloadId: PayloadId) => Promise; From 5679f5bdc35b2c2d8682f24e8d628d20d8d14650 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Wed, 2 Sep 2026 13:27:00 +0100 Subject: [PATCH 07/11] fix(builder): preserve payload request fork correlation --- .../builder/src/services/payloadSource.ts | 20 ++++++++++--------- .../test/unit/services/payloadSource.test.ts | 5 +++++ 2 files changed, 16 insertions(+), 9 deletions(-) diff --git a/packages/builder/src/services/payloadSource.ts b/packages/builder/src/services/payloadSource.ts index 2a84bfa2b438..53fda42bc228 100644 --- a/packages/builder/src/services/payloadSource.ts +++ b/packages/builder/src/services/payloadSource.ts @@ -19,13 +19,15 @@ export type ForkchoiceState = { export type PayloadAttributes = SSEPayloadAttributes["payloadAttributes"]; -export type BuildRequest = { - fork: F; - forkchoiceState: ForkchoiceState; - payloadAttributes: PayloadAttributes; - /** Logical custody set. The transport serializes it for Engine API; null means no custody service. */ - custodyColumns: ColumnIndex[] | null; -}; +export type BuildRequest = F extends ForkPostGloas + ? { + fork: F; + forkchoiceState: ForkchoiceState; + payloadAttributes: PayloadAttributes; + /** Logical custody set. The transport serializes it for Engine API; null means no custody service. */ + custodyColumns: ColumnIndex[] | null; + } + : never; export type BuildHandle = { sourceId: string; @@ -65,7 +67,7 @@ export interface PayloadSourceEngine { /** Source that prepares and retrieves complete execution payloads without owning build scheduling policy. */ export interface PayloadSource { readonly id: string; - prepare(request: BuildRequest): Promise>; + prepare(request: R): Promise>; getPayload(handle: BuildHandle): Promise>; } @@ -98,7 +100,7 @@ export class EnginePayloadSource implements PayloadSource { private readonly engine: PayloadSourceEngine ) {} - async prepare(request: BuildRequest): Promise> { + async prepare(request: R): Promise> { const {headBlockHash, safeBlockHash, finalizedBlockHash} = request.forkchoiceState; const payloadId = await this.engine.notifyForkchoiceUpdate( request.fork, diff --git a/packages/builder/test/unit/services/payloadSource.test.ts b/packages/builder/test/unit/services/payloadSource.test.ts index 563b28132235..2b4c650e08bb 100644 --- a/packages/builder/test/unit/services/payloadSource.test.ts +++ b/packages/builder/test/unit/services/payloadSource.test.ts @@ -30,6 +30,11 @@ describe("EnginePayloadSource", () => { payloadAttributes, custodyColumns, }; + + // @ts-expect-error Heze requests cannot use Gloas payload attributes. + const mismatchedRequest: BuildRequest = {...request, fork: ForkName.heze}; + void mismatchedRequest; + const handle: BuildHandle = {sourceId, fork: ForkName.gloas, payloadId}; let notifyForkchoiceUpdate: Mock; From b20eb6258fa3e8ca16e1b214c31d4c36df31e274 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Wed, 2 Sep 2026 14:40:16 +0100 Subject: [PATCH 08/11] fix(builder): propagate payload request cancellation --- .../builder/src/services/payloadSource.ts | 23 +++++--- .../test/unit/services/payloadSource.test.ts | 54 +++++++++++-------- 2 files changed, 46 insertions(+), 31 deletions(-) diff --git a/packages/builder/src/services/payloadSource.ts b/packages/builder/src/services/payloadSource.ts index 53fda42bc228..6f813f7b5fd9 100644 --- a/packages/builder/src/services/payloadSource.ts +++ b/packages/builder/src/services/payloadSource.ts @@ -59,16 +59,21 @@ export interface PayloadSourceEngine { safeBlockHash: RootHex, finalizedBlockHash: RootHex, payloadAttributes: PayloadAttributes, - custodyColumns: ColumnIndex[] | null + custodyColumns: ColumnIndex[] | null, + signal: AbortSignal ): Promise; - getPayload(fork: F, payloadId: PayloadId): Promise>; + getPayload( + fork: F, + payloadId: PayloadId, + signal: AbortSignal + ): Promise>; } /** Source that prepares and retrieves complete execution payloads without owning build scheduling policy. */ export interface PayloadSource { readonly id: string; - prepare(request: R): Promise>; - getPayload(handle: BuildHandle): Promise>; + prepare(request: R, signal: AbortSignal): Promise>; + getPayload(handle: BuildHandle, signal: AbortSignal): Promise>; } export enum PayloadSourceErrorCode { @@ -100,7 +105,7 @@ export class EnginePayloadSource implements PayloadSource { private readonly engine: PayloadSourceEngine ) {} - async prepare(request: R): Promise> { + async prepare(request: R, signal: AbortSignal): Promise> { const {headBlockHash, safeBlockHash, finalizedBlockHash} = request.forkchoiceState; const payloadId = await this.engine.notifyForkchoiceUpdate( request.fork, @@ -108,7 +113,8 @@ export class EnginePayloadSource implements PayloadSource { safeBlockHash, finalizedBlockHash, request.payloadAttributes, - request.custodyColumns + request.custodyColumns, + signal ); if (payloadId === null) { @@ -121,7 +127,7 @@ export class EnginePayloadSource implements PayloadSource { return {sourceId: this.id, fork: request.fork, payloadId}; } - async getPayload(handle: BuildHandle): Promise> { + async getPayload(handle: BuildHandle, signal: AbortSignal): Promise> { if (handle.sourceId !== this.id) { throw new PayloadSourceError( { @@ -135,7 +141,8 @@ export class EnginePayloadSource implements PayloadSource { const {executionPayload, executionPayloadValue, blobsBundle, executionRequests} = await this.engine.getPayload( handle.fork, - handle.payloadId + handle.payloadId, + signal ); if (blobsBundle === undefined) { diff --git a/packages/builder/test/unit/services/payloadSource.test.ts b/packages/builder/test/unit/services/payloadSource.test.ts index 2b4c650e08bb..dca24cc7b0da 100644 --- a/packages/builder/test/unit/services/payloadSource.test.ts +++ b/packages/builder/test/unit/services/payloadSource.test.ts @@ -36,6 +36,7 @@ describe("EnginePayloadSource", () => { void mismatchedRequest; const handle: BuildHandle = {sourceId, fork: ForkName.gloas, payloadId}; + const signal = new AbortController().signal; let notifyForkchoiceUpdate: Mock; let getPayload: Mock; @@ -51,7 +52,7 @@ describe("EnginePayloadSource", () => { it("prepares a payload and returns a source-bound handle", async () => { notifyForkchoiceUpdate.mockResolvedValue(payloadId); - const result = await source.prepare(request); + const result = await source.prepare(request, signal); expect(notifyForkchoiceUpdate).toHaveBeenCalledWith( ForkName.gloas, @@ -59,7 +60,8 @@ describe("EnginePayloadSource", () => { forkchoiceState.safeBlockHash, forkchoiceState.finalizedBlockHash, payloadAttributes, - custodyColumns + custodyColumns, + signal ); expect(result).toEqual(handle); }); @@ -67,7 +69,7 @@ describe("EnginePayloadSource", () => { it("preserves a null custody set", async () => { notifyForkchoiceUpdate.mockResolvedValue(payloadId); - await source.prepare({...request, custodyColumns: null}); + await source.prepare({...request, custodyColumns: null}, signal); expect(notifyForkchoiceUpdate).toHaveBeenCalledWith( ForkName.gloas, @@ -75,7 +77,8 @@ describe("EnginePayloadSource", () => { forkchoiceState.safeBlockHash, forkchoiceState.finalizedBlockHash, payloadAttributes, - null + null, + signal ); }); @@ -90,13 +93,16 @@ describe("EnginePayloadSource", () => { executionPayloadValue: 12_345_678_901_234_567_890n, }); - const result = await source.prepare({ - fork: ForkName.heze, - forkchoiceState, - payloadAttributes: hezePayloadAttributes, - custodyColumns, - }); - const builtPayload = await source.getPayload(result); + const result = await source.prepare( + { + fork: ForkName.heze, + forkchoiceState, + payloadAttributes: hezePayloadAttributes, + custodyColumns, + }, + signal + ); + const builtPayload = await source.getPayload(result, signal); expect(result.fork).toBe(ForkName.heze); expect(notifyForkchoiceUpdate).toHaveBeenCalledWith( @@ -105,9 +111,10 @@ describe("EnginePayloadSource", () => { forkchoiceState.safeBlockHash, forkchoiceState.finalizedBlockHash, hezePayloadAttributes, - custodyColumns + custodyColumns, + signal ); - expect(getPayload).toHaveBeenCalledWith(ForkName.heze, payloadId); + expect(getPayload).toHaveBeenCalledWith(ForkName.heze, payloadId, signal); expect(builtPayload.fork).toBe(ForkName.heze); expect(hezePayloadAttributes.inclusionListTransactions).toHaveLength(1); }); @@ -115,7 +122,7 @@ describe("EnginePayloadSource", () => { it("rejects a missing payload ID with a structured error", async () => { notifyForkchoiceUpdate.mockResolvedValue(null); - const error = await getPayloadSourceError(source.prepare(request)); + const error = await getPayloadSourceError(source.prepare(request, signal)); expect(error.type).toEqual({code: PayloadSourceErrorCode.NO_PAYLOAD_ID, sourceId}); }); @@ -128,16 +135,16 @@ describe("EnginePayloadSource", () => { ])("propagates %s errors from payload preparation", async (_name, error) => { notifyForkchoiceUpdate.mockRejectedValue(error); - await expect(source.prepare(request)).rejects.toBe(error); + await expect(source.prepare(request, signal)).rejects.toBe(error); }); it("retrieves a complete payload without rebuilding exact-width values", async () => { const result = getEnginePayloadResult(); getPayload.mockResolvedValue(result); - const builtPayload = await source.getPayload(handle); + const builtPayload = await source.getPayload(handle, signal); - expect(getPayload).toHaveBeenCalledWith(ForkName.gloas, payloadId); + expect(getPayload).toHaveBeenCalledWith(ForkName.gloas, payloadId, signal); expect(builtPayload.sourceId).toBe(sourceId); expect(builtPayload.fork).toBe(ForkName.gloas); expect(builtPayload.executionPayload).toBe(result.executionPayload); @@ -147,7 +154,7 @@ describe("EnginePayloadSource", () => { }); it("rejects a handle belonging to another source before calling the Engine API", async () => { - const error = await getPayloadSourceError(source.getPayload({...handle, sourceId: "engine-1"})); + const error = await getPayloadSourceError(source.getPayload({...handle, sourceId: "engine-1"}, signal)); expect(error.type).toEqual({ code: PayloadSourceErrorCode.SOURCE_MISMATCH, @@ -160,7 +167,7 @@ describe("EnginePayloadSource", () => { it("rejects a response without a blobs bundle", async () => { getPayload.mockResolvedValue({...getEnginePayloadResult(), blobsBundle: undefined}); - const error = await getPayloadSourceError(source.getPayload(handle)); + const error = await getPayloadSourceError(source.getPayload(handle, signal)); expect(error.type).toEqual({code: PayloadSourceErrorCode.MISSING_BLOBS_BUNDLE, sourceId, payloadId}); }); @@ -168,7 +175,7 @@ describe("EnginePayloadSource", () => { it("rejects a response without execution requests", async () => { getPayload.mockResolvedValue({...getEnginePayloadResult(), executionRequests: undefined}); - const error = await getPayloadSourceError(source.getPayload(handle)); + const error = await getPayloadSourceError(source.getPayload(handle, signal)); expect(error.type).toEqual({code: PayloadSourceErrorCode.MISSING_EXECUTION_REQUESTS, sourceId, payloadId}); }); @@ -177,7 +184,7 @@ describe("EnginePayloadSource", () => { const error = new TimeoutError("engine_getPayloadV6"); getPayload.mockRejectedValue(error); - await expect(source.getPayload(handle)).rejects.toBe(error); + await expect(source.getPayload(handle, signal)).rejects.toBe(error); }); }); @@ -196,10 +203,11 @@ type NotifyForkchoiceUpdate = ( safeBlockHash: RootHex, finalizedBlockHash: RootHex, payloadAttributes: PayloadAttributes, - custodyColumns: ColumnIndex[] | null + custodyColumns: ColumnIndex[] | null, + signal: AbortSignal ) => Promise; -type GetPayload = (fork: ForkPostGloas, payloadId: PayloadId) => Promise; +type GetPayload = (fork: ForkPostGloas, payloadId: PayloadId, signal: AbortSignal) => Promise; async function getPayloadSourceError(promise: Promise): Promise { try { From f61017874acda98e667fe54613dac84078e5cf74 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Wed, 2 Sep 2026 14:53:22 +0100 Subject: [PATCH 09/11] docs(builder): clarify payload transport ownership --- packages/builder/src/services/payloadSource.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/builder/src/services/payloadSource.ts b/packages/builder/src/services/payloadSource.ts index 6f813f7b5fd9..1de219e1be0b 100644 --- a/packages/builder/src/services/payloadSource.ts +++ b/packages/builder/src/services/payloadSource.ts @@ -51,7 +51,7 @@ export type EnginePayloadResult = { executionRequests?: ExecutionRequests; }; -/** Narrow Engine boundary whose transport owns serialization, retries, timeouts, and Builder-lifetime cancellation. */ +/** Narrow Engine boundary whose transport owns serialization, retries, and request execution. */ export interface PayloadSourceEngine { notifyForkchoiceUpdate( fork: F, From 950334308e410e63f35e7a5a358cc6bd70298819 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Wed, 2 Sep 2026 18:07:00 +0100 Subject: [PATCH 10/11] feat(builder): assemble stateless payload envelopes --- .../src/services/executionPayloadEnvelope.ts | 99 +++++++++++++++ .../services/executionPayloadEnvelope.test.ts | 115 ++++++++++++++++++ 2 files changed, 214 insertions(+) create mode 100644 packages/builder/src/services/executionPayloadEnvelope.ts create mode 100644 packages/builder/test/unit/services/executionPayloadEnvelope.test.ts diff --git a/packages/builder/src/services/executionPayloadEnvelope.ts b/packages/builder/src/services/executionPayloadEnvelope.ts new file mode 100644 index 000000000000..cdd564f0a968 --- /dev/null +++ b/packages/builder/src/services/executionPayloadEnvelope.ts @@ -0,0 +1,99 @@ +import type {BuilderIndex, RootHex, Slot, gloas} from "@lodestar/types"; +import {LodestarError, fromHex, toRootHex} from "@lodestar/utils"; +import type {BuiltPayload} from "./payloadSource.js"; + +export type SelectedBidIdentity = { + slot: Slot; + parentBlockHash: RootHex; + parentBlockRoot: RootHex; + blockHash: RootHex; +}; + +export type ExecutionPayloadEnvelopeInput = { + blockRoot: RootHex; + builderIndex: BuilderIndex; + selectedBid: SelectedBidIdentity; + payload: BuiltPayload; +}; + +export type ExecutionPayloadEnvelopeMaterial = { + envelope: gloas.ExecutionPayloadEnvelope; + kzgProofs: BuiltPayload["blobsBundle"]["proofs"]; + blobs: BuiltPayload["blobsBundle"]["blobs"]; +}; + +export enum ExecutionPayloadEnvelopeErrorCode { + SLOT_MISMATCH = "EXECUTION_PAYLOAD_ENVELOPE_ERROR_SLOT_MISMATCH", + PARENT_BLOCK_HASH_MISMATCH = "EXECUTION_PAYLOAD_ENVELOPE_ERROR_PARENT_BLOCK_HASH_MISMATCH", + BLOCK_HASH_MISMATCH = "EXECUTION_PAYLOAD_ENVELOPE_ERROR_BLOCK_HASH_MISMATCH", +} + +export type ExecutionPayloadEnvelopeErrorType = + | { + code: ExecutionPayloadEnvelopeErrorCode.SLOT_MISMATCH; + bidSlot: Slot; + payloadSlot: Slot; + } + | { + code: ExecutionPayloadEnvelopeErrorCode.PARENT_BLOCK_HASH_MISMATCH; + bidParentBlockHash: RootHex; + payloadParentBlockHash: RootHex; + } + | { + code: ExecutionPayloadEnvelopeErrorCode.BLOCK_HASH_MISMATCH; + bidBlockHash: RootHex; + payloadBlockHash: RootHex; + }; + +export class ExecutionPayloadEnvelopeError extends LodestarError {} + +export function createExecutionPayloadEnvelopeMaterial({ + blockRoot, + builderIndex, + selectedBid, + payload, +}: ExecutionPayloadEnvelopeInput): ExecutionPayloadEnvelopeMaterial { + const payloadSlot = payload.executionPayload.slotNumber; + if (payloadSlot !== selectedBid.slot) { + throw new ExecutionPayloadEnvelopeError( + {code: ExecutionPayloadEnvelopeErrorCode.SLOT_MISMATCH, bidSlot: selectedBid.slot, payloadSlot}, + `Selected bid slot does not match payload slot bidSlot=${selectedBid.slot} payloadSlot=${payloadSlot}` + ); + } + + const payloadParentBlockHash = toRootHex(payload.executionPayload.parentHash); + if (payloadParentBlockHash !== selectedBid.parentBlockHash) { + throw new ExecutionPayloadEnvelopeError( + { + code: ExecutionPayloadEnvelopeErrorCode.PARENT_BLOCK_HASH_MISMATCH, + bidParentBlockHash: selectedBid.parentBlockHash, + payloadParentBlockHash, + }, + `Selected bid parent does not match payload parent bidParentBlockHash=${selectedBid.parentBlockHash} payloadParentBlockHash=${payloadParentBlockHash}` + ); + } + + const payloadBlockHash = toRootHex(payload.executionPayload.blockHash); + if (payloadBlockHash !== selectedBid.blockHash) { + throw new ExecutionPayloadEnvelopeError( + { + code: ExecutionPayloadEnvelopeErrorCode.BLOCK_HASH_MISMATCH, + bidBlockHash: selectedBid.blockHash, + payloadBlockHash, + }, + `Selected bid block hash does not match payload bidBlockHash=${selectedBid.blockHash} payloadBlockHash=${payloadBlockHash}` + ); + } + + return { + envelope: { + payload: payload.executionPayload, + executionRequests: payload.executionRequests, + builderIndex, + beaconBlockRoot: fromHex(blockRoot), + parentBeaconBlockRoot: fromHex(selectedBid.parentBlockRoot), + }, + kzgProofs: payload.blobsBundle.proofs, + blobs: payload.blobsBundle.blobs, + }; +} diff --git a/packages/builder/test/unit/services/executionPayloadEnvelope.test.ts b/packages/builder/test/unit/services/executionPayloadEnvelope.test.ts new file mode 100644 index 000000000000..e1fa0f3f623f --- /dev/null +++ b/packages/builder/test/unit/services/executionPayloadEnvelope.test.ts @@ -0,0 +1,115 @@ +import {describe, expect, it} from "vitest"; +import {ForkName, type ForkPostGloas} from "@lodestar/params"; +import type {RootHex} from "@lodestar/types"; +import {ssz} from "@lodestar/types"; +import {fromHex, toRootHex} from "@lodestar/utils"; +import { + ExecutionPayloadEnvelopeError, + ExecutionPayloadEnvelopeErrorCode, + type SelectedBidIdentity, + createExecutionPayloadEnvelopeMaterial, +} from "../../../src/services/executionPayloadEnvelope.js"; +import type {BuiltPayload} from "../../../src/services/payloadSource.js"; + +const builderIndex = 7; +const blockRoot = root(8); + +describe("createExecutionPayloadEnvelopeMaterial", () => { + for (const fork of [ForkName.gloas, ForkName.heze] as const) { + it(`assembles exact ${fork} stateless envelope material`, () => { + const payload = createBuiltPayload(fork); + const selectedBid = bidIdentity(payload); + + const material = createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, payload}); + + expect(material.envelope).toEqual({ + payload: payload.executionPayload, + executionRequests: payload.executionRequests, + builderIndex, + beaconBlockRoot: fromHex(blockRoot), + parentBeaconBlockRoot: fromHex(selectedBid.parentBlockRoot), + }); + expect(material.kzgProofs).toBe(payload.blobsBundle.proofs); + expect(material.blobs).toBe(payload.blobsBundle.blobs); + }); + } + + it("rejects retained material for a different slot", () => { + const payload = createBuiltPayload(ForkName.gloas); + const selectedBid = {...bidIdentity(payload), slot: 11}; + + expectEnvelopeError(() => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, payload}), { + code: ExecutionPayloadEnvelopeErrorCode.SLOT_MISMATCH, + bidSlot: 11, + payloadSlot: payload.executionPayload.slotNumber, + }); + }); + + it("rejects retained material for a different parent block hash", () => { + const payload = createBuiltPayload(ForkName.gloas); + const selectedBid = {...bidIdentity(payload), parentBlockHash: root(9)}; + + expectEnvelopeError(() => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, payload}), { + code: ExecutionPayloadEnvelopeErrorCode.PARENT_BLOCK_HASH_MISMATCH, + bidParentBlockHash: selectedBid.parentBlockHash, + payloadParentBlockHash: toRootHex(payload.executionPayload.parentHash), + }); + }); + + it("rejects retained material for a different execution block hash", () => { + const payload = createBuiltPayload(ForkName.gloas); + const selectedBid = {...bidIdentity(payload), blockHash: root(9)}; + + expectEnvelopeError(() => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, payload}), { + code: ExecutionPayloadEnvelopeErrorCode.BLOCK_HASH_MISMATCH, + bidBlockHash: selectedBid.blockHash, + payloadBlockHash: toRootHex(payload.executionPayload.blockHash), + }); + }); +}); + +function createBuiltPayload(fork: F): BuiltPayload { + const forkTypes = fork === ForkName.heze ? ssz.heze : ssz.gloas; + const executionPayload = forkTypes.ExecutionPayload.defaultValue(); + executionPayload.slotNumber = 10; + executionPayload.parentHash = Buffer.alloc(32, 2); + executionPayload.blockHash = Buffer.alloc(32, 4); + const blobsBundle = forkTypes.BlobsBundle.defaultValue(); + blobsBundle.proofs.push(Buffer.alloc(48, 5)); + blobsBundle.blobs.push(Buffer.alloc(0)); + + return { + sourceId: "engine", + fork, + executionPayload, + executionRequests: forkTypes.ExecutionRequests.defaultValue(), + blobsBundle, + executionPayloadValue: 1n, + } as BuiltPayload; +} + +function bidIdentity(payload: BuiltPayload): SelectedBidIdentity { + return { + slot: payload.executionPayload.slotNumber, + parentBlockHash: toRootHex(payload.executionPayload.parentHash), + parentBlockRoot: root(3), + blockHash: toRootHex(payload.executionPayload.blockHash), + }; +} + +function root(byte: number): RootHex { + return toRootHex(Buffer.alloc(32, byte)); +} + +function expectEnvelopeError(fn: () => unknown, type: ExecutionPayloadEnvelopeError["type"]): void { + expect(fn).toThrowError(ExecutionPayloadEnvelopeError); + try { + fn(); + throw Error("Expected ExecutionPayloadEnvelopeError"); + } catch (error) { + if (!(error instanceof ExecutionPayloadEnvelopeError)) { + throw error; + } + expect(error.type).toEqual(type); + } +} From 1e99557c0ba48553c5c9f23d07840b40626a2cc5 Mon Sep 17 00:00:00 2001 From: Kris O'Shea Date: Thu, 3 Sep 2026 15:34:00 +0100 Subject: [PATCH 11/11] fix(builder): bind envelopes to retained parent root --- .../src/services/executionPayloadEnvelope.ts | 28 +++++++- .../services/executionPayloadEnvelope.test.ts | 65 ++++++++++++++----- 2 files changed, 74 insertions(+), 19 deletions(-) diff --git a/packages/builder/src/services/executionPayloadEnvelope.ts b/packages/builder/src/services/executionPayloadEnvelope.ts index cdd564f0a968..8b74c35a3854 100644 --- a/packages/builder/src/services/executionPayloadEnvelope.ts +++ b/packages/builder/src/services/executionPayloadEnvelope.ts @@ -1,4 +1,4 @@ -import type {BuilderIndex, RootHex, Slot, gloas} from "@lodestar/types"; +import type {BuilderIndex, Root, RootHex, Slot, gloas} from "@lodestar/types"; import {LodestarError, fromHex, toRootHex} from "@lodestar/utils"; import type {BuiltPayload} from "./payloadSource.js"; @@ -13,7 +13,10 @@ export type ExecutionPayloadEnvelopeInput = { blockRoot: RootHex; builderIndex: BuilderIndex; selectedBid: SelectedBidIdentity; - payload: BuiltPayload; + storedPayload: { + parentBlockRoot: Root; + payload: BuiltPayload; + }; }; export type ExecutionPayloadEnvelopeMaterial = { @@ -24,11 +27,17 @@ export type ExecutionPayloadEnvelopeMaterial = { export enum ExecutionPayloadEnvelopeErrorCode { SLOT_MISMATCH = "EXECUTION_PAYLOAD_ENVELOPE_ERROR_SLOT_MISMATCH", + PARENT_BLOCK_ROOT_MISMATCH = "EXECUTION_PAYLOAD_ENVELOPE_ERROR_PARENT_BLOCK_ROOT_MISMATCH", PARENT_BLOCK_HASH_MISMATCH = "EXECUTION_PAYLOAD_ENVELOPE_ERROR_PARENT_BLOCK_HASH_MISMATCH", BLOCK_HASH_MISMATCH = "EXECUTION_PAYLOAD_ENVELOPE_ERROR_BLOCK_HASH_MISMATCH", } export type ExecutionPayloadEnvelopeErrorType = + | { + code: ExecutionPayloadEnvelopeErrorCode.PARENT_BLOCK_ROOT_MISMATCH; + bidParentBlockRoot: RootHex; + storedParentBlockRoot: RootHex; + } | { code: ExecutionPayloadEnvelopeErrorCode.SLOT_MISMATCH; bidSlot: Slot; @@ -51,8 +60,21 @@ export function createExecutionPayloadEnvelopeMaterial({ blockRoot, builderIndex, selectedBid, - payload, + storedPayload, }: ExecutionPayloadEnvelopeInput): ExecutionPayloadEnvelopeMaterial { + const storedParentBlockRoot = toRootHex(storedPayload.parentBlockRoot); + if (storedParentBlockRoot !== selectedBid.parentBlockRoot) { + throw new ExecutionPayloadEnvelopeError( + { + code: ExecutionPayloadEnvelopeErrorCode.PARENT_BLOCK_ROOT_MISMATCH, + bidParentBlockRoot: selectedBid.parentBlockRoot, + storedParentBlockRoot, + }, + `Selected bid beacon parent does not match retained payload bidParentBlockRoot=${selectedBid.parentBlockRoot} storedParentBlockRoot=${storedParentBlockRoot}` + ); + } + + const {payload} = storedPayload; const payloadSlot = payload.executionPayload.slotNumber; if (payloadSlot !== selectedBid.slot) { throw new ExecutionPayloadEnvelopeError( diff --git a/packages/builder/test/unit/services/executionPayloadEnvelope.test.ts b/packages/builder/test/unit/services/executionPayloadEnvelope.test.ts index e1fa0f3f623f..d9826e2b00a9 100644 --- a/packages/builder/test/unit/services/executionPayloadEnvelope.test.ts +++ b/packages/builder/test/unit/services/executionPayloadEnvelope.test.ts @@ -6,6 +6,7 @@ import {fromHex, toRootHex} from "@lodestar/utils"; import { ExecutionPayloadEnvelopeError, ExecutionPayloadEnvelopeErrorCode, + type ExecutionPayloadEnvelopeInput, type SelectedBidIdentity, createExecutionPayloadEnvelopeMaterial, } from "../../../src/services/executionPayloadEnvelope.js"; @@ -19,8 +20,9 @@ describe("createExecutionPayloadEnvelopeMaterial", () => { it(`assembles exact ${fork} stateless envelope material`, () => { const payload = createBuiltPayload(fork); const selectedBid = bidIdentity(payload); + const storedPayload = retain(payload, selectedBid.parentBlockRoot); - const material = createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, payload}); + const material = createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, storedPayload}); expect(material.envelope).toEqual({ payload: payload.executionPayload, @@ -37,34 +39,61 @@ describe("createExecutionPayloadEnvelopeMaterial", () => { it("rejects retained material for a different slot", () => { const payload = createBuiltPayload(ForkName.gloas); const selectedBid = {...bidIdentity(payload), slot: 11}; + const storedPayload = retain(payload, selectedBid.parentBlockRoot); - expectEnvelopeError(() => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, payload}), { - code: ExecutionPayloadEnvelopeErrorCode.SLOT_MISMATCH, - bidSlot: 11, - payloadSlot: payload.executionPayload.slotNumber, - }); + expectEnvelopeError( + () => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, storedPayload}), + { + code: ExecutionPayloadEnvelopeErrorCode.SLOT_MISMATCH, + bidSlot: 11, + payloadSlot: payload.executionPayload.slotNumber, + } + ); + }); + + it("rejects retained material for a different parent block root", () => { + const payload = createBuiltPayload(ForkName.gloas); + const selectedBid = bidIdentity(payload); + const storedPayload = retain(payload, root(9)); + + expectEnvelopeError( + () => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, storedPayload}), + { + code: ExecutionPayloadEnvelopeErrorCode.PARENT_BLOCK_ROOT_MISMATCH, + bidParentBlockRoot: selectedBid.parentBlockRoot, + storedParentBlockRoot: root(9), + } + ); }); it("rejects retained material for a different parent block hash", () => { const payload = createBuiltPayload(ForkName.gloas); const selectedBid = {...bidIdentity(payload), parentBlockHash: root(9)}; + const storedPayload = retain(payload, selectedBid.parentBlockRoot); - expectEnvelopeError(() => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, payload}), { - code: ExecutionPayloadEnvelopeErrorCode.PARENT_BLOCK_HASH_MISMATCH, - bidParentBlockHash: selectedBid.parentBlockHash, - payloadParentBlockHash: toRootHex(payload.executionPayload.parentHash), - }); + expectEnvelopeError( + () => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, storedPayload}), + { + code: ExecutionPayloadEnvelopeErrorCode.PARENT_BLOCK_HASH_MISMATCH, + bidParentBlockHash: selectedBid.parentBlockHash, + payloadParentBlockHash: toRootHex(payload.executionPayload.parentHash), + } + ); }); it("rejects retained material for a different execution block hash", () => { const payload = createBuiltPayload(ForkName.gloas); const selectedBid = {...bidIdentity(payload), blockHash: root(9)}; + const storedPayload = retain(payload, selectedBid.parentBlockRoot); - expectEnvelopeError(() => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, payload}), { - code: ExecutionPayloadEnvelopeErrorCode.BLOCK_HASH_MISMATCH, - bidBlockHash: selectedBid.blockHash, - payloadBlockHash: toRootHex(payload.executionPayload.blockHash), - }); + expectEnvelopeError( + () => createExecutionPayloadEnvelopeMaterial({blockRoot, builderIndex, selectedBid, storedPayload}), + { + code: ExecutionPayloadEnvelopeErrorCode.BLOCK_HASH_MISMATCH, + bidBlockHash: selectedBid.blockHash, + payloadBlockHash: toRootHex(payload.executionPayload.blockHash), + } + ); }); }); @@ -97,6 +126,10 @@ function bidIdentity(payload: BuiltPayload): SelectedBidIdentity { }; } +function retain(payload: BuiltPayload, parentBlockRoot: RootHex): ExecutionPayloadEnvelopeInput["storedPayload"] { + return {parentBlockRoot: fromHex(parentBlockRoot), payload}; +} + function root(byte: number): RootHex { return toRootHex(Buffer.alloc(32, byte)); }