diff --git a/packages/beacon-node/src/chain/blocks/importBlock.ts b/packages/beacon-node/src/chain/blocks/importBlock.ts index 4657361240e7..7bdf387f9906 100644 --- a/packages/beacon-node/src/chain/blocks/importBlock.ts +++ b/packages/beacon-node/src/chain/blocks/importBlock.ts @@ -51,7 +51,7 @@ export async function importBlock( opts: ImportBlockOpts ): Promise { const {blockInput, postState, parentBlockSlot, executionStatus} = fullyVerifiedBlock; - const {block} = blockInput; + const {block, serializedData} = blockInput; const blockRoot = this.config.getForkTypes(block.message.slot).BeaconBlock.hashTreeRoot(block.message); const blockRootHex = toHexString(blockRoot); const currentEpoch = computeEpochAtSlot(this.forkChoice.getTime()); @@ -60,8 +60,14 @@ export async function importBlock( const blockDelaySec = (fullyVerifiedBlock.seenTimestampSec - postState.genesisTime) % this.config.SECONDS_PER_SLOT; // 1. Persist block to hot DB (pre-emptively) - - await this.db.block.add(block); + if (serializedData) { + // skip serializing data if we already have it + this.metrics?.importBlock.persistBlockWithSerializedDataCount.inc(); + await this.db.block.putBinary(this.db.block.getId(block), serializedData); + } else { + this.metrics?.importBlock.persistBlockNoSerializedDataCount.inc(); + await this.db.block.add(block); + } this.logger.debug("Persisted block to hot DB", { slot: block.message.slot, root: blockRootHex, @@ -245,7 +251,7 @@ export async function importBlock( // Only track "recent" blocks. Otherwise sync can distort this metrics heavily. // We want to track recent blocks coming from gossip, unknown block sync, and API. if (delaySec < 64 * this.config.SECONDS_PER_SLOT) { - this.metrics.elapsedTimeTillBecomeHead.observe(delaySec); + this.metrics.importBlock.elapsedTimeTillBecomeHead.observe(delaySec); } } diff --git a/packages/beacon-node/src/chain/blocks/index.ts b/packages/beacon-node/src/chain/blocks/index.ts index 04096b02c402..ab8d54e983f6 100644 --- a/packages/beacon-node/src/chain/blocks/index.ts +++ b/packages/beacon-node/src/chain/blocks/index.ts @@ -1,4 +1,4 @@ -import {allForks} from "@lodestar/types"; +import {WithOptionalBytes, allForks} from "@lodestar/types"; import {toHex} from "@lodestar/utils"; import {JobItemQueue} from "../../util/queue/index.js"; import {Metrics} from "../../metrics/metrics.js"; @@ -18,10 +18,10 @@ const QUEUE_MAX_LENGTH = 256; * BlockProcessor processes block jobs in a queued fashion, one after the other. */ export class BlockProcessor { - readonly jobQueue: JobItemQueue<[BlockInput[], ImportBlockOpts], void>; + readonly jobQueue: JobItemQueue<[WithOptionalBytes[], ImportBlockOpts], void>; constructor(chain: BeaconChain, metrics: Metrics | null, opts: BlockProcessOpts, signal: AbortSignal) { - this.jobQueue = new JobItemQueue<[BlockInput[], ImportBlockOpts], void>( + this.jobQueue = new JobItemQueue<[WithOptionalBytes[], ImportBlockOpts], void>( (job, importOpts) => { return processBlocks.call(chain, job, {...opts, ...importOpts}); }, @@ -30,7 +30,7 @@ export class BlockProcessor { ); } - async processBlocksJob(job: BlockInput[], opts: ImportBlockOpts = {}): Promise { + async processBlocksJob(job: WithOptionalBytes[], opts: ImportBlockOpts = {}): Promise { await this.jobQueue.push(job, opts); } } @@ -47,7 +47,7 @@ export class BlockProcessor { */ export async function processBlocks( this: BeaconChain, - blocks: BlockInput[], + blocks: WithOptionalBytes[], opts: BlockProcessOpts & ImportBlockOpts ): Promise { if (blocks.length === 0) { diff --git a/packages/beacon-node/src/chain/blocks/types.ts b/packages/beacon-node/src/chain/blocks/types.ts index 32d341681164..8fa01bbd131a 100644 --- a/packages/beacon-node/src/chain/blocks/types.ts +++ b/packages/beacon-node/src/chain/blocks/types.ts @@ -1,6 +1,6 @@ import {CachedBeaconStateAllForks, computeEpochAtSlot} from "@lodestar/state-transition"; import {MaybeValidExecutionStatus} from "@lodestar/fork-choice"; -import {allForks, deneb, Slot} from "@lodestar/types"; +import {allForks, deneb, Slot, WithOptionalBytes} from "@lodestar/types"; import {ForkSeq, MIN_EPOCHS_FOR_BLOB_SIDECARS_REQUESTS} from "@lodestar/params"; import {ChainForkConfig} from "@lodestar/config"; @@ -92,7 +92,7 @@ export type ImportBlockOpts = { * A wrapper around a `SignedBeaconBlock` that indicates that this block is fully verified and ready to import */ export type FullyVerifiedBlock = { - blockInput: BlockInput; + blockInput: WithOptionalBytes; postState: CachedBeaconStateAllForks; parentBlockSlot: Slot; proposerBalanceDelta: number; diff --git a/packages/beacon-node/src/chain/blocks/verifyBlocksSanityChecks.ts b/packages/beacon-node/src/chain/blocks/verifyBlocksSanityChecks.ts index 03afc8fa138c..6df2c263d12f 100644 --- a/packages/beacon-node/src/chain/blocks/verifyBlocksSanityChecks.ts +++ b/packages/beacon-node/src/chain/blocks/verifyBlocksSanityChecks.ts @@ -1,7 +1,7 @@ import {computeStartSlotAtEpoch} from "@lodestar/state-transition"; import {ChainForkConfig} from "@lodestar/config"; import {IForkChoice, ProtoBlock} from "@lodestar/fork-choice"; -import {Slot} from "@lodestar/types"; +import {Slot, WithOptionalBytes} from "@lodestar/types"; import {toHexString} from "@lodestar/utils"; import {IClock} from "../../util/clock.js"; import {BlockError, BlockErrorCode} from "../errors/index.js"; @@ -21,14 +21,14 @@ import {BlockInput, ImportBlockOpts} from "./types.js"; */ export function verifyBlocksSanityChecks( chain: {forkChoice: IForkChoice; clock: IClock; config: ChainForkConfig}, - blocks: BlockInput[], + blocks: WithOptionalBytes[], opts: ImportBlockOpts -): {relevantBlocks: BlockInput[]; parentSlots: Slot[]; parentBlock: ProtoBlock | null} { +): {relevantBlocks: WithOptionalBytes[]; parentSlots: Slot[]; parentBlock: ProtoBlock | null} { if (blocks.length === 0) { throw Error("Empty partiallyVerifiedBlocks"); } - const relevantBlocks: BlockInput[] = []; + const relevantBlocks: WithOptionalBytes[] = []; const parentSlots: Slot[] = []; let parentBlock: ProtoBlock | null = null; diff --git a/packages/beacon-node/src/chain/chain.ts b/packages/beacon-node/src/chain/chain.ts index 9a01cc6d0e50..5adb62b55375 100644 --- a/packages/beacon-node/src/chain/chain.ts +++ b/packages/beacon-node/src/chain/chain.ts @@ -12,7 +12,19 @@ import { PubkeyIndexMap, } from "@lodestar/state-transition"; import {BeaconConfig} from "@lodestar/config"; -import {allForks, UintNum64, Root, phase0, Slot, RootHex, Epoch, ValidatorIndex, deneb, Wei} from "@lodestar/types"; +import { + allForks, + UintNum64, + Root, + phase0, + Slot, + RootHex, + Epoch, + ValidatorIndex, + deneb, + Wei, + WithOptionalBytes, +} from "@lodestar/types"; import {CheckpointWithHex, ExecutionStatus, IForkChoice, ProtoBlock} from "@lodestar/fork-choice"; import {ProcessShutdownCallback} from "@lodestar/validator"; import {Logger, pruneSetToMax, toHex} from "@lodestar/utils"; @@ -450,11 +462,11 @@ export class BeaconChain implements IBeaconChain { return blobsSidecar; } - async processBlock(block: BlockInput, opts?: ImportBlockOpts): Promise { + async processBlock(block: WithOptionalBytes, opts?: ImportBlockOpts): Promise { return this.blockProcessor.processBlocksJob([block], opts); } - async processChainSegment(blocks: BlockInput[], opts?: ImportBlockOpts): Promise { + async processChainSegment(blocks: WithOptionalBytes[], opts?: ImportBlockOpts): Promise { return this.blockProcessor.processBlocksJob(blocks, opts); } diff --git a/packages/beacon-node/src/chain/interface.ts b/packages/beacon-node/src/chain/interface.ts index 70111ec4123b..6cf2ea8c2d2e 100644 --- a/packages/beacon-node/src/chain/interface.ts +++ b/packages/beacon-node/src/chain/interface.ts @@ -1,4 +1,16 @@ -import {allForks, UintNum64, Root, phase0, Slot, RootHex, Epoch, ValidatorIndex, deneb, Wei} from "@lodestar/types"; +import { + allForks, + UintNum64, + Root, + phase0, + Slot, + RootHex, + Epoch, + ValidatorIndex, + deneb, + Wei, + WithOptionalBytes, +} from "@lodestar/types"; import {CachedBeaconStateAllForks, Index2PubkeyCache, PubkeyIndexMap} from "@lodestar/state-transition"; import {BeaconConfig} from "@lodestar/config"; import {CompositeTypeAny, TreeView, Type} from "@chainsafe/ssz"; @@ -122,9 +134,9 @@ export interface IBeaconChain { produceBlindedBlock(blockAttributes: BlockAttributes): Promise<{block: allForks.BlindedBeaconBlock; blockValue: Wei}>; /** Process a block until complete */ - processBlock(block: BlockInput, opts?: ImportBlockOpts): Promise; + processBlock(block: WithOptionalBytes, opts?: ImportBlockOpts): Promise; /** Process a chain of blocks until complete */ - processChainSegment(blocks: BlockInput[], opts?: ImportBlockOpts): Promise; + processChainSegment(blocks: WithOptionalBytes[], opts?: ImportBlockOpts): Promise; getStatus(): phase0.Status; diff --git a/packages/beacon-node/src/metrics/metrics/lodestar.ts b/packages/beacon-node/src/metrics/metrics/lodestar.ts index 6aca4ea6c6ed..d82953d81b27 100644 --- a/packages/beacon-node/src/metrics/metrics/lodestar.ts +++ b/packages/beacon-node/src/metrics/metrics/lodestar.ts @@ -496,11 +496,21 @@ export function createLodestarMetrics( buckets: [0.05, 0.1, 0.2, 0.5, 1, 1.5, 2, 4], }), }, - elapsedTimeTillBecomeHead: register.histogram({ - name: "lodestar_gossip_block_elapsed_time_till_become_head", - help: "Time elapsed between block slot time and the time block becomes head", - buckets: [0.5, 1, 2, 4, 6, 12], - }), + importBlock: { + persistBlockNoSerializedDataCount: register.gauge({ + name: "lodestar_import_block_persist_block_no_serialized_data_count", + help: "Count persisting block with no serialized data", + }), + persistBlockWithSerializedDataCount: register.gauge({ + name: "lodestar_import_block_persist_block_with_serialized_data_count", + help: "Count persisting block with serialized data", + }), + elapsedTimeTillBecomeHead: register.histogram({ + name: "lodestar_gossip_block_elapsed_time_till_become_head", + help: "Time elapsed between block slot time and the time block becomes head", + buckets: [0.5, 1, 2, 4, 6, 12], + }), + }, engineNotifyNewPayloadResult: register.gauge<"result">({ name: "lodestar_execution_engine_notify_new_payload_result_total", help: "The total result of calling notifyNewPayload execution engine api", diff --git a/packages/beacon-node/src/network/processor/gossipHandlers.ts b/packages/beacon-node/src/network/processor/gossipHandlers.ts index 998401b6a106..a9c672df6efd 100644 --- a/packages/beacon-node/src/network/processor/gossipHandlers.ts +++ b/packages/beacon-node/src/network/processor/gossipHandlers.ts @@ -2,7 +2,7 @@ import {peerIdFromString} from "@libp2p/peer-id"; import {toHexString} from "@chainsafe/ssz"; import {BeaconConfig} from "@lodestar/config"; import {Logger, prettyBytes} from "@lodestar/utils"; -import {Root, Slot, ssz} from "@lodestar/types"; +import {Root, Slot, ssz, WithBytes} from "@lodestar/types"; import {ForkName, ForkSeq} from "@lodestar/params"; import {Metrics} from "../../metrics/index.js"; import {OpSource} from "../../metrics/validatorMonitor.js"; @@ -124,7 +124,11 @@ export function getGossipHandlers(modules: ValidatorFnsModules, options: GossipH } } - function handleValidBeaconBlock(blockInput: BlockInput, peerIdStr: string, seenTimestampSec: number): void { + function handleValidBeaconBlock( + blockInput: WithBytes, + peerIdStr: string, + seenTimestampSec: number + ): void { const signedBlock = blockInput.block; // Handler - MUST NOT `await`, to allow validation result to be propagated @@ -178,7 +182,7 @@ export function getGossipHandlers(modules: ValidatorFnsModules, options: GossipH const blockInput = getBlockInput.preDeneb(config, signedBlock); await validateBeaconBlock(blockInput, topic.fork, peerIdStr, seenTimestampSec); - handleValidBeaconBlock(blockInput, peerIdStr, seenTimestampSec); + handleValidBeaconBlock({...blockInput, serializedData}, peerIdStr, seenTimestampSec); }, [GossipType.beacon_block_and_blobs_sidecar]: async ({serializedData}, topic, peerIdStr, seenTimestampSec) => { @@ -193,7 +197,7 @@ export function getGossipHandlers(modules: ValidatorFnsModules, options: GossipH const blockInput = getBlockInput.postDeneb(config, beaconBlock, blobsSidecar); await validateBeaconBlock(blockInput, topic.fork, peerIdStr, seenTimestampSec); validateGossipBlobsSidecar(beaconBlock, blobsSidecar); - handleValidBeaconBlock(blockInput, peerIdStr, seenTimestampSec); + handleValidBeaconBlock({...blockInput, serializedData}, peerIdStr, seenTimestampSec); }, [GossipType.beacon_aggregate_and_proof]: async ({serializedData}, topic, _peer, seenTimestampSec) => { diff --git a/packages/types/src/types.ts b/packages/types/src/types.ts index 48987e2c52eb..b73c0a0d1b29 100644 --- a/packages/types/src/types.ts +++ b/packages/types/src/types.ts @@ -20,3 +20,5 @@ export enum BlockSource { export type SlotRootHex = {slot: Slot; root: RootHex}; export type SlotOptionalRoot = {slot: Slot; root?: RootHex}; +export type WithBytes> = T & {serializedData: Uint8Array}; +export type WithOptionalBytes> = T & {serializedData?: Uint8Array};