From a518a59d3531f8f63dc0b8579429cde4aab8b14c Mon Sep 17 00:00:00 2001 From: dapplion <35266934+dapplion@users.noreply.github.com> Date: Tue, 14 Feb 2023 15:40:36 +0800 Subject: [PATCH 1/3] Track validator_monitor_validators_in_sync_committee --- .../src/api/impl/validator/index.ts | 6 ++++++ .../src/metrics/metrics/lodestar.ts | 6 +++++- .../src/metrics/validatorMonitor.ts | 20 ++++++++++++++++++- 3 files changed, 30 insertions(+), 2 deletions(-) diff --git a/packages/beacon-node/src/api/impl/validator/index.ts b/packages/beacon-node/src/api/impl/validator/index.ts index cb2bff614194..7a109777bc70 100644 --- a/packages/beacon-node/src/api/impl/validator/index.ts +++ b/packages/beacon-node/src/api/impl/validator/index.ts @@ -647,6 +647,12 @@ export function getValidatorApi({ } network.prepareSyncCommitteeSubnets(subs); + + if (metrics) { + for (const subscription of subscriptions) { + metrics.registerLocalValidatorInSyncCommittee(subscription.validatorIndex, subscription.untilEpoch); + } + } }, async prepareBeaconProposer(proposers) { diff --git a/packages/beacon-node/src/metrics/metrics/lodestar.ts b/packages/beacon-node/src/metrics/metrics/lodestar.ts index c0ebcdd8c19e..c93db80f974d 100644 --- a/packages/beacon-node/src/metrics/metrics/lodestar.ts +++ b/packages/beacon-node/src/metrics/metrics/lodestar.ts @@ -746,7 +746,11 @@ export function createLodestarMetrics( validatorsConnected: register.gauge({ name: "validator_monitor_validators", help: "Count of validators that are specifically monitored by this beacon node", - labelNames: ["index"], + }), + + validatorsInSyncCommittee: register.gauge({ + name: "validator_monitor_validators_in_sync_committee", + help: "Count of validators monitored by this beacon node that are part of sync committee", }), // Validator Monitor Metrics (per-epoch summaries) diff --git a/packages/beacon-node/src/metrics/validatorMonitor.ts b/packages/beacon-node/src/metrics/validatorMonitor.ts index 960bc1089596..c71d97234795 100644 --- a/packages/beacon-node/src/metrics/validatorMonitor.ts +++ b/packages/beacon-node/src/metrics/validatorMonitor.ts @@ -20,6 +20,7 @@ export enum OpSource { export interface IValidatorMonitor { registerLocalValidator(index: number): void; + registerLocalValidatorInSyncCommittee(index: number, untilEpoch: number): void; registerValidatorStatuses(currentEpoch: Epoch, statuses: IAttesterStatus[], balances?: number[]): void; registerBeaconBlock(src: OpSource, seenTimestampSec: Seconds, block: allForks.BeaconBlock): void; registerImportedBlock(block: allForks.BeaconBlock, data: {proposerBalanceDelta: number}): void; @@ -160,6 +161,7 @@ type MonitoredValidator = { index: number; /// A history of the validator over time. */ summaries: Map; + inSyncCommitteeUntilEpoch: number; }; export function createValidatorMonitor( @@ -176,7 +178,14 @@ export function createValidatorMonitor( return { registerLocalValidator(index) { if (!validators.has(index)) { - validators.set(index, {index, summaries: new Map()}); + validators.set(index, {index, summaries: new Map(), inSyncCommitteeUntilEpoch: -1}); + } + }, + + registerLocalValidatorInSyncCommittee(index, untilEpoch) { + const validator = validators.get(index); + if (validator) { + validator.inSyncCommitteeUntilEpoch = Math.max(untilEpoch, validator.inSyncCommitteeUntilEpoch ?? -1); } }, @@ -445,6 +454,8 @@ export function createValidatorMonitor( metrics.validatorMonitor.prevEpochAttestationBlockInclusions.reset(); metrics.validatorMonitor.prevEpochAttestationBlockMinInclusionDistance.reset(); + let validatorsInSyncCommittee = 0; + for (const validator of validators.values()) { const summary = validator.summaries.get(previousEpoch); if (!summary) { @@ -472,7 +483,14 @@ export function createValidatorMonitor( metrics.validatorMonitor.prevEpochAggregatesTotal.observe(summary.aggregates); if (summary.aggregateMinDelay !== null) metrics.validatorMonitor.prevEpochAggregatesMinDelaySeconds.observe(summary.aggregateMinDelay); + + // Sync committee + if (validator.inSyncCommitteeUntilEpoch >= epoch) { + validatorsInSyncCommittee++; + } } + + metrics.validatorMonitor.validatorsInSyncCommittee.set(validatorsInSyncCommittee); }, }; } From 77661f6d3e916804fc9607bbaf8697bbdd8ee5e3 Mon Sep 17 00:00:00 2001 From: dapplion <35266934+dapplion@users.noreply.github.com> Date: Tue, 14 Feb 2023 16:03:42 +0800 Subject: [PATCH 2/3] Add sync committee participation metrics --- .../src/chain/blocks/importBlock.ts | 11 ++++- .../src/metrics/metrics/lodestar.ts | 8 ++++ .../src/metrics/validatorMonitor.ts | 45 ++++++++++++++++--- 3 files changed, 57 insertions(+), 7 deletions(-) diff --git a/packages/beacon-node/src/chain/blocks/importBlock.ts b/packages/beacon-node/src/chain/blocks/importBlock.ts index b24667041cd3..13b00f8e5644 100644 --- a/packages/beacon-node/src/chain/blocks/importBlock.ts +++ b/packages/beacon-node/src/chain/blocks/importBlock.ts @@ -1,5 +1,5 @@ -import {capella, ssz, allForks} from "@lodestar/types"; -import {MAX_SEED_LOOKAHEAD, SLOTS_PER_EPOCH} from "@lodestar/params"; +import {capella, ssz, allForks, altair} from "@lodestar/types"; +import {ForkSeq, MAX_SEED_LOOKAHEAD, SLOTS_PER_EPOCH} from "@lodestar/params"; import {toHexString} from "@chainsafe/ssz"; import { CachedBeaconStateAltair, @@ -384,6 +384,13 @@ export async function importBlock( this.metrics?.parentBlockDistance.observe(block.message.slot - parentBlockSlot); this.metrics?.proposerBalanceDeltaAny.observe(fullyVerifiedBlock.proposerBalanceDelta); this.metrics?.registerImportedBlock(block.message, fullyVerifiedBlock); + if (this.config.getForkSeq(block.message.slot) >= ForkSeq.altair) { + this.metrics?.registerImportedBlockSyncAggregate( + blockEpoch, + (block as altair.SignedBeaconBlock).message.body.syncAggregate, + fullyVerifiedBlock.postState.epochCtx.currentSyncCommitteeIndexed.validatorIndices + ); + } const advancedSlot = this.clock.slotWithFutureTolerance(REPROCESS_MIN_TIME_TO_NEXT_SLOT_SEC); diff --git a/packages/beacon-node/src/metrics/metrics/lodestar.ts b/packages/beacon-node/src/metrics/metrics/lodestar.ts index c93db80f974d..01f87802a45a 100644 --- a/packages/beacon-node/src/metrics/metrics/lodestar.ts +++ b/packages/beacon-node/src/metrics/metrics/lodestar.ts @@ -854,6 +854,14 @@ export function createLodestarMetrics( help: "The min delay between when the validator should send the aggregate and when it was received", buckets: [0.1, 0.25, 0.5, 1, 2, 5, 10], }), + prevEpochSyncCommitteeHits: register.gauge({ + name: "validator_monitor_prev_epoch_sync_committee_hits", + help: "Count of times in prev epoch connected validators participated in imported block's syncAggregate", + }), + prevEpochSyncCommitteeMisses: register.gauge({ + name: "validator_monitor_prev_epoch_sync_committee_misses", + help: "Count of times in prev epoch connected validators fail to participate in imported block's syncAggregate", + }), // Validator Monitor Metrics (real-time) diff --git a/packages/beacon-node/src/metrics/validatorMonitor.ts b/packages/beacon-node/src/metrics/validatorMonitor.ts index c71d97234795..b2236d2104fc 100644 --- a/packages/beacon-node/src/metrics/validatorMonitor.ts +++ b/packages/beacon-node/src/metrics/validatorMonitor.ts @@ -1,6 +1,6 @@ import {computeEpochAtSlot, IAttesterStatus, parseAttesterFlags} from "@lodestar/state-transition"; import {ILogger} from "@lodestar/utils"; -import {allForks} from "@lodestar/types"; +import {allForks, altair} from "@lodestar/types"; import {IChainForkConfig} from "@lodestar/config"; import {MIN_ATTESTATION_INCLUSION_DELAY, SLOTS_PER_EPOCH} from "@lodestar/params"; import {Epoch, Slot, ValidatorIndex} from "@lodestar/types"; @@ -20,10 +20,15 @@ export enum OpSource { export interface IValidatorMonitor { registerLocalValidator(index: number): void; - registerLocalValidatorInSyncCommittee(index: number, untilEpoch: number): void; + registerLocalValidatorInSyncCommittee(index: number, untilEpoch: Epoch): void; registerValidatorStatuses(currentEpoch: Epoch, statuses: IAttesterStatus[], balances?: number[]): void; registerBeaconBlock(src: OpSource, seenTimestampSec: Seconds, block: allForks.BeaconBlock): void; registerImportedBlock(block: allForks.BeaconBlock, data: {proposerBalanceDelta: number}): void; + registerImportedBlockSyncAggregate( + epoch: Epoch, + syncAggregate: altair.SyncAggregate, + syncCommitteeIndices: number[] + ): void; submitUnaggregatedAttestation( seenTimestampSec: number, indexedAttestation: IndexedAttestation, @@ -122,6 +127,10 @@ type EpochSummary = { aggregates: number; /** The delay between when the aggregate should have been produced and when it was observed. */ aggregateMinDelay: Seconds | null; + /** Count of times validator expected in sync aggregate participated */ + syncCommitteeHits: number; + /** Count of times validator expected in sync aggregate failed to participate */ + syncCommitteeMisses: number; }; function withEpochSummary(validator: MonitoredValidator, epoch: Epoch, fn: (summary: EpochSummary) => void): void { @@ -138,6 +147,8 @@ function withEpochSummary(validator: MonitoredValidator, epoch: Epoch, fn: (summ aggregates: 0, aggregateMinDelay: null, attestationCorrectHead: null, + syncCommitteeHits: 0, + syncCommitteeMisses: 0, }; validator.summaries.set(epoch, summary); } @@ -291,6 +302,21 @@ export function createValidatorMonitor( } }, + registerImportedBlockSyncAggregate(epoch, syncAggregate, syncCommitteeIndices) { + for (let i = 0; i < syncCommitteeIndices.length; i++) { + const validator = validators.get(syncCommitteeIndices[i]); + if (validator) { + withEpochSummary(validator, epoch, (summary) => { + if (syncAggregate.syncCommitteeBits.get(i)) { + summary.syncCommitteeHits++; + } else { + summary.syncCommitteeMisses++; + } + }); + } + } + }, + submitUnaggregatedAttestation(seenTimestampSec, indexedAttestation, subnet, sentPeers) { const data = indexedAttestation.data; // Returns the duration between when the attestation `data` could be produced (1/3rd through the slot) and `seenTimestamp`. @@ -455,8 +481,16 @@ export function createValidatorMonitor( metrics.validatorMonitor.prevEpochAttestationBlockMinInclusionDistance.reset(); let validatorsInSyncCommittee = 0; + let prevEpochSyncCommitteeHits = 0; + let prevEpochSyncCommitteeMisses = 0; for (const validator of validators.values()) { + // Participation in sync committee + if (validator.inSyncCommitteeUntilEpoch >= epoch) { + validatorsInSyncCommittee++; + } + + // Prev-epoch summary const summary = validator.summaries.get(previousEpoch); if (!summary) { continue; @@ -485,12 +519,13 @@ export function createValidatorMonitor( metrics.validatorMonitor.prevEpochAggregatesMinDelaySeconds.observe(summary.aggregateMinDelay); // Sync committee - if (validator.inSyncCommitteeUntilEpoch >= epoch) { - validatorsInSyncCommittee++; - } + prevEpochSyncCommitteeHits += summary.syncCommitteeHits; + prevEpochSyncCommitteeMisses += summary.syncCommitteeMisses; } metrics.validatorMonitor.validatorsInSyncCommittee.set(validatorsInSyncCommittee); + metrics.validatorMonitor.prevEpochSyncCommitteeHits.set(prevEpochSyncCommitteeHits); + metrics.validatorMonitor.prevEpochSyncCommitteeMisses.set(prevEpochSyncCommitteeMisses); }, }; } From 13c2d8ddc61ff22cbeb273bbc72093e35027233c Mon Sep 17 00:00:00 2001 From: dapplion <35266934+dapplion@users.noreply.github.com> Date: Tue, 14 Feb 2023 16:28:09 +0800 Subject: [PATCH 3/3] Add sync contribution aggregate metrics --- .../src/api/impl/validator/index.ts | 7 +- .../src/chain/blocks/importBlock.ts | 2 +- .../syncCommitteeContributionAndProof.ts | 16 +++-- .../src/metrics/metrics/lodestar.ts | 9 +++ .../src/metrics/validatorMonitor.ts | 69 +++++++++++++------ .../src/network/gossip/handlers/index.ts | 5 +- 6 files changed, 75 insertions(+), 33 deletions(-) diff --git a/packages/beacon-node/src/api/impl/validator/index.ts b/packages/beacon-node/src/api/impl/validator/index.ts index 7a109777bc70..5bde52040d9c 100644 --- a/packages/beacon-node/src/api/impl/validator/index.ts +++ b/packages/beacon-node/src/api/impl/validator/index.ts @@ -564,12 +564,15 @@ export function getValidatorApi({ contributionAndProofs.map(async (contributionAndProof, i) => { try { // TODO: Validate in batch - const {syncCommitteeParticipants} = await validateSyncCommitteeGossipContributionAndProof( + const {syncCommitteeParticipantIndices} = await validateSyncCommitteeGossipContributionAndProof( chain, contributionAndProof, true // skip known participants check ); - chain.syncContributionAndProofPool.add(contributionAndProof.message, syncCommitteeParticipants); + chain.syncContributionAndProofPool.add( + contributionAndProof.message, + syncCommitteeParticipantIndices.length + ); await network.gossip.publishContributionAndProof(contributionAndProof); } catch (e) { errors.push(e as Error); diff --git a/packages/beacon-node/src/chain/blocks/importBlock.ts b/packages/beacon-node/src/chain/blocks/importBlock.ts index 13b00f8e5644..c3c259b8be9b 100644 --- a/packages/beacon-node/src/chain/blocks/importBlock.ts +++ b/packages/beacon-node/src/chain/blocks/importBlock.ts @@ -385,7 +385,7 @@ export async function importBlock( this.metrics?.proposerBalanceDeltaAny.observe(fullyVerifiedBlock.proposerBalanceDelta); this.metrics?.registerImportedBlock(block.message, fullyVerifiedBlock); if (this.config.getForkSeq(block.message.slot) >= ForkSeq.altair) { - this.metrics?.registerImportedBlockSyncAggregate( + this.metrics?.registerSyncAggregateInBlock( blockEpoch, (block as altair.SignedBeaconBlock).message.body.syncAggregate, fullyVerifiedBlock.postState.epochCtx.currentSyncCommitteeIndexed.validatorIndices diff --git a/packages/beacon-node/src/chain/validation/syncCommitteeContributionAndProof.ts b/packages/beacon-node/src/chain/validation/syncCommitteeContributionAndProof.ts index 98266cdaa1ac..3f9351f2d0d3 100644 --- a/packages/beacon-node/src/chain/validation/syncCommitteeContributionAndProof.ts +++ b/packages/beacon-node/src/chain/validation/syncCommitteeContributionAndProof.ts @@ -17,7 +17,7 @@ export async function validateSyncCommitteeGossipContributionAndProof( chain: IBeaconChain, signedContributionAndProof: altair.SignedContributionAndProof, skipValidationKnownParticipants = false -): Promise<{syncCommitteeParticipants: number}> { +): Promise<{syncCommitteeParticipantIndices: ValidatorIndex[]}> { const contributionAndProof = signedContributionAndProof.message; const {contribution, aggregatorIndex} = contributionAndProof; const {subcommitteeIndex, slot} = contribution; @@ -53,8 +53,8 @@ export async function validateSyncCommitteeGossipContributionAndProof( } // [REJECT] The contribution has participants -- that is, any(contribution.aggregation_bits) - const syncCommitteeIndices = getContributionIndices(headState as CachedBeaconStateAltair, contribution); - if (!syncCommitteeIndices.length) { + const syncCommitteeParticipantIndices = getContributionIndices(headState as CachedBeaconStateAltair, contribution); + if (syncCommitteeParticipantIndices.length === 0) { throw new SyncCommitteeError(GossipAction.REJECT, { code: SyncCommitteeErrorCode.NO_PARTICIPANT, }); @@ -73,7 +73,9 @@ export async function validateSyncCommitteeGossipContributionAndProof( // i.e. state.validators[contribution_and_proof.aggregator_index].pubkey in get_sync_subcommittee_pubkeys(state, contribution.subcommittee_index). // > Checked in validateGossipSyncCommitteeExceptSig() - const pubkeys = syncCommitteeIndices.map((validatorIndex) => headState.epochCtx.index2pubkey[validatorIndex]); + const participantPubkeys = syncCommitteeParticipantIndices.map( + (validatorIndex) => headState.epochCtx.index2pubkey[validatorIndex] + ); const signatureSets = [ // [REJECT] The contribution_and_proof.selection_proof is a valid signature of the SyncAggregatorSelectionData // derived from the contribution by the validator with index contribution_and_proof.aggregator_index. @@ -84,7 +86,7 @@ export async function validateSyncCommitteeGossipContributionAndProof( // [REJECT] The aggregate signature is valid for the message beacon_block_root and aggregate pubkey derived from // the participation info in aggregation_bits for the subcommittee specified by the contribution.subcommittee_index. - getSyncCommitteeContributionSignatureSet(headState as CachedBeaconStateAltair, contribution, pubkeys), + getSyncCommitteeContributionSignatureSet(headState as CachedBeaconStateAltair, contribution, participantPubkeys), ]; if (!(await chain.bls.verifySignatureSets(signatureSets, {batchable: true}))) { @@ -94,9 +96,9 @@ export async function validateSyncCommitteeGossipContributionAndProof( } // no need to add to seenSyncCommittteeContributionCache here, gossip handler will do that - chain.seenContributionAndProof.add(contributionAndProof, syncCommitteeIndices.length); + chain.seenContributionAndProof.add(contributionAndProof, syncCommitteeParticipantIndices.length); - return {syncCommitteeParticipants: syncCommitteeIndices.length}; + return {syncCommitteeParticipantIndices}; } /** diff --git a/packages/beacon-node/src/metrics/metrics/lodestar.ts b/packages/beacon-node/src/metrics/metrics/lodestar.ts index 01f87802a45a..4c34563ff944 100644 --- a/packages/beacon-node/src/metrics/metrics/lodestar.ts +++ b/packages/beacon-node/src/metrics/metrics/lodestar.ts @@ -862,6 +862,11 @@ export function createLodestarMetrics( name: "validator_monitor_prev_epoch_sync_committee_misses", help: "Count of times in prev epoch connected validators fail to participate in imported block's syncAggregate", }), + prevEpochSyncSignatureAggregateInclusions: register.histogram({ + name: "validator_monitor_prev_epoch_sync_signature_aggregate_inclusions", + help: "The count of times a sync signature was seen inside an aggregate", + buckets: [0, 1, 2, 3, 5, 10], + }), // Validator Monitor Metrics (real-time) @@ -915,6 +920,10 @@ export function createLodestarMetrics( help: "The excess slots (beyond the minimum delay) between the attestation slot and the block slot", buckets: [0.1, 0.25, 0.5, 1, 2, 5, 10], }), + syncSignatureInAggregateTotal: register.gauge({ + name: "validator_monitor_sync_signature_in_aggregate_total", + help: "Number of times a sync signature has been seen in an aggregate", + }), beaconBlockTotal: register.gauge<"src">({ name: "validator_monitor_beacon_block_total", help: "Total number of beacon blocks seen", diff --git a/packages/beacon-node/src/metrics/validatorMonitor.ts b/packages/beacon-node/src/metrics/validatorMonitor.ts index b2236d2104fc..e526257df5cb 100644 --- a/packages/beacon-node/src/metrics/validatorMonitor.ts +++ b/packages/beacon-node/src/metrics/validatorMonitor.ts @@ -24,11 +24,6 @@ export interface IValidatorMonitor { registerValidatorStatuses(currentEpoch: Epoch, statuses: IAttesterStatus[], balances?: number[]): void; registerBeaconBlock(src: OpSource, seenTimestampSec: Seconds, block: allForks.BeaconBlock): void; registerImportedBlock(block: allForks.BeaconBlock, data: {proposerBalanceDelta: number}): void; - registerImportedBlockSyncAggregate( - epoch: Epoch, - syncAggregate: altair.SyncAggregate, - syncCommitteeIndices: number[] - ): void; submitUnaggregatedAttestation( seenTimestampSec: number, indexedAttestation: IndexedAttestation, @@ -47,6 +42,11 @@ export interface IValidatorMonitor { indexedAttestation: IndexedAttestation ): void; registerAttestationInBlock(indexedAttestation: IndexedAttestation, parentSlot: Slot, correctHead: boolean): void; + registerGossipSyncContributionAndProof( + syncContributionAndProof: altair.ContributionAndProof, + syncCommitteeParticipantIndices: ValidatorIndex[] + ): void; + registerSyncAggregateInBlock(epoch: Epoch, syncAggregate: altair.SyncAggregate, syncCommitteeIndices: number[]): void; scrapeMetrics(slotClock: Slot): void; } @@ -131,6 +131,8 @@ type EpochSummary = { syncCommitteeHits: number; /** Count of times validator expected in sync aggregate failed to participate */ syncCommitteeMisses: number; + /** Number of times a validator's sync signature was seen in an aggregate */ + syncSignatureAggregateInclusions: number; }; function withEpochSummary(validator: MonitoredValidator, epoch: Epoch, fn: (summary: EpochSummary) => void): void { @@ -149,6 +151,7 @@ function withEpochSummary(validator: MonitoredValidator, epoch: Epoch, fn: (summ attestationCorrectHead: null, syncCommitteeHits: 0, syncCommitteeMisses: 0, + syncSignatureAggregateInclusions: 0, }; validator.summaries.set(epoch, summary); } @@ -302,21 +305,6 @@ export function createValidatorMonitor( } }, - registerImportedBlockSyncAggregate(epoch, syncAggregate, syncCommitteeIndices) { - for (let i = 0; i < syncCommitteeIndices.length; i++) { - const validator = validators.get(syncCommitteeIndices[i]); - if (validator) { - withEpochSummary(validator, epoch, (summary) => { - if (syncAggregate.syncCommitteeBits.get(i)) { - summary.syncCommitteeHits++; - } else { - summary.syncCommitteeMisses++; - } - }); - } - } - }, - submitUnaggregatedAttestation(seenTimestampSec, indexedAttestation, subnet, sentPeers) { const data = indexedAttestation.data; // Returns the duration between when the attestation `data` could be produced (1/3rd through the slot) and `seenTimestamp`. @@ -452,6 +440,36 @@ export function createValidatorMonitor( } }, + registerGossipSyncContributionAndProof(syncContributionAndProof, syncCommitteeParticipantIndices) { + const epoch = computeEpochAtSlot(syncContributionAndProof.contribution.slot); + + for (const index of syncCommitteeParticipantIndices) { + const validator = validators.get(index); + if (validator) { + metrics.validatorMonitor.syncSignatureInAggregateTotal.inc(); + + withEpochSummary(validator, epoch, (summary) => { + summary.syncSignatureAggregateInclusions += 1; + }); + } + } + }, + + registerSyncAggregateInBlock(epoch, syncAggregate, syncCommitteeIndices) { + for (let i = 0; i < syncCommitteeIndices.length; i++) { + const validator = validators.get(syncCommitteeIndices[i]); + if (validator) { + withEpochSummary(validator, epoch, (summary) => { + if (syncAggregate.syncCommitteeBits.get(i)) { + summary.syncCommitteeHits++; + } else { + summary.syncCommitteeMisses++; + } + }); + } + } + }, + /** * Scrape `self` for metrics. * Should be called whenever Prometheus is scraping. @@ -479,6 +497,7 @@ export function createValidatorMonitor( metrics.validatorMonitor.prevEpochAttestationAggregateInclusions.reset(); metrics.validatorMonitor.prevEpochAttestationBlockInclusions.reset(); metrics.validatorMonitor.prevEpochAttestationBlockMinInclusionDistance.reset(); + metrics.validatorMonitor.prevEpochSyncSignatureAggregateInclusions.reset(); let validatorsInSyncCommittee = 0; let prevEpochSyncCommitteeHits = 0; @@ -486,7 +505,8 @@ export function createValidatorMonitor( for (const validator of validators.values()) { // Participation in sync committee - if (validator.inSyncCommitteeUntilEpoch >= epoch) { + const validatorInSyncCommittee = validator.inSyncCommitteeUntilEpoch >= epoch; + if (validatorInSyncCommittee) { validatorsInSyncCommittee++; } @@ -521,6 +541,13 @@ export function createValidatorMonitor( // Sync committee prevEpochSyncCommitteeHits += summary.syncCommitteeHits; prevEpochSyncCommitteeMisses += summary.syncCommitteeMisses; + + // Only observe if included in sync committee to prevent distorting metrics + if (validatorInSyncCommittee) { + metrics.validatorMonitor.prevEpochSyncSignatureAggregateInclusions.observe( + summary.syncSignatureAggregateInclusions + ); + } } metrics.validatorMonitor.validatorsInSyncCommittee.set(validatorsInSyncCommittee); diff --git a/packages/beacon-node/src/network/gossip/handlers/index.ts b/packages/beacon-node/src/network/gossip/handlers/index.ts index 0261f6637f64..ebc6afd01018 100644 --- a/packages/beacon-node/src/network/gossip/handlers/index.ts +++ b/packages/beacon-node/src/network/gossip/handlers/index.ts @@ -301,7 +301,7 @@ export function getGossipHandlers(modules: ValidatorFnsModules, options: GossipH }, [GossipType.sync_committee_contribution_and_proof]: async (contributionAndProof) => { - const {syncCommitteeParticipants} = await validateSyncCommitteeGossipContributionAndProof( + const {syncCommitteeParticipantIndices} = await validateSyncCommitteeGossipContributionAndProof( chain, contributionAndProof ).catch((e) => { @@ -312,9 +312,10 @@ export function getGossipHandlers(modules: ValidatorFnsModules, options: GossipH }); // Handler + metrics?.registerGossipSyncContributionAndProof(contributionAndProof.message, syncCommitteeParticipantIndices); try { - chain.syncContributionAndProofPool.add(contributionAndProof.message, syncCommitteeParticipants); + chain.syncContributionAndProofPool.add(contributionAndProof.message, syncCommitteeParticipantIndices.length); } catch (e) { logger.error("Error adding to contributionAndProof pool", {}, e as Error); }