From a91af31df6e69df7a9cc66c0b461c1bf7d8dea48 Mon Sep 17 00:00:00 2001 From: Nico Flaig Date: Wed, 25 Feb 2026 10:07:22 +0000 Subject: [PATCH 1/4] fix: skip prometheus metrics trackProtocolStream for identify --- .../beacon-node/src/network/libp2p/index.ts | 32 +++++++++++++++---- 1 file changed, 25 insertions(+), 7 deletions(-) diff --git a/packages/beacon-node/src/network/libp2p/index.ts b/packages/beacon-node/src/network/libp2p/index.ts index 9aa53d6fd10b..12c282210732 100644 --- a/packages/beacon-node/src/network/libp2p/index.ts +++ b/packages/beacon-node/src/network/libp2p/index.ts @@ -73,6 +73,30 @@ export async function createNodeJsLibp2p( noiseCrypto.chaCha20Poly1305Encrypt = asCrypto.chaCha20Poly1305Encrypt; } + const libp2pMetrics = nodeJsLibp2pOpts.metrics + ? ((components: LodestarComponents) => { + const metrics = prometheusMetrics({ + collectDefaultMetrics: false, + preserveExistingMetrics: true, + registry: nodeJsLibp2pOpts.metricsRegistry, + })(components); + + // Work around identify EOF race: + // `trackProtocolStream` attaches a `message` listener immediately after protocol + // negotiation. For `/ipfs/id/1.0.0`, identify() adds its own reader later and can + // miss the first response frame when metrics listener drains events first. + const originalTrackProtocolStream = metrics.trackProtocolStream.bind(metrics); + metrics.trackProtocolStream = ((stream) => { + if (stream.protocol === "/ipfs/id/1.0.0") { + return; + } + originalTrackProtocolStream(stream); + }) as typeof metrics.trackProtocolStream; + + return metrics; + }) + : undefined; + return createLibp2p({ privateKey, nodeInfo: { @@ -101,13 +125,7 @@ export async function createNodeJsLibp2p( ], streamMuxers: [mplex({disconnectThreshold})], peerDiscovery, - metrics: nodeJsLibp2pOpts.metrics - ? prometheusMetrics({ - collectDefaultMetrics: false, - preserveExistingMetrics: true, - registry: nodeJsLibp2pOpts.metricsRegistry, - }) - : undefined, + metrics: libp2pMetrics, connectionManager: { // dialer config maxParallelDials: 100, From 50743ad30e421096d66c5e22cdd1c62c9693f1d0 Mon Sep 17 00:00:00 2001 From: Nico Flaig Date: Wed, 25 Feb 2026 10:15:39 +0000 Subject: [PATCH 2/4] Revert #8955 --- .../src/network/peers/peerManager.ts | 65 +------- .../e2e/network/peers/peerManager.test.ts | 153 ------------------ 2 files changed, 4 insertions(+), 214 deletions(-) diff --git a/packages/beacon-node/src/network/peers/peerManager.ts b/packages/beacon-node/src/network/peers/peerManager.ts index 86f19c979a6d..ca8c8892dd6d 100644 --- a/packages/beacon-node/src/network/peers/peerManager.ts +++ b/packages/beacon-node/src/network/peers/peerManager.ts @@ -1,4 +1,4 @@ -import {Connection, type IdentifyResult, PeerId, PrivateKey} from "@libp2p/interface"; +import {Connection, PeerId, PrivateKey} from "@libp2p/interface"; import {BitArray} from "@chainsafe/ssz"; import {BeaconConfig} from "@lodestar/config"; import {LoggerNode} from "@lodestar/logger/node"; @@ -162,9 +162,6 @@ export class PeerManager { // A single map of connected peers with all necessary data to handle PINGs, STATUS, and metrics private connectedPeers: Map; - /** Track one in-flight identify call per peer/connection id */ - private readonly identifyInProgress = new Map(); - private opts: PeerManagerOpts; private intervals: NodeJS.Timeout[] = []; @@ -195,7 +192,6 @@ export class PeerManager { this.libp2p.services.components.events.addEventListener(Libp2pEvent.connectionOpen, this.onLibp2pPeerConnect); this.libp2p.services.components.events.addEventListener(Libp2pEvent.connectionClose, this.onLibp2pPeerDisconnect); - this.libp2p.services.components.events.addEventListener(Libp2pEvent.peerIdentify, this.onPeerIdentify); this.networkEventBus.on(NetworkEvent.reqRespRequest, this.onRequest); this.lastStatus = this.statusCache.get(); @@ -239,7 +235,6 @@ export class PeerManager { Libp2pEvent.connectionClose, this.onLibp2pPeerDisconnect ); - this.libp2p.services.components.events.removeEventListener(Libp2pEvent.peerIdentify, this.onPeerIdentify); this.networkEventBus.off(NetworkEvent.reqRespRequest, this.onRequest); for (const interval of this.intervals) clearInterval(interval); } @@ -490,25 +485,7 @@ export class PeerManager { // peers that close identify right after connection open or turn out to be // irrelevant. if (peerData?.agentVersion === null) { - const peerIdStr = peer.toString(); - const connection = getConnection(this.libp2p, peerIdStr); - if (!connection || connection.status !== "open") { - this.logger.debug("Peer has no open connection for identify", {peerId: prettyPrintPeerId(peer)}); - return; - } - - const identifyKey = connection.id; - if (this.identifyInProgress.get(peerIdStr) === identifyKey) { - return; - } - - this.identifyInProgress.set(peerIdStr, identifyKey); - void this.identifyPeer(peerIdStr, prettyPrintPeerId(peer), connection, identifyKey).finally(() => { - // Clear only if this identify attempt is still the active one for this peer - if (this.identifyInProgress.get(peerIdStr) === identifyKey) { - this.identifyInProgress.delete(peerIdStr); - } - }); + void this.identifyPeer(peer.toString(), prettyPrintPeerId(peer), getConnection(this.libp2p, peer.toString())); } } } @@ -867,7 +844,6 @@ export class PeerManager { // remove the ping and status timer for the peer this.connectedPeers.delete(peerIdStr); - this.identifyInProgress.delete(peerIdStr); this.logger.verbose(logMessage, logContext); this.networkEventBus.emit(NetworkEvent.peerDisconnected, {peer: peerIdStr}); @@ -885,47 +861,14 @@ export class PeerManager { } } - /** - * Consume successful identify results from libp2p events. - * This captures agentVersion from identify-push or successful inbound/outbound identify, - * even if our explicit identify request failed earlier. - */ - private onPeerIdentify = (evt: CustomEvent): void => { - const {peerId, agentVersion} = evt.detail; - if (!agentVersion) return; - - const peerIdStr = peerId.toString(); - const peerData = this.connectedPeers.get(peerIdStr); - if (!peerData) return; - - peerData.agentVersion = agentVersion; - peerData.agentClient = getKnownClientFromAgentVersion(agentVersion); - this.identifyInProgress.delete(peerIdStr); - }; - - private async identifyPeer( - peerIdStr: string, - peerIdPretty: string, - connection: Connection, - identifyKey: string - ): Promise { - if (this.identifyInProgress.get(peerIdStr) !== identifyKey) { - return; - } - - if (connection.status !== "open") { + private async identifyPeer(peerIdStr: string, peerIdPretty: string, connection?: Connection): Promise { + if (!connection || connection.status !== "open") { this.logger.debug("Peer has no open connection for identify", {peerId: peerIdPretty}); return; } try { const result = await this.libp2p.services.identify.identify(connection); - - // A newer identify attempt may have superseded this one (e.g. reconnect). - if (this.identifyInProgress.get(peerIdStr) !== identifyKey) { - return; - } - const agentVersion = result.agentVersion; if (agentVersion) { const connectedPeerData = this.connectedPeers.get(peerIdStr); diff --git a/packages/beacon-node/test/e2e/network/peers/peerManager.test.ts b/packages/beacon-node/test/e2e/network/peers/peerManager.test.ts index 2cde954688f6..edb5048803d9 100644 --- a/packages/beacon-node/test/e2e/network/peers/peerManager.test.ts +++ b/packages/beacon-node/test/e2e/network/peers/peerManager.test.ts @@ -343,157 +343,4 @@ describe("network / peers / PeerManager", () => { expect(peerData?.agentClient).toBe(ClientKind.Nimbus); }); - it("Should deduplicate in-flight identify requests for the same connection", async () => { - const {libp2p, peerManager, statusCache, networkEventBus} = await mockModules(); - - let resolveIdentify!: (value: {agentVersion: string}) => void; - const identifyPromise = new Promise<{agentVersion: string}>((resolve) => { - resolveIdentify = resolve; - }); - - vi.spyOn(libp2p.services.identify, "identify").mockImplementation( - () => identifyPromise as ReturnType - ); - - const connection = { - id: "connection-1", - direction: "inbound", - status: "open", - remotePeer: peerId1, - close: async () => {}, - abort: () => {}, - } as unknown as Connection; - - getConnectionsMap(libp2p).set(peerId1.toString(), {key: peerId1, value: [connection]}); - await peerManager["onLibp2pPeerConnect"](new CustomEvent("evt", {detail: connection})); - - const remoteStatus = statusCache.get(); - networkEventBus.emit(NetworkEvent.reqRespRequest, { - request: {method: ReqRespMethod.Status, body: remoteStatus}, - peer: peerId1, - peerClient: "Unknown", - }); - networkEventBus.emit(NetworkEvent.reqRespRequest, { - request: {method: ReqRespMethod.Status, body: remoteStatus}, - peer: peerId1, - peerClient: "Unknown", - }); - - await sleep(0); - expect(libp2p.services.identify.identify).toHaveBeenCalledTimes(1); - - resolveIdentify({agentVersion: "Prysm/v6.0.0"}); - await sleep(0); - - const peerData = peerManager["connectedPeers"].get(peerId1.toString()); - expect(peerData?.agentVersion).toBe("Prysm/v6.0.0"); - expect(peerData?.agentClient).toBe(ClientKind.Prysm); - }); - - it("Should allow a new identify attempt after reconnect and ignore stale previous result", async () => { - const {libp2p, peerManager, statusCache, networkEventBus} = await mockModules(); - - let resolveFirstIdentify!: (value: {agentVersion: string}) => void; - const firstIdentifyPromise = new Promise<{agentVersion: string}>((resolve) => { - resolveFirstIdentify = resolve; - }); - - vi.spyOn(libp2p.services.identify, "identify") - .mockImplementationOnce(() => firstIdentifyPromise as ReturnType) - .mockImplementationOnce( - () => Promise.resolve({agentVersion: "Teku/v24.9.0"}) as ReturnType - ); - - const connection1 = { - id: "connection-1", - direction: "inbound", - status: "open", - remotePeer: peerId1, - close: async () => {}, - abort: () => {}, - } as unknown as Connection; - - getConnectionsMap(libp2p).set(peerId1.toString(), {key: peerId1, value: [connection1]}); - await peerManager["onLibp2pPeerConnect"](new CustomEvent("evt", {detail: connection1})); - - const remoteStatus = statusCache.get(); - networkEventBus.emit(NetworkEvent.reqRespRequest, { - request: {method: ReqRespMethod.Status, body: remoteStatus}, - peer: peerId1, - peerClient: "Unknown", - }); - await sleep(0); - - const closedConnection1 = {...connection1, status: "closed"} as Connection; - getConnectionsMap(libp2p).set(peerId1.toString(), {key: peerId1, value: [closedConnection1]}); - await peerManager["onLibp2pPeerDisconnect"](new CustomEvent("evt", {detail: closedConnection1})); - - const connection2 = { - id: "connection-2", - direction: "inbound", - status: "open", - remotePeer: peerId1, - close: async () => {}, - abort: () => {}, - } as unknown as Connection; - getConnectionsMap(libp2p).set(peerId1.toString(), {key: peerId1, value: [connection2]}); - await peerManager["onLibp2pPeerConnect"](new CustomEvent("evt", {detail: connection2})); - - networkEventBus.emit(NetworkEvent.reqRespRequest, { - request: {method: ReqRespMethod.Status, body: remoteStatus}, - peer: peerId1, - peerClient: "Unknown", - }); - await sleep(0); - - expect(libp2p.services.identify.identify).toHaveBeenCalledTimes(2); - - // Resolve old identify last; it must not overwrite new connection's identify result. - resolveFirstIdentify({agentVersion: "Lighthouse/v6.0.1"}); - await sleep(0); - - const peerData = peerManager["connectedPeers"].get(peerId1.toString()); - expect(peerData?.agentVersion).toBe("Teku/v24.9.0"); - expect(peerData?.agentClient).toBe(ClientKind.Teku); - }); - - it("Should update agentVersion via peer:identify event even if explicit identify fails", async () => { - const {libp2p, peerManager, statusCache, networkEventBus} = await mockModules(); - - vi.spyOn(libp2p.services.identify, "identify").mockRejectedValue(new Error("Unexpected EOF")); - - const connection = { - id: "connection-1", - direction: "inbound", - status: "open", - remotePeer: peerId1, - close: async () => {}, - abort: () => {}, - } as unknown as Connection; - - getConnectionsMap(libp2p).set(peerId1.toString(), {key: peerId1, value: [connection]}); - await peerManager["onLibp2pPeerConnect"](new CustomEvent("evt", {detail: connection})); - - const remoteStatus = statusCache.get(); - networkEventBus.emit(NetworkEvent.reqRespRequest, { - request: {method: ReqRespMethod.Status, body: remoteStatus}, - peer: peerId1, - peerClient: "Unknown", - }); - await sleep(0); - - libp2p.services.components.events.dispatchEvent( - new CustomEvent("peer:identify", { - detail: { - peerId: peerId1, - agentVersion: "Lighthouse/v6.0.1", - }, - }) - ); - await sleep(0); - - const peerData = peerManager["connectedPeers"].get(peerId1.toString()); - expect(peerData?.agentVersion).toBe("Lighthouse/v6.0.1"); - expect(peerData?.agentClient).toBe(ClientKind.Lighthouse); - }); }); From 97e4e97b6f46f78f6add57d960d465e46a234f3b Mon Sep 17 00:00:00 2001 From: Nico Flaig Date: Wed, 25 Feb 2026 10:26:54 +0000 Subject: [PATCH 3/4] lint --- packages/beacon-node/src/network/libp2p/index.ts | 4 ++-- .../beacon-node/test/e2e/network/peers/peerManager.test.ts | 1 - 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/packages/beacon-node/src/network/libp2p/index.ts b/packages/beacon-node/src/network/libp2p/index.ts index 12c282210732..179b48f297eb 100644 --- a/packages/beacon-node/src/network/libp2p/index.ts +++ b/packages/beacon-node/src/network/libp2p/index.ts @@ -74,7 +74,7 @@ export async function createNodeJsLibp2p( } const libp2pMetrics = nodeJsLibp2pOpts.metrics - ? ((components: LodestarComponents) => { + ? (components: LodestarComponents) => { const metrics = prometheusMetrics({ collectDefaultMetrics: false, preserveExistingMetrics: true, @@ -94,7 +94,7 @@ export async function createNodeJsLibp2p( }) as typeof metrics.trackProtocolStream; return metrics; - }) + } : undefined; return createLibp2p({ diff --git a/packages/beacon-node/test/e2e/network/peers/peerManager.test.ts b/packages/beacon-node/test/e2e/network/peers/peerManager.test.ts index edb5048803d9..deecd6e6ee25 100644 --- a/packages/beacon-node/test/e2e/network/peers/peerManager.test.ts +++ b/packages/beacon-node/test/e2e/network/peers/peerManager.test.ts @@ -342,5 +342,4 @@ describe("network / peers / PeerManager", () => { expect(peerData?.agentVersion).toBe("Nimbus/v25.0.0"); expect(peerData?.agentClient).toBe(ClientKind.Nimbus); }); - }); From 32930e33a25e7961b0829b3f0f46255ed64cb2c1 Mon Sep 17 00:00:00 2001 From: Nico Flaig Date: Wed, 25 Feb 2026 10:27:43 +0000 Subject: [PATCH 4/4] Remove unused event --- packages/beacon-node/src/constants/network.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/packages/beacon-node/src/constants/network.ts b/packages/beacon-node/src/constants/network.ts index e26bbe69ee6b..c786bf25a3bc 100644 --- a/packages/beacon-node/src/constants/network.ts +++ b/packages/beacon-node/src/constants/network.ts @@ -30,5 +30,4 @@ export const GOODBYE_KNOWN_CODES: Record = { export enum Libp2pEvent { connectionOpen = "connection:open", connectionClose = "connection:close", - peerIdentify = "peer:identify", }