diff --git a/packages/beacon-node/src/network/core/networkCore.ts b/packages/beacon-node/src/network/core/networkCore.ts index c2ed2159bc6a..c100bd50128e 100644 --- a/packages/beacon-node/src/network/core/networkCore.ts +++ b/packages/beacon-node/src/network/core/networkCore.ts @@ -285,6 +285,9 @@ export class NetworkCore implements INetworkCore { this.clock.off(ClockEvent.epoch, this.onEpoch); + await this.peerManager.quiesce(); + this.logger.debug("network peerManager quiesced"); + // Must goodbye and disconnect before stopping libp2p await this.peerManager.goodbyeAndDisconnectAllPeers(); this.logger.debug("network sent goodbye to all peers"); @@ -302,6 +305,10 @@ export class NetworkCore implements INetworkCore { // a handle that never closes and `Worker.terminate()` then never resolves. Diffing this list // between a clean and a stuck shutdown should narrow down which handle it is this.logger.debug("network lib2p closed", {activeResources: formatActiveResources()}); + const processReport = process.report.getReport() as {libuv: Array<{type: string}>}; + this.logger.debug("network lib2p tcp handles", { + handles: JSON.stringify(processReport.libuv.filter((handle) => handle.type === "tcp")), + }); this.closed = true; } diff --git a/packages/beacon-node/src/network/peers/peerManager.ts b/packages/beacon-node/src/network/peers/peerManager.ts index 928a5b6e643a..49b1148511fb 100644 --- a/packages/beacon-node/src/network/peers/peerManager.ts +++ b/packages/beacon-node/src/network/peers/peerManager.ts @@ -164,6 +164,7 @@ export class PeerManager { private connectedPeers: Map; private opts: PeerManagerOpts; private intervals: NodeJS.Timeout[] = []; + private quiesced = false; constructor(modules: PeerManagerModules, opts: PeerManagerOpts, discovery: PeerDiscovery | null) { const {networkConfig} = modules; @@ -228,15 +229,23 @@ export class PeerManager { return new PeerManager(modules, opts, discovery); } - async close(): Promise { + async quiesce(): Promise { + if (this.quiesced) return; + + for (const interval of this.intervals) clearInterval(interval); + this.intervals = []; await this.discovery?.stop(); + this.quiesced = true; + } + + async close(): Promise { + await this.quiesce(); this.libp2p.services.components.events.removeEventListener(Libp2pEvent.connectionOpen, this.onLibp2pPeerConnect); this.libp2p.services.components.events.removeEventListener( Libp2pEvent.connectionClose, this.onLibp2pPeerDisconnect ); this.networkEventBus.off(NetworkEvent.reqRespRequest, this.onRequest); - for (const interval of this.intervals) clearInterval(interval); } /** diff --git a/packages/beacon-node/test/unit/network/core/networkCore.test.ts b/packages/beacon-node/test/unit/network/core/networkCore.test.ts new file mode 100644 index 000000000000..1af877f3e9f6 --- /dev/null +++ b/packages/beacon-node/test/unit/network/core/networkCore.test.ts @@ -0,0 +1,51 @@ +import {describe, expect, it, vi} from "vitest"; +import {NetworkCore} from "../../../../src/network/core/networkCore.js"; + +describe("network / core / NetworkCore", () => { + it("quiesces peer management before disconnecting peers", async () => { + const calls: string[] = []; + const modules = { + libp2p: {stop: vi.fn(async () => calls.push("libp2p.stop"))}, + gossip: {stop: vi.fn(async () => calls.push("gossip.stop"))}, + reqResp: { + stop: vi.fn(async () => calls.push("reqResp.stop")), + unregisterAllProtocols: vi.fn(async () => calls.push("reqResp.unregister")), + }, + attnetsService: {close: vi.fn(() => calls.push("attnets.close"))}, + syncnetsService: {close: vi.fn(() => calls.push("syncnets.close"))}, + peerManager: { + quiesce: vi.fn(async () => calls.push("peerManager.quiesce")), + goodbyeAndDisconnectAllPeers: vi.fn(async () => calls.push("peerManager.goodbye")), + close: vi.fn(async () => calls.push("peerManager.close")), + }, + networkConfig: {}, + peersData: {}, + metadata: {}, + logger: {debug: vi.fn()}, + config: {}, + clock: { + on: vi.fn(), + off: vi.fn(() => calls.push("clock.off")), + }, + statusCache: {}, + metrics: null, + opts: {}, + } as unknown as ConstructorParameters[0]; + + const network = new NetworkCore(modules); + await network.close(); + + expect(calls).toEqual([ + "clock.off", + "peerManager.quiesce", + "peerManager.goodbye", + "peerManager.close", + "gossip.stop", + "reqResp.stop", + "reqResp.unregister", + "attnets.close", + "syncnets.close", + "libp2p.stop", + ]); + }); +});