Skip to content
Merged
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
116 changes: 80 additions & 36 deletions packages/datadog-plugin-ws/test/index.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,8 @@ const assert = require('node:assert')
const { once } = require('node:events')

const dc = require('dc-polyfill')
const { after, afterEach, before, beforeEach, describe, it } = require('mocha')
const setSocketCh = dc.channel('tracing:ws:server:connect:setSocket')

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking, but with DataDog/dc-polyfill#27 this isn't really needed anymore

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That is indeed redundant from a fix perspective, but it's still good practice to call channel only once at the top, so no harm in keeping it like that.

const { afterEach, beforeEach, describe, it } = require('mocha')

const agent = require('../../dd-trace/test/plugins/agent')
const { storage } = require('../../datadog-core')
Expand All @@ -26,6 +27,7 @@ function findSpan (traces, predicate) {
}

function closeWsServer (server) {
if (!server) return
for (const ws of server.clients) {
ws.terminate()
}
Expand All @@ -36,58 +38,72 @@ describe('Plugin', () => {
let WebSocket
let wsServer
let connectionReceived
let clientPort = 6015
let clientPort
let client
let messageReceived
let route

describe('ws', () => {
withVersions('ws', 'ws', '>=8.0.0', version => {
describe('regression tests', () => {
before(async () => {
let regressionServer
let regressionSocket

beforeEach(async () => {
await agent.load(['ws'], [{
service: 'some',
traceWebsocketMessagesEnabled: true,
}])
WebSocket = require(`../../../versions/ws@${version}`).get()
})

afterEach(() => {
regressionSocket?.terminate()
regressionSocket = undefined
})

afterEach(async () => {
await closeWsServer(regressionServer)
})

afterEach(() => {
regressionServer = undefined
})

afterEach(async () => {
await agent.close({ ritmReset: false, wipe: true })
})

it('should not crash when sending on a socket without spanContext', async () => {
const server = new WebSocket.Server({ port: 16015 })
const connectionPromise = once(server, 'connection')
regressionServer = new WebSocket.Server({ port: 0 })
await once(regressionServer, 'listening')
const { port } = regressionServer.address()
const connectionPromise = once(regressionServer, 'connection')

const socket = new WebSocket('ws://localhost:16015')
regressionSocket = new WebSocket(`ws://localhost:${port}`)
const [serverSocket] = await connectionPromise
await once(socket, 'open')
await once(regressionSocket, 'open')

assert.strictEqual(socket.spanContext, undefined)
assert.strictEqual(regressionSocket.spanContext, undefined)

const messagePromise = once(serverSocket, 'message')
await new Promise((resolve, reject) => {
socket.send('test message', {}, (err) => err ? reject(err) : resolve())
regressionSocket.send('test message', {}, (err) => err ? reject(err) : resolve())
})
await messagePromise

socket.close()
await once(socket, 'close')
server.close()
})

it('should emit original error in case close is called before connection is established', async () => {
const socket = new WebSocket('wss://localhost:12345')
regressionSocket = new WebSocket('wss://localhost:12345')

const errorPromise = once(socket, 'error')
socket.close()
const errorPromise = once(regressionSocket, 'error')
regressionSocket.close()

const error = await errorPromise

// Some versions emit an array with an error, some directly emit the error
assert.strictEqual(error?.[0]?.message, 'WebSocket was closed before the connection was established')
})

after(async () => {
await agent.close({ ritmReset: false, wipe: true })
})
})

describe('when using WebSocket', () => {
Expand All @@ -106,22 +122,27 @@ describe('Plugin', () => {
}])
WebSocket = require(`../../../versions/ws@${version}`).get()

wsServer = new WebSocket.Server({ port: clientPort })
wsServer = new WebSocket.Server({ port: 0 })
await once(wsServer, 'listening')
clientPort = wsServer.address().port
})

afterEach(async () => {
clientPort++
afterEach(() => {
if (client) {
client.removeAllListeners('error')
client.on('error', () => {})
}
})

afterEach(async () => {
await closeWsServer(wsServer)
})

afterEach(async () => {
await agent.close({ ritmReset: false, wipe: true })
})

it('should not retain the connection span during socket setup', async () => {
const setSocketCh = dc.channel('tracing:ws:server:connect:setSocket')
let resolve
const promise = new Promise((_resolve) => {
resolve = _resolve
Expand Down Expand Up @@ -464,17 +485,23 @@ describe('Plugin', () => {
}])
WebSocket = require(`../../../versions/ws@${version}`).get()

wsServer = new WebSocket.Server({ port: clientPort })
wsServer = new WebSocket.Server({ port: 0 })
await once(wsServer, 'listening')
clientPort = wsServer.address().port
})

afterEach(async () => {
clientPort++
afterEach(() => {
if (client) {
client.removeAllListeners('error')
client.on('error', () => {})
}
})

afterEach(async () => {
await closeWsServer(wsServer)
})

afterEach(async () => {
await agent.close({ ritmReset: false, wipe: true })
})

Expand Down Expand Up @@ -576,17 +603,22 @@ describe('Plugin', () => {
}])
WebSocket = require(`../../../versions/ws@${version}`).get()

wsServer = new WebSocket.Server({ port: clientPort })
wsServer = new WebSocket.Server({ port: 0 })
await once(wsServer, 'listening')
})

afterEach(async () => {
clientPort++
afterEach(() => {
if (client) {
client.removeAllListeners('error')
client.on('error', () => {})
}
})

afterEach(async () => {
await closeWsServer(wsServer)
})

afterEach(async () => {
await agent.close({ ritmReset: false, wipe: true })
})

Expand Down Expand Up @@ -625,17 +657,23 @@ describe('Plugin', () => {
}])
WebSocket = require(`../../../versions/ws@${version}`).get()

wsServer = new WebSocket.Server({ port: clientPort })
wsServer = new WebSocket.Server({ port: 0 })
await once(wsServer, 'listening')
clientPort = wsServer.address().port
})

afterEach(async () => {
clientPort++
afterEach(() => {
if (client) {
client.removeAllListeners('error')
client.on('error', () => {})
}
})

afterEach(async () => {
await closeWsServer(wsServer)
})

afterEach(async () => {
await agent.close({ ritmReset: false, wipe: true })
})

Expand Down Expand Up @@ -721,22 +759,28 @@ describe('Plugin', () => {
}])
WebSocket = require(`../../../versions/ws@${version}`).get()

wsServer = new WebSocket.Server({ port: clientPort })
wsServer = new WebSocket.Server({ port: 0 })
await once(wsServer, 'listening')
clientPort = wsServer.address().port

parentHeaders = {}
tracer.trace('test.parent', parentSpan => {
tracer.inject(parentSpan, 'http_headers', parentHeaders)
})
})

afterEach(async () => {
clientPort++
afterEach(() => {
if (client) {
client.removeAllListeners('error')
client.on('error', () => {})
}
})

afterEach(async () => {
await closeWsServer(wsServer)
})

afterEach(async () => {
await agent.close({ ritmReset: false, wipe: true })
})

Expand Down
Loading