Skip to content

Commit 1dfb74e

Browse files
feat: add WebSockets metrics (#2649)
Creates two new metric counter groups for WebSockets: `libp2p_websockets_dialer_events_total` and `libp2p_websockets_listener_events_total`. --------- Co-authored-by: achingbrain <[email protected]>
1 parent 7939dbd commit 1dfb74e

File tree

3 files changed

+90
-7
lines changed

3 files changed

+90
-7
lines changed

packages/transport-websockets/src/index.ts

Lines changed: 28 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,7 @@ import { isBrowser, isWebWorker } from 'wherearewe'
6767
import * as filters from './filters.js'
6868
import { createListener } from './listener.js'
6969
import { socketToMaConn } from './socket-to-conn.js'
70-
import type { Transport, MultiaddrFilter, CreateListenerOptions, DialTransportOptions, Listener, AbortOptions, ComponentLogger, Logger, Connection, OutboundConnectionUpgradeEvents } from '@libp2p/interface'
70+
import type { Transport, MultiaddrFilter, CreateListenerOptions, DialTransportOptions, Listener, AbortOptions, ComponentLogger, Logger, Connection, OutboundConnectionUpgradeEvents, Metrics, CounterGroup } from '@libp2p/interface'
7171
import type { Multiaddr } from '@multiformats/multiaddr'
7272
import type { Server } from 'http'
7373
import type { DuplexWebSocket } from 'it-ws/duplex'
@@ -82,6 +82,11 @@ export interface WebSocketsInit extends AbortOptions, WebSocketOptions {
8282

8383
export interface WebSocketsComponents {
8484
logger: ComponentLogger
85+
metrics?: Metrics
86+
}
87+
88+
export interface WebSocketsMetrics {
89+
dialerEvents: CounterGroup
8590
}
8691

8792
export type WebSocketsDialEvents =
@@ -92,11 +97,23 @@ class WebSockets implements Transport<WebSocketsDialEvents> {
9297
private readonly log: Logger
9398
private readonly init?: WebSocketsInit
9499
private readonly logger: ComponentLogger
100+
private readonly metrics?: WebSocketsMetrics
101+
private readonly components: WebSocketsComponents
95102

96103
constructor (components: WebSocketsComponents, init?: WebSocketsInit) {
97104
this.log = components.logger.forComponent('libp2p:websockets')
98105
this.logger = components.logger
106+
this.components = components
99107
this.init = init
108+
109+
if (components.metrics != null) {
110+
this.metrics = {
111+
dialerEvents: components.metrics.registerCounterGroup('libp2p_websockets_dialer_events_total', {
112+
label: 'event',
113+
help: 'Total count of WebSockets dialer events by type'
114+
})
115+
}
116+
}
100117
}
101118

102119
readonly [transportSymbol] = true
@@ -113,7 +130,8 @@ class WebSockets implements Transport<WebSocketsDialEvents> {
113130

114131
const socket = await this._connect(ma, options)
115132
const maConn = socketToMaConn(socket, ma, {
116-
logger: this.logger
133+
logger: this.logger,
134+
metrics: this.metrics?.dialerEvents
117135
})
118136
this.log('new outbound connection %s', maConn.remoteAddr)
119137

@@ -136,13 +154,18 @@ class WebSockets implements Transport<WebSocketsDialEvents> {
136154
// https://developer.mozilla.org/en-US/docs/Web/API/WebSocket/error_event
137155
const err = new CodeError(`Could not connect to ${ma.toString()}`, 'ERR_CONNECTION_FAILED')
138156
this.log.error('connection error:', err)
157+
this.metrics?.dialerEvents.increment({ error: true })
139158
errorPromise.reject(err)
140159
})
141160

142161
try {
143162
options.onProgress?.(new CustomProgressEvent('websockets:open-connection'))
144163
await raceSignal(Promise.race([rawSocket.connected(), errorPromise.promise]), options.signal)
145164
} catch (err: any) {
165+
if (options.signal?.aborted === true) {
166+
this.metrics?.dialerEvents.increment({ abort: true })
167+
}
168+
146169
rawSocket.close()
147170
.catch(err => {
148171
this.log.error('error closing raw socket', err)
@@ -152,6 +175,7 @@ class WebSockets implements Transport<WebSocketsDialEvents> {
152175
}
153176

154177
this.log('connected %s', ma)
178+
this.metrics?.dialerEvents.increment({ connect: true })
155179
return rawSocket
156180
}
157181

@@ -162,7 +186,8 @@ class WebSockets implements Transport<WebSocketsDialEvents> {
162186
*/
163187
createListener (options: CreateListenerOptions): Listener {
164188
return createListener({
165-
logger: this.logger
189+
logger: this.logger,
190+
metrics: this.components.metrics
166191
}, {
167192
...this.init,
168193
...options

packages/transport-websockets/src/listener.ts

Lines changed: 49 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,43 +1,57 @@
11
import os from 'os'
2-
import { TypedEventEmitter, CustomEvent } from '@libp2p/interface'
2+
import { TypedEventEmitter } from '@libp2p/interface'
33
import { ipPortToMultiaddr as toMultiaddr } from '@libp2p/utils/ip-port-to-multiaddr'
44
import { multiaddr, protocols } from '@multiformats/multiaddr'
55
import { createServer } from 'it-ws/server'
66
import { socketToMaConn } from './socket-to-conn.js'
7-
import type { ComponentLogger, Logger, Connection, Listener, ListenerEvents, CreateListenerOptions } from '@libp2p/interface'
7+
import type { ComponentLogger, Logger, Connection, Listener, ListenerEvents, CreateListenerOptions, CounterGroup, MetricGroup, Metrics } from '@libp2p/interface'
88
import type { Multiaddr } from '@multiformats/multiaddr'
99
import type { Server } from 'http'
1010
import type { DuplexWebSocket } from 'it-ws/duplex'
1111
import type { WebSocketServer } from 'it-ws/server'
1212

1313
export interface WebSocketListenerComponents {
1414
logger: ComponentLogger
15+
metrics?: Metrics
1516
}
1617

1718
export interface WebSocketListenerInit extends CreateListenerOptions {
1819
server?: Server
1920
}
2021

22+
export interface WebSocketListenerMetrics {
23+
status: MetricGroup
24+
errors: CounterGroup
25+
events: CounterGroup
26+
}
27+
2128
class WebSocketListener extends TypedEventEmitter<ListenerEvents> implements Listener {
2229
private readonly connections: Set<DuplexWebSocket>
2330
private listeningMultiaddr?: Multiaddr
2431
private readonly server: WebSocketServer
2532
private readonly log: Logger
33+
private metrics?: WebSocketListenerMetrics
34+
private addr: string
2635

2736
constructor (components: WebSocketListenerComponents, init: WebSocketListenerInit) {
2837
super()
2938

3039
this.log = components.logger.forComponent('libp2p:websockets:listener')
40+
const metrics = components.metrics
3141
// Keep track of open connections to destroy when the listener is closed
3242
this.connections = new Set<DuplexWebSocket>()
3343

3444
const self = this // eslint-disable-line @typescript-eslint/no-this-alias
3545

46+
this.addr = 'unknown'
47+
3648
this.server = createServer({
3749
...init,
3850
onConnection: (stream: DuplexWebSocket) => {
3951
const maConn = socketToMaConn(stream, toMultiaddr(stream.remoteAddress ?? '', stream.remotePort ?? 0), {
40-
logger: components.logger
52+
logger: components.logger,
53+
metrics: this.metrics?.events,
54+
metricPrefix: `${this.addr} `
4155
})
4256
this.log('new inbound connection %s', maConn.remoteAddr)
4357

@@ -62,6 +76,7 @@ class WebSocketListener extends TypedEventEmitter<ListenerEvents> implements Lis
6276
})
6377
.catch(async err => {
6478
this.log.error('inbound connection failed to upgrade', err)
79+
this.metrics?.errors.increment({ [`${this.addr} inbound_upgrade`]: true })
6580

6681
await maConn.close().catch(err => {
6782
this.log.error('inbound connection failed to close after upgrade failed', err)
@@ -71,15 +86,46 @@ class WebSocketListener extends TypedEventEmitter<ListenerEvents> implements Lis
7186
this.log.error('inbound connection failed to upgrade', err)
7287
maConn.close().catch(err => {
7388
this.log.error('inbound connection failed to close after upgrade failed', err)
89+
this.metrics?.errors.increment({ [`${this.addr} inbound_closing_failed`]: true })
7490
})
7591
}
7692
}
7793
})
7894

7995
this.server.on('listening', () => {
96+
if (metrics != null) {
97+
const { host, port } = this.listeningMultiaddr?.toOptions() ?? {}
98+
this.addr = `${host}:${port}`
99+
100+
metrics.registerMetricGroup('libp2p_websockets_inbound_connections_total', {
101+
label: 'address',
102+
help: 'Current active connections in WebSocket listener',
103+
calculate: () => {
104+
return {
105+
[this.addr]: this.connections.size
106+
}
107+
}
108+
})
109+
110+
this.metrics = {
111+
status: metrics?.registerMetricGroup('libp2p_websockets_listener_status_info', {
112+
label: 'address',
113+
help: 'Current status of the WebSocket listener socket'
114+
}),
115+
errors: metrics?.registerMetricGroup('libp2p_websockets_listener_errors_total', {
116+
label: 'address',
117+
help: 'Total count of WebSocket listener errors by type'
118+
}),
119+
events: metrics?.registerMetricGroup('libp2p_websockets_listener_events_total', {
120+
label: 'address',
121+
help: 'Total count of WebSocket listener events by type'
122+
})
123+
}
124+
}
80125
this.dispatchEvent(new CustomEvent('listening'))
81126
})
82127
this.server.on('error', (err: Error) => {
128+
this.metrics?.errors.increment({ [`${this.addr} listen_error`]: true })
83129
this.dispatchEvent(new CustomEvent('error', {
84130
detail: err
85131
}))

packages/transport-websockets/src/socket-to-conn.ts

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,22 @@
11
import { CodeError } from '@libp2p/interface'
22
import { CLOSE_TIMEOUT } from './constants.js'
3-
import type { AbortOptions, ComponentLogger, MultiaddrConnection } from '@libp2p/interface'
3+
import type { AbortOptions, ComponentLogger, CounterGroup, MultiaddrConnection } from '@libp2p/interface'
44
import type { Multiaddr } from '@multiformats/multiaddr'
55
import type { DuplexWebSocket } from 'it-ws/duplex'
66

77
export interface SocketToConnOptions {
88
localAddr?: Multiaddr
99
logger: ComponentLogger
10+
metrics?: CounterGroup
11+
metricPrefix?: string
1012
}
1113

1214
// Convert a stream into a MultiaddrConnection
1315
// https://github.com/libp2p/interface-transport#multiaddrconnection
1416
export function socketToMaConn (stream: DuplexWebSocket, remoteAddr: Multiaddr, options: SocketToConnOptions): MultiaddrConnection {
1517
const log = options.logger.forComponent('libp2p:websockets:maconn')
18+
const metrics = options.metrics
19+
const metricPrefix = options.metricPrefix ?? ''
1620

1721
const maConn: MultiaddrConnection = {
1822
log,
@@ -81,10 +85,18 @@ export function socketToMaConn (stream: DuplexWebSocket, remoteAddr: Multiaddr,
8185

8286
stream.destroy()
8387
maConn.timeline.close = Date.now()
88+
89+
// ws WebSocket.terminate does not accept an Error arg to emit an 'error'
90+
// event on destroy like other node streams so we can't update a metric
91+
// with an event listener
92+
// https://github.com/websockets/ws/issues/1752#issuecomment-622380981
93+
metrics?.increment({ [`${metricPrefix}error`]: true })
8494
}
8595
}
8696

8797
stream.socket.addEventListener('close', () => {
98+
metrics?.increment({ [`${metricPrefix}close`]: true })
99+
88100
// In instances where `close` was not explicitly called,
89101
// such as an iterable stream ending, ensure we have set the close
90102
// timeline

0 commit comments

Comments
 (0)