Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions packages/beacon-node/src/network/core/networkCore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand All @@ -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;
}
Expand Down
13 changes: 11 additions & 2 deletions packages/beacon-node/src/network/peers/peerManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,7 @@ export class PeerManager {
private connectedPeers: Map<PeerIdStr, PeerData>;
private opts: PeerManagerOpts;
private intervals: NodeJS.Timeout[] = [];
private quiesced = false;

constructor(modules: PeerManagerModules, opts: PeerManagerOpts, discovery: PeerDiscovery | null) {
const {networkConfig} = modules;
Expand Down Expand Up @@ -228,15 +229,23 @@ export class PeerManager {
return new PeerManager(modules, opts, discovery);
}

async close(): Promise<void> {
async quiesce(): Promise<void> {
if (this.quiesced) return;

for (const interval of this.intervals) clearInterval(interval);
this.intervals = [];
await this.discovery?.stop();
this.quiesced = true;
}

async close(): Promise<void> {
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);
}

/**
Expand Down
51 changes: 51 additions & 0 deletions packages/beacon-node/test/unit/network/core/networkCore.test.ts
Original file line number Diff line number Diff line change
@@ -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<typeof NetworkCore>[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",
]);
});
});
Loading