diff --git a/packages/beacon-node/src/api/impl/beacon/pool/index.ts b/packages/beacon-node/src/api/impl/beacon/pool/index.ts index 8edef0a8fe4f..47e5011f6261 100644 --- a/packages/beacon-node/src/api/impl/beacon/pool/index.ts +++ b/packages/beacon-node/src/api/impl/beacon/pool/index.ts @@ -65,10 +65,12 @@ export function getBeaconPoolApi({ beaconBlockRoot ); - const insertOutcome = chain.attestationPool.add(attestation); + if (network.attnetsService.shouldProcess(subnet, slot)) { + const insertOutcome = chain.attestationPool.add(attestation); + metrics?.opPool.attestationPoolInsertOutcome.inc({insertOutcome}); + } const sentPeers = await network.gossip.publishBeaconAttestation(attestation, subnet); metrics?.submitUnaggregatedAttestation(seenTimestampSec, indexedAttestation, subnet, sentPeers); - metrics?.opPool.attestationPoolInsertOutcome.inc({insertOutcome}); } catch (e) { errors.push(e as Error); logger.error( diff --git a/packages/beacon-node/src/metrics/metrics/lodestar.ts b/packages/beacon-node/src/metrics/metrics/lodestar.ts index 01db193822a2..7aab734b0ef1 100644 --- a/packages/beacon-node/src/metrics/metrics/lodestar.ts +++ b/packages/beacon-node/src/metrics/metrics/lodestar.ts @@ -305,6 +305,10 @@ export function createLodestarMetrics( help: "Count of unsubscribe_subnets calls", labelNames: ["subnet", "src"], }), + aggregatorSlotSubnetCount: register.gauge({ + name: "lodestar_attnets_service_aggregator_slot_subnet_total", + help: "Count of aggregator per slot and subnet", + }), }, syncnetsService: { diff --git a/packages/beacon-node/src/network/interface.ts b/packages/beacon-node/src/network/interface.ts index 2711eef83d4b..5ba8998ba748 100644 --- a/packages/beacon-node/src/network/interface.ts +++ b/packages/beacon-node/src/network/interface.ts @@ -13,7 +13,7 @@ import {INetworkEventBus} from "./events.js"; import {GossipBeaconNode} from "./gossip/index.js"; import {PeerAction, PeerScoreStats} from "./peers/index.js"; import {IReqRespBeaconNode} from "./reqresp/ReqRespBeaconNode.js"; -import {CommitteeSubscription} from "./subnets/index.js"; +import {AttnetsService, CommitteeSubscription} from "./subnets/index.js"; export type PeerSearchOptions = { supportsProtocols?: string[]; @@ -27,6 +27,7 @@ export interface INetwork { events: INetworkEventBus; reqResp: IReqRespBeaconNode; + attnetsService: AttnetsService; gossip: GossipBeaconNode; getEnr(): Promise; diff --git a/packages/beacon-node/src/network/subnets/attnetsService.ts b/packages/beacon-node/src/network/subnets/attnetsService.ts index 85be3b7fadb5..78d643b40ef6 100644 --- a/packages/beacon-node/src/network/subnets/attnetsService.ts +++ b/packages/beacon-node/src/network/subnets/attnetsService.ts @@ -8,7 +8,7 @@ import { SLOTS_PER_EPOCH, } from "@lodestar/params"; import {Epoch, Slot, ssz} from "@lodestar/types"; -import {Logger, randBetween} from "@lodestar/utils"; +import {Logger, MapDef, randBetween} from "@lodestar/utils"; import {shuffle} from "../../util/shuffle.js"; import {ChainEvent, IBeaconChain} from "../../chain/index.js"; import {GossipTopic, GossipType} from "../gossip/index.js"; @@ -45,6 +45,11 @@ export class AttnetsService implements IAttnetsService { private subscriptionsCommittee = new SubnetMap(); /** Same as `subscriptionsCommittee` but for long-lived subnets. May overlap with `subscriptionsCommittee` */ private subscriptionsRandom = new SubnetMap(); + /** + * Map of an aggregator at a slot and subnet + * Used to determine if we should process an attestation. + */ + private aggregatorSlotSubnet = new MapDef>(() => new Set()); /** * A collection of seen validators. These dictate how many random subnets we should be @@ -119,6 +124,7 @@ export class AttnetsService implements IAttnetsService { if (isAggregator) { // need exact slot here subnetsToSubscribe.push({subnet, toSlot: slot}); + this.aggregatorSlotSubnet.getOrDefault(slot).add(subnet); } } @@ -141,7 +147,10 @@ export class AttnetsService implements IAttnetsService { * Check if a subscription is still active before handling a gossip object */ shouldProcess(subnet: number, slot: Slot): boolean { - return this.subscriptionsCommittee.isActiveAtSlot(subnet, slot); + if (!this.aggregatorSlotSubnet.has(slot)) { + return false; + } + return this.aggregatorSlotSubnet.getOrDefault(slot).has(subnet); } /** Call ONLY ONCE: Two epoch before the fork, re-subscribe all existing random subscriptions to the new fork */ @@ -184,6 +193,7 @@ export class AttnetsService implements IAttnetsService { try { const slot = computeStartSlotAtEpoch(epoch); this.pruneExpiredKnownValidators(slot); + this.pruneExpiredAggregator(slot); } catch (e) { this.logger.error("Error on AttnetsService.onEpoch", {epoch}, e as Error); } @@ -246,6 +256,18 @@ export class AttnetsService implements IAttnetsService { if (deletedKnownValidators) this.rebalanceRandomSubnets(); } + /** + * No need to track aggregator for past slots. + * @param currentSlot + */ + private pruneExpiredAggregator(currentSlot: Slot): void { + for (const slot of this.aggregatorSlotSubnet.keys()) { + if (currentSlot > slot) { + this.aggregatorSlotSubnet.delete(slot); + } + } + } + /** * Called when we have new validators or expired validators. * knownValidators should be updated before this function. @@ -350,5 +372,10 @@ export class AttnetsService implements IAttnetsService { metrics.attnetsService.committeeSubnets.set(this.committeeSubnets.size); metrics.attnetsService.subscriptionsCommittee.set(this.subscriptionsCommittee.size); metrics.attnetsService.subscriptionsRandom.set(this.subscriptionsRandom.size); + let aggregatorCount = 0; + for (const subnets of this.aggregatorSlotSubnet.values()) { + aggregatorCount += subnets.size; + } + metrics.attnetsService.aggregatorSlotSubnetCount.set(aggregatorCount); } } diff --git a/packages/beacon-node/test/unit/network/attestationService.test.ts b/packages/beacon-node/test/unit/network/subnets/attnetsService.test.ts similarity index 90% rename from packages/beacon-node/test/unit/network/attestationService.test.ts rename to packages/beacon-node/test/unit/network/subnets/attnetsService.test.ts index 0cb25d744c99..d91bc1e8fd6e 100644 --- a/packages/beacon-node/test/unit/network/attestationService.test.ts +++ b/packages/beacon-node/test/unit/network/subnets/attnetsService.test.ts @@ -8,14 +8,14 @@ import { } from "@lodestar/params"; import {createBeaconConfig} from "@lodestar/config"; import {BeaconStateAllForks, getCurrentSlot} from "@lodestar/state-transition"; -import {MockBeaconChain} from "../../utils/mocks/chain/chain.js"; -import {generateState} from "../../utils/state.js"; -import {testLogger} from "../../utils/logger.js"; -import {MetadataController} from "../../../src/network/metadata.js"; -import {Eth2Gossipsub, GossipType} from "../../../src/network/gossip/index.js"; -import {AttnetsService, CommitteeSubscription, ShuffleFn} from "../../../src/network/subnets/index.js"; -import {ChainEvent, IBeaconChain} from "../../../src/chain/index.js"; -import {ZERO_HASH} from "../../../src/constants/index.js"; +import {MockBeaconChain} from "../../../utils/mocks/chain/chain.js"; +import {generateState} from "../../../utils/state.js"; +import {testLogger} from "../../../utils/logger.js"; +import {MetadataController} from "../../../../src/network/metadata.js"; +import {Eth2Gossipsub, GossipType} from "../../../../src/network/gossip/index.js"; +import {AttnetsService, CommitteeSubscription, ShuffleFn} from "../../../../src/network/subnets/index.js"; +import {ChainEvent, IBeaconChain} from "../../../../src/chain/index.js"; +import {ZERO_HASH} from "../../../../src/constants/index.js"; describe("AttnetsService", function () { const COMMITTEE_SUBNET_SUBSCRIPTION = 10; @@ -193,6 +193,7 @@ describe("AttnetsService", function () { randomSubnet = COMMITTEE_SUBNET_SUBSCRIPTION; const aggregatorSubscription: CommitteeSubscription = {...subscription, isAggregator: true}; service.addCommitteeSubscriptions([aggregatorSubscription]); + expect(service.shouldProcess(subscription.subnet, subscription.slot)).to.be.true; expect(service.getActiveSubnets()).to.be.deep.equal([{subnet: COMMITTEE_SUBNET_SUBSCRIPTION, toSlot: 101}]); // committee subnet is same to random subnet expect(gossipStub.subscribeTopic).to.be.calledOnce; @@ -202,4 +203,10 @@ describe("AttnetsService", function () { // don't unsubscribe bc random subnet is still there expect(gossipStub.unsubscribeTopic).to.be.not.called; }); + + it("should not process if no aggregator at dutied slot", () => { + expect(subscription.isAggregator).to.be.false; + service.addCommitteeSubscriptions([subscription]); + expect(service.shouldProcess(subscription.subnet, subscription.slot)).to.be.false; + }); });