Skip to content
47 changes: 45 additions & 2 deletions packages/beacon-node/src/chain/blocks/blockInput/blockInput.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import {ForkName, ForkPostFulu, ForkPreDeneb, ForkPreGloas} from "@lodestar/params";
import {ForkName, ForkPostFulu, ForkPreDeneb, ForkPreGloas, NUMBER_OF_COLUMNS} from "@lodestar/params";
import {BeaconBlockBody, BlobIndex, ColumnIndex, SignedBeaconBlock, Slot, deneb, fulu} from "@lodestar/types";
import {fromHex, prettyBytes, toRootHex, withTimeout} from "@lodestar/utils";
import {VersionedHashes} from "../../../execution/index.js";
Expand Down Expand Up @@ -561,6 +561,7 @@ type BlockInputColumnsState =
| {
hasBlock: true;
hasAllData: true;
hasComputedAllData: boolean;
versionedHashes: VersionedHashes;
block: SignedBeaconBlock<ForkColumnsDA>;
source: SourceMeta;
Expand All @@ -569,18 +570,21 @@ type BlockInputColumnsState =
| {
hasBlock: true;
hasAllData: false;
hasComputedAllData: false;
versionedHashes: VersionedHashes;
block: SignedBeaconBlock<ForkColumnsDA>;
source: SourceMeta;
}
| {
hasBlock: false;
hasAllData: true;
hasComputedAllData: boolean;
versionedHashes: VersionedHashes;
}
| {
hasBlock: false;
hasAllData: false;
hasComputedAllData: false;
versionedHashes: VersionedHashes;
};
/**
Expand All @@ -598,6 +602,12 @@ export class BlockInputColumns extends AbstractBlockInput<ForkColumnsDA, fulu.Da
private columnsCache = new Map<ColumnIndex, ColumnWithSource>();
private readonly sampledColumns: ColumnIndex[];
private readonly custodyColumns: ColumnIndex[];
/**
* This promise resolves when all sampled columns are available
*
* This is different from `dataPromise` which resolves when all data is available or could become available (e.g. through reconstruction)
*/
protected computedDataPromise = createPromise<fulu.DataColumnSidecars>();

private constructor(
init: BlockInputInit,
Expand Down Expand Up @@ -626,6 +636,7 @@ export class BlockInputColumns extends AbstractBlockInput<ForkColumnsDA, fulu.Da
const state = {
hasBlock: true,
hasAllData,
hasComputedAllData: hasAllData,
versionedHashes: props.block.message.body.blobKzgCommitments.map(kzgCommitmentToVersionedHash),
block: props.block,
source: {
Comment thread
wemeetagain marked this conversation as resolved.
Expand All @@ -649,6 +660,7 @@ export class BlockInputColumns extends AbstractBlockInput<ForkColumnsDA, fulu.Da
blockInput.blockPromise.resolve(props.block);
if (hasAllData) {
blockInput.dataPromise.resolve([]);
blockInput.computedDataPromise.resolve([]);
}
return blockInput;
}
Expand All @@ -661,6 +673,7 @@ export class BlockInputColumns extends AbstractBlockInput<ForkColumnsDA, fulu.Da
const state: BlockInputColumnsState = {
hasBlock: false,
hasAllData,
hasComputedAllData: hasAllData as false,
Comment thread
wemeetagain marked this conversation as resolved.
Comment thread
nflaig marked this conversation as resolved.
versionedHashes: props.columnSidecar.kzgCommitments.map(kzgCommitmentToVersionedHash),
};
const init: BlockInputInit = {
Expand All @@ -674,6 +687,7 @@ export class BlockInputColumns extends AbstractBlockInput<ForkColumnsDA, fulu.Da
const blockInput = new BlockInputColumns(init, state, props.sampledColumns, props.custodyColumns);
if (hasAllData) {
blockInput.dataPromise.resolve([]);
blockInput.computedDataPromise.resolve([]);
}
return blockInput;
}
Expand Down Expand Up @@ -722,11 +736,14 @@ export class BlockInputColumns extends AbstractBlockInput<ForkColumnsDA, fulu.Da
const hasAllData =
(props.block.message.body as BeaconBlockBody<ForkPostFulu & ForkPreGloas>).blobKzgCommitments.length === 0 ||
this.state.hasAllData;
const hasComputedAllData =
props.block.message.body.blobKzgCommitments.length === 0 || this.state.hasComputedAllData;

this.state = {
...this.state,
hasBlock: true,
hasAllData,
hasComputedAllData,
block: props.block,
source: {
source: props.source,
Expand Down Expand Up @@ -774,17 +791,32 @@ export class BlockInputColumns extends AbstractBlockInput<ForkColumnsDA, fulu.Da
this.columnsCache.set(columnSidecar.index, {columnSidecar, source, seenTimestampSec, peerIdStr});

const sampledColumns = this.getSampledColumns();
const hasAllData = this.state.hasAllData || sampledColumns.length === this.sampledColumns.length;
const hasAllData =
// already hasAllData
this.state.hasAllData ||
// has all sampled columns
sampledColumns.length === this.sampledColumns.length ||
// has enough columns to reconstruct the rest
this.columnsCache.size >= NUMBER_OF_COLUMNS / 2;

const hasComputedAllData =
// has all sampled columns
sampledColumns.length === this.sampledColumns.length;
Comment thread
wemeetagain marked this conversation as resolved.

this.state = {
...this.state,
hasAllData: hasAllData || this.state.hasAllData,
hasComputedAllData: hasComputedAllData || this.state.hasComputedAllData,
timeCompleteSec: hasAllData ? seenTimestampSec : undefined,
} as BlockInputColumnsState;

if (hasAllData && sampledColumns !== null) {
this.dataPromise.resolve(sampledColumns);
}

if (hasComputedAllData && sampledColumns !== null) {
this.computedDataPromise.resolve(sampledColumns);
}
}

hasColumn(columnIndex: number): boolean {
Expand Down Expand Up @@ -859,4 +891,15 @@ export class BlockInputColumns extends AbstractBlockInput<ForkColumnsDA, fulu.Da
versionedHashes: this.state.versionedHashes,
};
}

hasComputedAllData(): boolean {
return this.state.hasComputedAllData;
}

waitForComputedAllData(timeout: number, signal?: AbortSignal): Promise<fulu.DataColumnSidecars> {
if (!this.state.hasComputedAllData) {
return withTimeout(() => this.computedDataPromise.promise, timeout, signal);
}
return Promise.resolve(this.getSampledColumns());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,15 @@ export async function writeBlockInputToDb(this: BeaconChain, blocksInputs: IBloc

// NOTE: Old data is pruned on archive
if (isBlockInputColumns(blockInput)) {
if (!blockInput.hasComputedAllData()) {
// Supernodes may only have a subset of the data columns by the time the block begins to be imported
// because full data availability can be assumed after NUMBER_OF_COLUMNS / 2 columns are available.
// Here, however, all data columns must be fully available/reconstructed before persisting to the DB.
await blockInput.waitForComputedAllData(BLOB_AVAILABILITY_TIMEOUT).catch(() => {
this.logger.debug("Failed to wait for computed all data", {slot, blockRoot: blockRootHex});
});
}

const {custodyColumns} = this.custodyConfig;
const blobsLen = (block.message as fulu.BeaconBlock).body.blobKzgCommitments.length;
let dataColumnsLen: number;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -576,7 +576,7 @@ function getSequentialHandlers(modules: ValidatorFnsModules, options: GossipHand
break;
}

if (!blockInput.hasAllData()) {
if (!blockInput.hasComputedAllData()) {
// immediately attempt fetch of data columns from execution engine
chain.getBlobsTracker.triggerGetBlobs(blockInput);
// if we've received at least half of the columns, trigger reconstruction of the rest
Expand Down
Loading