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
11 changes: 9 additions & 2 deletions packages/libp2p/src/connection-manager/dial-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ export class DialQueue {
private readonly connections: PeerMap<Connection[]>
private readonly log: Logger
private readonly resolvers: Record<string, MultiaddrResolver>
private readonly activeDials: Set<Promise<Connection>>

constructor (components: DialQueueComponents, init: DialerInit = {}) {
this.addressSorter = init.addressSorter
Expand All @@ -91,6 +92,7 @@ export class DialQueue {
this.log = components.logger.forComponent('libp2p:connection-manager:dial-queue')
this.components = components
this.resolvers = init.resolvers ?? defaultOptions.resolvers
this.activeDials = new Set()

this.shutDownController = new AbortController()
setMaxListeners(Infinity, this.shutDownController.signal)
Expand All @@ -117,9 +119,10 @@ export class DialQueue {
/**
* Clears any pending dials
*/
stop (): void {
async stop (): Promise<void> {
this.shutDownController.abort()
this.queue.abort()
await Promise.allSettled(this.activeDials)
}

/**
Expand Down Expand Up @@ -199,9 +202,13 @@ export class DialQueue {
])
setMaxListeners(Infinity, signal)

const dialPromise = this.dialPeer(options, signal)
this.activeDials.add(dialPromise)

try {
return await this.dialPeer(options, signal)
return await dialPromise
} finally {
this.activeDials.delete(dialPromise)
// clean up abort signals/controllers
signal.clear()
}
Expand Down
10 changes: 8 additions & 2 deletions packages/libp2p/src/connection-manager/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -376,12 +376,15 @@ export class DefaultConnectionManager implements ConnectionManager, Startable {
this.log('started')
}

beforeStop (): void {
this.started = false
}

/**
* Stops the Connection Manager
*/
async stop (): Promise<void> {
this.events.removeEventListener('connection:open', this.onConnect)
this.events.removeEventListener('connection:close', this.onDisconnect)
this.started = false

await stop(
this.reconnectQueue,
Expand Down Expand Up @@ -413,6 +416,9 @@ export class DefaultConnectionManager implements ConnectionManager, Startable {
await Promise.all(tasks)
this.connections.clear()

this.events.removeEventListener('connection:open', this.onConnect)
this.events.removeEventListener('connection:close', this.onDisconnect)

this.log('stopped')
}

Expand Down
39 changes: 37 additions & 2 deletions packages/libp2p/test/connection-manager/dial-queue.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,12 +39,47 @@ describe('dial queue', () => {
}
})

afterEach(() => {
afterEach(async () => {
if (dialer != null) {
dialer.stop()
await dialer.stop()
}
})

it('should wait for active dials to finish when stopped', async () => {
const dialStarted = pDefer<void>()
const dialCleanup = pDefer<void>()

components.transportManager.dialTransportForMultiaddr.returns(stubInterface<Transport>())
components.transportManager.dial.callsFake(async (_ma, options) => {
dialStarted.resolve()

await new Promise<void>(resolve => {
options?.signal?.addEventListener('abort', () => { resolve() }, { once: true })
})
await dialCleanup.promise

throw new Error('dial aborted')
})

dialer = new DialQueue(components)
const dialPromise = dialer.dial(multiaddr('/ip4/127.0.0.1/tcp/1234'))
const dialResult = expect(dialPromise).to.eventually.be.rejected()

await dialStarted.promise

let stopped = false
const stopPromise = dialer.stop().then(() => {
stopped = true
})

await Promise.resolve()
expect(stopped).to.be.false()

dialCleanup.resolve()
await stopPromise
await dialResult
})

it('should end when a single multiaddr dials succeeds', async () => {
const connection = stubInterface<Connection>()
const deferredConn = pDefer<Connection>()
Expand Down
18 changes: 18 additions & 0 deletions packages/libp2p/test/connection-manager/index.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -525,4 +525,22 @@ describe('Connection Manager', () => {

expect(connectionManager.getConnections(remotePeer)).to.have.lengthOf(0)
})

it('should close connections opened after shutdown starts', async () => {
const connectionManager = new DefaultConnectionManager(components)
await connectionManager.start()

const connection = stubInterface<Connection>({
remotePeer: peerIdFromPrivateKey(await generateKeyPair('Ed25519')),
status: 'open'
})

connectionManager.beforeStop()
components.events.safeDispatchEvent('connection:open', { detail: connection })

await pWaitFor(() => connection.close.called)
expect(connectionManager.getConnections()).to.have.lengthOf(0)

await connectionManager.stop()
})
})
11 changes: 8 additions & 3 deletions packages/transport-tcp/src/listener.ts
Original file line number Diff line number Diff line change
Expand Up @@ -217,8 +217,10 @@ export class TCPListener extends TypedEventEmitter<ListenerEvents> implements Li
.catch(async err => {
this.log.error('inbound connection upgrade failed - %e', err)
this.metrics.errors?.increment({ [`${this.addr} inbound_upgrade`]: true })
this.sockets.delete(socket)
const socketClosed = socket.closed ? undefined : pEvent(socket, 'close', { rejectionEvents: [] })
maConn.abort(err)
await socketClosed
this.sockets.delete(socket)
})
}

Expand Down Expand Up @@ -281,9 +283,12 @@ export class TCPListener extends TypedEventEmitter<ListenerEvents> implements Li
// synchronously close any open connections - should be done after closing
// the server socket in case new sockets are opened during the shutdown
this.sockets.forEach(socket => {
if (socket.readable) {
if (!socket.closed) {
events.push(pEvent(socket, 'close', options))
socket.destroy()

if (!socket.destroyed) {
socket.destroy()
}
}
})

Expand Down
6 changes: 6 additions & 0 deletions packages/transport-tcp/src/socket-to-conn.ts
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,12 @@ class TCPSocketMultiaddrConnection extends AbstractMultiaddrConnection {
}

sendReset (): void {
// Node.js cannot reset a socket while a graceful shutdown request is pending
if (this.socket.writableEnded) {
this.socket.destroy()
return
}

this.socket.resetAndDestroy()
}

Expand Down
44 changes: 40 additions & 4 deletions packages/transport-tcp/src/tcp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,25 +30,28 @@
import net from 'net'
import { AbortError, TimeoutError, serviceCapabilities, transportSymbol } from '@libp2p/interface'
import { TCP as TCPMatcher } from '@multiformats/multiaddr-matcher'
import { pEvent } from 'p-event'
import { CustomProgressEvent } from 'progress-events'
import { TCPListener } from './listener.js'
import { toMultiaddrConnection } from './socket-to-conn.js'
import { multiaddrToNetConfig } from './utils.js'
import type { TCPComponents, TCPCreateListenerOptions, TCPDialEvents, TCPDialOptions, TCPMetrics, TCPOptions } from './index.js'
import type { Logger, Connection, Transport, Listener, MultiaddrConnection } from '@libp2p/interface'
import type { Logger, Connection, Transport, Listener, MultiaddrConnection, Startable } from '@libp2p/interface'
import type { Multiaddr } from '@multiformats/multiaddr'
import type { Socket, IpcSocketConnectOpts, TcpSocketConnectOpts } from 'net'

export class TCP implements Transport<TCPDialEvents> {
export class TCP implements Transport<TCPDialEvents>, Startable {
private readonly opts: TCPOptions
private readonly metrics?: TCPMetrics
private readonly components: TCPComponents
private readonly log: Logger
private readonly outboundSockets: Set<Socket>

constructor (components: TCPComponents, options: TCPOptions = {}) {
this.log = components.logger.forComponent('libp2p:tcp')
this.opts = options
this.components = components
this.outboundSockets = new Set()

if (components.metrics != null) {
this.metrics = {
Expand All @@ -72,6 +75,24 @@ export class TCP implements Transport<TCPDialEvents> {
'@libp2p/transport'
]

start (): void {}

stop (): void {}

async afterStop (): Promise<void> {
const socketClosedPromises: Array<Promise<unknown>> = []

for (const socket of this.outboundSockets) {
socketClosedPromises.push(pEvent(socket, 'close', { rejectionEvents: [] }))

if (!socket.destroyed) {
socket.destroy()
}
}

await Promise.all(socketClosedPromises)
}

async dial (ma: Multiaddr, options: TCPDialOptions): Promise<Connection> {
options.keepAlive = options.keepAlive ?? true
options.noDelay = options.noDelay ?? true
Expand All @@ -93,7 +114,9 @@ export class TCP implements Transport<TCPDialEvents> {
})
} catch (err: any) {
this.metrics?.errors.increment({ outbound_to_connection: true })
const socketClosed = socket.closed ? undefined : pEvent(socket, 'close', { rejectionEvents: [] })
socket.destroy(err)
await socketClosed
throw err
}

Expand All @@ -103,7 +126,9 @@ export class TCP implements Transport<TCPDialEvents> {
} catch (err: any) {
this.metrics?.errors.increment({ outbound_upgrade: true })
this.log.error('error upgrading outbound connection - %e', err)
const socketClosed = socket.closed ? undefined : pEvent(socket, 'close', { rejectionEvents: [] })
maConn.abort(err)
await socketClosed
throw err
}
}
Expand All @@ -123,6 +148,10 @@ export class TCP implements Transport<TCPDialEvents> {

this.log('dialing %a with opts %o', ma, cOpts)
rawSocket = net.connect(cOpts)
this.outboundSockets.add(rawSocket)
rawSocket.once('close', () => {
this.outboundSockets.delete(rawSocket)
})

const onError = (err: Error): void => {
this.log.error('dial to %a errored - %e', ma, err)
Expand Down Expand Up @@ -175,8 +204,15 @@ export class TCP implements Transport<TCPDialEvents> {

options.signal.addEventListener('abort', onAbort)
})
.catch(err => {
rawSocket?.destroy()
.catch(async err => {
if (rawSocket != null && !rawSocket.closed) {
const socketClosed = new Promise<void>(resolve => {
rawSocket.once('close', resolve)
})
rawSocket.destroy()
await socketClosed
}

throw err
})
}
Expand Down
Loading