From 433b809cc76178684a5e2d0f1ae90528f2f96bab Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Mon, 10 Aug 2026 19:48:06 +0200 Subject: [PATCH 1/7] =?UTF-8?q?=E2=9C=85=20Add=20a=20shared=20mockWebSocke?= =?UTF-8?q?t=20test=20utility?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Lift the fake WebSocket inlined in the WebSocket observable spec into a shared core test helper, so upcoming WebSocket tickets drive instrumentation through one test double instead of growing their own copies. The utility follows the XHR mock's shape: it swaps the `WebSocket` global and registers its own cleanup, including resetting the observable singleton. The observable spec's test cases and assertions are unchanged. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/browser/webSocketObservable.spec.ts | 121 ++++-------------- .../test/emulate/mockWebSocket.ts | 90 +++++++++++++ packages/browser-core/test/index.ts | 1 + 3 files changed, 115 insertions(+), 97 deletions(-) create mode 100644 packages/browser-core/test/emulate/mockWebSocket.ts diff --git a/packages/browser-core/src/browser/webSocketObservable.spec.ts b/packages/browser-core/src/browser/webSocketObservable.spec.ts index 05ef1ccf96..019cc90e6b 100644 --- a/packages/browser-core/src/browser/webSocketObservable.spec.ts +++ b/packages/browser-core/src/browser/webSocketObservable.spec.ts @@ -1,94 +1,21 @@ -import { registerCleanupTask } from '../../test' +import { createMockWebSocket, mockWebSocket, MockWebSocket, registerCleanupTask } from '../../test' import type { Subscription } from '../tools/observable' import { setAllowUntrustedEvents } from './addEventListener' import type { WebSocketContext } from './webSocketObservable' import { initWebSocketObservable, resetWebSocketObservable } from './webSocketObservable' -// A minimal stand-in for the native `WebSocket` constructor. We do not connect to a real server in -// unit tests; instead we expose helpers to simulate the browser dispatching events on the instance. -class FakeWebSocket extends EventTarget { - static readonly CONNECTING = 0 - static readonly OPEN = 1 - static readonly CLOSING = 2 - static readonly CLOSED = 3 - - url: string - protocol = '' - bufferedAmount = 0 - readyState: number = FakeWebSocket.CONNECTING - onmessage: ((event: MessageEvent) => void) | null = null - onopen: ((event: Event) => void) | null = null - onclose: ((event: CloseEvent) => void) | null = null - - constructor(url: string | URL, protocols?: string | string[]) { - super() - this.url = resolveWebSocketUrl(String(url)) - if (typeof protocols === 'string') { - this.protocol = protocols - } - } - - send(_data: string | ArrayBufferLike | Blob | ArrayBufferView): void { - // no-op; tests will set `bufferedAmount` before calling send to verify it is sampled. - } - - close(_code?: number, _reason?: string): void { - this.readyState = FakeWebSocket.CLOSED - } - - simulateOpen() { - this.readyState = FakeWebSocket.OPEN - const event = new Event('open') - this.dispatchEvent(event) - this.onopen?.(event) - } - - simulateMessage(data: unknown) { - const event = new MessageEvent('message', { data }) - this.dispatchEvent(event) - this.onmessage?.(event) - } - - simulateClose(code: number, reason: string, wasClean: boolean) { - this.readyState = FakeWebSocket.CLOSED - // CloseEvent is not always constructable in test environments; use a plain Event with assigned fields. - const event = Object.assign(new Event('close'), { code, reason, wasClean }) as CloseEvent - this.dispatchEvent(event) - this.onclose?.(event) - } -} - -// Mimics how a real browser resolves the URL passed to the `WebSocket` constructor: relative URLs -// are resolved against the document location, and `http(s)` schemes are translated to `ws(s)`. -function resolveWebSocketUrl(url: string): string { - const resolved = new URL(url, location.href) - if (resolved.protocol === 'http:') { - resolved.protocol = 'ws:' - } else if (resolved.protocol === 'https:') { - resolved.protocol = 'wss:' - } - return resolved.href -} - -type FakeWebSocketConstructor = typeof FakeWebSocket - -const windowAsWebSocketHost = window as unknown as { WebSocket: FakeWebSocketConstructor } - describe('webSocketObservable', () => { - let originalWebSocket: FakeWebSocketConstructor let contexts: WebSocketContext[] let subscription: Subscription | undefined beforeEach(() => { - originalWebSocket = windowAsWebSocketHost.WebSocket - windowAsWebSocketHost.WebSocket = FakeWebSocket + mockWebSocket() contexts = [] registerCleanupTask(() => { subscription?.unsubscribe() subscription = undefined resetWebSocketObservable() - windowAsWebSocketHost.WebSocket = originalWebSocket }) }) @@ -111,7 +38,7 @@ describe('webSocketObservable', () => { describe('connecting context', () => { it('emits a "connecting" context when a WebSocket is constructed', () => { const url = 'wss://example.com/socket' - const ws = new windowAsWebSocketHost.WebSocket(url) + const ws = createMockWebSocket(url) const connectingContexts = getContexts('connecting') expect(connectingContexts.length).toBe(1) @@ -121,7 +48,7 @@ describe('webSocketObservable', () => { }) it('reports the resolved instance.url rather than the raw constructor argument', () => { - const ws = new windowAsWebSocketHost.WebSocket('/socket') + const ws = createMockWebSocket('/socket') const connectingContext = getContexts('connecting')[0] expect(connectingContext.url).not.toBe('/socket') @@ -129,7 +56,7 @@ describe('webSocketObservable', () => { }) it('does not include protocols in the "connecting" context when omitted', () => { - new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + createMockWebSocket('wss://example.com/socket') expect(getContexts('connecting')[0].protocols).toBeUndefined() }) @@ -137,7 +64,7 @@ describe('webSocketObservable', () => { it('includes string protocols in the "connecting" context', () => { const url = 'wss://example.com/socket' const protocols = 'chat.v1' - new windowAsWebSocketHost.WebSocket(url, protocols) + createMockWebSocket(url, protocols) expect(getContexts('connecting')[0].protocols).toBe(protocols) }) @@ -145,7 +72,7 @@ describe('webSocketObservable', () => { it('includes array protocols in the "connecting" context', () => { const url = 'wss://example.com/socket' const protocols = ['chat.v1', 'json'] - new windowAsWebSocketHost.WebSocket(url, protocols) + createMockWebSocket(url, protocols) expect(getContexts('connecting')[0].protocols).toEqual(protocols) }) @@ -153,7 +80,7 @@ describe('webSocketObservable', () => { describe('preservation of native behavior', () => { it('does not clobber a customer-set onmessage handler', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const customerHandler = jasmine.createSpy() ws.onmessage = customerHandler @@ -164,7 +91,7 @@ describe('webSocketObservable', () => { }) it('does not clobber a customer-set onopen handler', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const customerHandler = jasmine.createSpy() ws.onopen = customerHandler @@ -175,7 +102,7 @@ describe('webSocketObservable', () => { }) it('does not clobber a customer-set onclose handler', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const customerHandler = jasmine.createSpy() ws.onclose = customerHandler @@ -188,7 +115,7 @@ describe('webSocketObservable', () => { describe('open context', () => { it('emits an "open" context when the WebSocket opens', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const negotiatedProtocol = 'chat.v1' ws.protocol = negotiatedProtocol ws.simulateOpen() @@ -201,7 +128,7 @@ describe('webSocketObservable', () => { }) it('emits an "open" context with empty protocol when no sub-protocol negotiated', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') ws.simulateOpen() const openContexts = getContexts('open') @@ -212,7 +139,7 @@ describe('webSocketObservable', () => { describe('message-in context', () => { it('emits "message-in" with byte-length size for string payloads', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') ws.simulateOpen() const payload = 'hello world' ws.simulateMessage(payload) @@ -223,7 +150,7 @@ describe('webSocketObservable', () => { }) it('emits "message-in" with UTF-8 byte length for multi-byte strings', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') ws.simulateOpen() // 'é' is 2 bytes in UTF-8 and 'あ' is 3 bytes; total is 5 bytes for 2 chars const payload = 'éあ' @@ -233,7 +160,7 @@ describe('webSocketObservable', () => { }) it('emits "message-in" with byteLength for ArrayBuffer payloads', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') ws.simulateOpen() const byteLength = 16 ws.simulateMessage(new ArrayBuffer(byteLength)) @@ -242,7 +169,7 @@ describe('webSocketObservable', () => { }) it('emits "message-in" with byteLength for ArrayBufferView payloads', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') ws.simulateOpen() const viewByteLength = 12 ws.simulateMessage(new Uint8Array(new ArrayBuffer(32), 4, viewByteLength)) @@ -251,7 +178,7 @@ describe('webSocketObservable', () => { }) it('emits "message-in" with size for Blob payloads', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') ws.simulateOpen() const blob = new Blob(['hello']) ws.simulateMessage(blob) @@ -262,7 +189,7 @@ describe('webSocketObservable', () => { describe('message-out context', () => { it('emits "message-out" with size and bufferedAmountPreSend for string payloads', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const bufferedAmountPreSend = 42 ws.bufferedAmount = bufferedAmountPreSend const payload = 'hello' @@ -276,7 +203,7 @@ describe('webSocketObservable', () => { }) it('emits "message-out" with byteLength for ArrayBuffer payloads', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const byteLength = 8 ws.send(new ArrayBuffer(byteLength)) @@ -284,7 +211,7 @@ describe('webSocketObservable', () => { }) it('emits "message-out" with byteLength for ArrayBufferView payloads', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const viewByteLength = 10 ws.send(new Uint8Array(new ArrayBuffer(20), 2, viewByteLength)) @@ -292,7 +219,7 @@ describe('webSocketObservable', () => { }) it('emits "message-out" with size for Blob payloads', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const blob = new Blob(['hello world']) ws.send(blob) @@ -302,7 +229,7 @@ describe('webSocketObservable', () => { describe('closed context', () => { it('emits a "closed" context with code, reason, and wasClean', () => { - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') const closeCode = 1000 const closeReason = 'bye' const wasClean = true @@ -322,14 +249,14 @@ describe('webSocketObservable', () => { subscription?.unsubscribe() subscription = undefined - expect(windowAsWebSocketHost.WebSocket).toBe(FakeWebSocket) + expect(window.WebSocket as unknown).toBe(MockWebSocket) }) it('does not emit any further events after all subscribers unsubscribe', () => { subscription?.unsubscribe() subscription = undefined - const ws = new windowAsWebSocketHost.WebSocket('wss://example.com/socket') + const ws = createMockWebSocket('wss://example.com/socket') ws.simulateOpen() ws.send('hello') diff --git a/packages/browser-core/test/emulate/mockWebSocket.ts b/packages/browser-core/test/emulate/mockWebSocket.ts new file mode 100644 index 0000000000..1e7a02ac3e --- /dev/null +++ b/packages/browser-core/test/emulate/mockWebSocket.ts @@ -0,0 +1,90 @@ +import { registerCleanupTask } from '../registerCleanupTask' + +/** + * Replace the `WebSocket` global with a test double. Specs drive it the way an application + * would: construct a socket through {@link createMockWebSocket}, call `send()` on it, and + * simulate the events the browser would dispatch. + */ +export function mockWebSocket() { + const originalWebSocket = window.WebSocket + + window.WebSocket = MockWebSocket as unknown as typeof WebSocket + + registerCleanupTask(() => { + window.WebSocket = originalWebSocket + }) +} + +/** + * Construct a socket through the `WebSocket` global, so that instrumentation installed on it + * applies, while keeping the mock type to drive the connection. + */ +export function createMockWebSocket(url: string | URL, protocols?: string | string[]): MockWebSocket { + return new (window.WebSocket as unknown as typeof MockWebSocket)(url, protocols) +} + +// A minimal stand-in for the native `WebSocket` constructor. We do not connect to a real server in +// unit tests; instead we expose helpers to simulate the browser dispatching events on the instance. +export class MockWebSocket extends EventTarget { + static readonly CONNECTING = 0 + static readonly OPEN = 1 + static readonly CLOSING = 2 + static readonly CLOSED = 3 + + url: string + protocol = '' + bufferedAmount = 0 + readyState: number = MockWebSocket.CONNECTING + onmessage: ((event: MessageEvent) => void) | null = null + onopen: ((event: Event) => void) | null = null + onclose: ((event: CloseEvent) => void) | null = null + + constructor(url: string | URL, protocols?: string | string[]) { + super() + this.url = resolveWebSocketUrl(String(url)) + if (typeof protocols === 'string') { + this.protocol = protocols + } + } + + send(_data: string | ArrayBufferLike | Blob | ArrayBufferView): void { + // no-op; tests will set `bufferedAmount` before calling send to verify it is sampled. + } + + close(_code?: number, _reason?: string): void { + this.readyState = MockWebSocket.CLOSED + } + + simulateOpen() { + this.readyState = MockWebSocket.OPEN + const event = new Event('open') + this.dispatchEvent(event) + this.onopen?.(event) + } + + simulateMessage(data: unknown) { + const event = new MessageEvent('message', { data }) + this.dispatchEvent(event) + this.onmessage?.(event) + } + + simulateClose(code: number, reason: string, wasClean: boolean) { + this.readyState = MockWebSocket.CLOSED + // CloseEvent is not always constructable in test environments; use a plain Event with assigned fields. + const event = Object.assign(new Event('close'), { code, reason, wasClean }) as CloseEvent + this.dispatchEvent(event) + this.onclose?.(event) + } +} + +// Mimics how a real browser resolves the URL passed to the `WebSocket` constructor: relative URLs +// are resolved against the document location, and `http(s)` schemes are translated to `ws(s)`. +function resolveWebSocketUrl(url: string): string { + const resolved = new URL(url, location.href) + if (resolved.protocol === 'http:') { + resolved.protocol = 'ws:' + } else if (resolved.protocol === 'https:') { + resolved.protocol = 'wss:' + } + return resolved.href +} diff --git a/packages/browser-core/test/index.ts b/packages/browser-core/test/index.ts index ebaecdba33..2a8d794cad 100644 --- a/packages/browser-core/test/index.ts +++ b/packages/browser-core/test/index.ts @@ -17,6 +17,7 @@ export * from './emulate/mockEventBridge' export * from './emulate/mockFlushController' export * from './emulate/mockFetch' export * from './emulate/mockXhr' +export * from './emulate/mockWebSocket' export * from './emulate/mockEventTarget' export * from './emulate/mockRequestIdleCallback' export * from './emulate/mockTelemetry' From 70aebc133111434acd368acdc2d369f6121aa9d3 Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Mon, 10 Aug 2026 21:44:45 +0200 Subject: [PATCH 2/7] =?UTF-8?q?=E2=9C=A8=20Buffer=20WebSocket=20activity?= =?UTF-8?q?=20from=20SDK=20load?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WebSocket joins fetch, XHR, console and runtime errors as a buffered data source, so a socket constructed before init() is recorded rather than lost. The subscription is unconditional: the WebSocket opt-in is only known at init(), and consulting it before instrumenting would miss every connection opened before then — which is the point of the feature. The opt-in is left to the consumer of the source. The full WebSocketContext union crosses the buffer unchanged, with no coalescing in the core buffering layer. Nothing consumes the new source yet, so there is no customer-visible change. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/domain/bufferedData.spec.ts | 46 ++++++++++++++++++- .../browser-core/src/domain/bufferedData.ts | 8 ++++ 2 files changed, 53 insertions(+), 1 deletion(-) diff --git a/packages/browser-core/src/domain/bufferedData.spec.ts b/packages/browser-core/src/domain/bufferedData.spec.ts index d7238297a6..bda8a5eb78 100644 --- a/packages/browser-core/src/domain/bufferedData.spec.ts +++ b/packages/browser-core/src/domain/bufferedData.spec.ts @@ -1,10 +1,20 @@ import { clocksNow } from '@datadog/js-core/time' import { ConsoleApiName } from '@datadog/js-core/util' import type { MockFetch } from '../../test' -import { collectAsyncCalls, mockFetch, mockXhr, registerCleanupTask, replaceMockable, withXhr } from '../../test' +import { + collectAsyncCalls, + createMockWebSocket, + mockFetch, + mockWebSocket, + mockXhr, + registerCleanupTask, + replaceMockable, + withXhr, +} from '../../test' import { Observable } from '../tools/observable' import { resetFetchObservable } from '../browser/fetchObservable' import { resetXhrObservable } from '../browser/xhrObservable' +import { resetWebSocketObservable } from '../browser/webSocketObservable' import { noop } from '../tools/utils/functionUtils' import { resetConsoleObservable } from './console/consoleObservable' import type { BufferedData } from './bufferedData' @@ -136,6 +146,40 @@ describe('startBufferingData', () => { ]) }) + it('collects web socket activity', async () => { + mockWebSocket() + const { observable, stop } = startBufferingData() + const collected: BufferedData[] = [] + const bufferedDataCollectedSpy = jasmine.createSpy() + + registerCleanupTask(() => { + stop() + resetWebSocketObservable() + }) + + observable.subscribe((data) => { + if (data.type === BufferedDataType.WEB_SOCKET) { + collected.push(data) + bufferedDataCollectedSpy() + } + }) + + const webSocket = createMockWebSocket('wss://fake-url/') + + await collectAsyncCalls(bufferedDataCollectedSpy, 1) + + expect(collected).toEqual([ + { + type: BufferedDataType.WEB_SOCKET, + data: jasmine.objectContaining({ + state: 'connecting', + url: 'wss://fake-url/', + instance: webSocket as unknown as WebSocket, + }), + }, + ]) + }) + it('collects console logs', (done) => { spyOn(console, 'error').and.callFake(noop) const { observable, stop } = startBufferingData() diff --git a/packages/browser-core/src/domain/bufferedData.ts b/packages/browser-core/src/domain/bufferedData.ts index ba22f9d542..6ba31b4ce9 100644 --- a/packages/browser-core/src/domain/bufferedData.ts +++ b/packages/browser-core/src/domain/bufferedData.ts @@ -6,6 +6,8 @@ import type { FetchContext } from '../browser/fetchObservable' import { initFetchObservable } from '../browser/fetchObservable' import type { XhrContext } from '../browser/xhrObservable' import { initXhrObservable } from '../browser/xhrObservable' +import type { WebSocketContext } from '../browser/webSocketObservable' +import { initWebSocketObservable } from '../browser/webSocketObservable' import { addTelemetryDebug } from './telemetry' import type { RawError } from './error/error.types' import { trackRuntimeError } from './error/trackRuntimeError' @@ -19,6 +21,7 @@ export const enum BufferedDataType { FETCH, XHR, CONSOLE, + WEB_SOCKET, } export type BufferedData = @@ -26,6 +29,7 @@ export type BufferedData = | { type: BufferedDataType.FETCH; data: FetchContext } | { type: BufferedDataType.XHR; data: XhrContext } | { type: BufferedDataType.CONSOLE; data: ConsoleLog } + | { type: BufferedDataType.WEB_SOCKET; data: WebSocketContext } export function startBufferingData() { const observable = new BufferedObservable(BUFFER_LIMIT, (count) => { @@ -51,6 +55,10 @@ export function startBufferingData() { subscribe(BufferedDataType.FETCH, initFetchObservable()) subscribe(BufferedDataType.XHR, initXhrObservable()) subscribe(BufferedDataType.CONSOLE, initConsoleObservable(Object.values(ConsoleApiName))) + // Subscribed unconditionally: the WebSocket opt-in is only known at init(), and consulting it + // before instrumenting would miss every connection opened before then. The opt-in is left to the + // consumer of this data source. + subscribe(BufferedDataType.WEB_SOCKET, initWebSocketObservable()) return { observable, From 29abcf7d0145ca3575a8f4d8a3338bf674146f31 Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 09:56:22 +0200 Subject: [PATCH 3/7] =?UTF-8?q?=E2=9C=85=20Drive=20WebSocket=20collection?= =?UTF-8?q?=20tests=20through=20instrumentation?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../resource/webSocketCollection.spec.ts | 320 +++++++----------- 1 file changed, 128 insertions(+), 192 deletions(-) diff --git a/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts b/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts index c73773d9fc..c98214a4ed 100644 --- a/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts +++ b/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts @@ -1,7 +1,13 @@ -import type { WebSocketContext } from '@datadog/browser-core' -import { initWebSocketObservable, Observable } from '@datadog/browser-core' -import { mockClock, registerCleanupTask, type Clock } from '@datadog/browser-core/test' -import type { ClocksState, Duration, RelativeTime } from '@datadog/js-core/time' +import { initWebSocketObservable, resetAllowUntrustedEvents, setAllowUntrustedEvents } from '@datadog/browser-core' +import { + createMockWebSocket, + mockClock, + mockWebSocket, + registerCleanupTask, + type Clock, + type MockWebSocket, +} from '@datadog/browser-core/test' +import type { Duration, RelativeTime } from '@datadog/js-core/time' import { elapsed, relativeToClocks } from '@datadog/js-core/time' import { mockViewHistory } from '../../../test' import { VitalType } from '../../rawRumEvent.types' @@ -20,17 +26,16 @@ const UUID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{ describe('webSocketCollection', () => { let lifeCycle: LifeCycle - let wsObservable: Observable let webSocketCompleteEvents: WebSocketCompleteEvent[] - let wsInstance: WebSocket let clock: Clock beforeEach(() => { clock = mockClock() + mockWebSocket() + setAllowUntrustedEvents(true) lifeCycle = new LifeCycle() - wsObservable = new Observable() webSocketCompleteEvents = [] - wsInstance = {} as WebSocket + registerCleanupTask(resetAllowUntrustedEvents) lifeCycle.subscribe(LifeCycleEventType.WEBSOCKET_COMPLETED, (webSocket) => { webSocketCompleteEvents.push(webSocket) }) @@ -44,77 +49,61 @@ describe('webSocketCollection', () => { viewHistory = mockViewHistory(), addDurationVital: (vital: DurationVital) => void = jasmine.createSpy() ) { - return trackWebSocket(lifeCycle, wsObservable, viewHistory, addDurationVital) + const tracker = trackWebSocket(lifeCycle, initWebSocketObservable(), viewHistory, addDurationVital) + registerCleanupTask(tracker.stop) + return tracker } - function notifyConnecting( - startRelative = 0, - url = 'wss://example.com/socket', - startClocks?: ClocksState, - protocols?: string | string[] - ) { - wsObservable.notify({ - state: 'connecting', - instance: wsInstance, - url, - ...(protocols !== undefined ? { protocols } : {}), - startClocks: startClocks ?? relativeToClocks(clock.relative(startRelative)), - }) + function notifyConnecting(startRelative = 0, url = 'wss://example.com/socket', protocols?: string | string[]) { + setClock(startRelative) + return createMockWebSocket(url, protocols) } - function notifyOpen(openRelative = 10, protocol = '', openClocks?: ClocksState) { - wsObservable.notify({ - state: 'open', - instance: wsInstance, - openClocks: openClocks ?? relativeToClocks(clock.relative(openRelative)), - protocol, - }) + function notifyOpen(socket: MockWebSocket, openRelative = 10, protocol = '') { + setClock(openRelative) + socket.protocol = protocol + socket.simulateOpen() } - function notifyMessageIn(at: number, size: number) { - wsObservable.notify({ state: 'message-in', instance: wsInstance, size, at: relativeToClocks(clock.relative(at)) }) + function notifyMessageIn(socket: MockWebSocket, at: number, size: number) { + setClock(at) + socket.simulateMessage('x'.repeat(size)) } - function notifyMessageOut(at: number, size: number, bufferedAmountPreSend = 0) { - wsObservable.notify({ - state: 'message-out', - instance: wsInstance, - size, - bufferedAmountPreSend, - at: relativeToClocks(clock.relative(at)), - }) + function notifyMessageOut(socket: MockWebSocket, at: number, size: number, bufferedAmountPreSend = 0) { + setClock(at) + socket.bufferedAmount = bufferedAmountPreSend + socket.send('x'.repeat(size)) } - function notifyClosed(at: number, code: number, reason: string, wasClean: boolean, atClocks?: ClocksState) { - wsObservable.notify({ - state: 'closed', - instance: wsInstance, - code, - reason, - wasClean, - at: atClocks ?? relativeToClocks(clock.relative(at)), - }) + function notifyClosed(socket: MockWebSocket, at: number, code: number, reason: string, wasClean: boolean) { + setClock(at) + socket.simulateClose(code, reason, wasClean) + } + + function setClock(relative: number) { + clock.setDate(new Date(clock.timeStamp(relative))) } describe('handshakeSucceeded', () => { it('is true when the open event fired before completion', () => { const tracker = startTracking() - notifyConnecting() - notifyOpen(10) - notifyClosed(40, 1000, 'bye', true) + const closedSocket = notifyConnecting() + notifyOpen(closedSocket, 10) + notifyClosed(closedSocket, 40, 1000, 'bye', true) expect(webSocketCompleteEvents[0].handshakeSucceeded).toBeTrue() - notifyConnecting() - notifyOpen(10) + const flushedSocket = notifyConnecting() + notifyOpen(flushedSocket, 10) tracker.flushOpenConnections('session_end') expect(webSocketCompleteEvents[1].handshakeSucceeded).toBeTrue() }) it('is false when the open event never fired before completion', () => { const tracker = startTracking() - notifyConnecting() - notifyClosed(25, 1006, 'abnormal', false) + const closedSocket = notifyConnecting() + notifyClosed(closedSocket, 25, 1006, 'abnormal', false) expect(webSocketCompleteEvents[0].handshakeSucceeded).toBeFalse() notifyConnecting() @@ -133,11 +122,11 @@ describe('webSocketCollection', () => { const closeReason = 'bye' startTracking() - notifyConnecting(0, url) - notifyOpen(10, protocol) - notifyMessageIn(20, messageInSize) - notifyMessageOut(30, messageOutSize, bufferedAmount) - notifyClosed(40, closeCode, closeReason, true) + const socket = notifyConnecting(0, url) + notifyOpen(socket, 10, protocol) + notifyMessageIn(socket, 20, messageInSize) + notifyMessageOut(socket, 30, messageOutSize, bufferedAmount) + notifyClosed(socket, 40, closeCode, closeReason, true) expect(webSocketCompleteEvents.length).toBe(1) const webSocket = webSocketCompleteEvents[0] @@ -157,67 +146,57 @@ describe('webSocketCollection', () => { const url = 'wss://example.com:8443/path/socket%3Froom?token=secret&tenant=acme' startTracking() - notifyConnecting(0, url) - notifyClosed(10, 1000, 'bye', true) + const socket = notifyConnecting(0, url) + notifyClosed(socket, 10, 1000, 'bye', true) expect(webSocketCompleteEvents[0].url).toBe('wss://example.com:8443/path/socket%3Froom') }) it('keeps the server-negotiated protocol on the completed event', () => { startTracking() - notifyConnecting(0, 'wss://example.com/socket', undefined, ['auth-token', 'chat.v1']) - notifyOpen(5, 'chat.v1') - notifyClosed(10, 1000, 'bye', true) + const socket = notifyConnecting(0, 'wss://example.com/socket', ['auth-token', 'chat.v1']) + notifyOpen(socket, 5, 'chat.v1') + notifyClosed(socket, 10, 1000, 'bye', true) expect(webSocketCompleteEvents[0].protocol).toBe('chat.v1') }) it('generates a unique connection_id per connection', () => { startTracking() - notifyConnecting() - notifyClosed(1, 1000, 'reason_a', true) + const firstSocket = notifyConnecting() + notifyClosed(firstSocket, 1, 1000, 'reason_a', true) const firstId = webSocketCompleteEvents[0].connectionId - wsInstance = {} as WebSocket - notifyConnecting() - notifyClosed(1, 1000, 'reason_b', true) + const secondSocket = notifyConnecting() + notifyClosed(secondSocket, 1, 1000, 'reason_b', true) expect(webSocketCompleteEvents[1].connectionId).not.toBe(firstId) }) it('tracks overlapping connections independently with unique connection_ids', () => { - const wsA = {} as WebSocket - const wsB = {} as WebSocket const urlA = 'wss://example.com/socket-a' const urlB = 'wss://example.com/socket-b' startTracking() - wsInstance = wsA - notifyConnecting(0, urlA) - notifyOpen(5) - - wsInstance = wsB - notifyConnecting(10, urlB) - notifyOpen(15) + const socketA = notifyConnecting(0, urlA) + notifyOpen(socketA, 5) - wsInstance = wsA - notifyMessageIn(20, 10) + const socketB = notifyConnecting(10, urlB) + notifyOpen(socketB, 15) - wsInstance = wsB - notifyMessageIn(25, 20) + notifyMessageIn(socketA, 20, 10) + notifyMessageIn(socketB, 25, 20) - wsInstance = wsA - notifyClosed(30, 1000, 'bye-a', true) + notifyClosed(socketA, 30, 1000, 'bye-a', true) expect(webSocketCompleteEvents.length).toBe(1) expect(webSocketCompleteEvents[0].url).toBe(urlA) expect(webSocketCompleteEvents[0].messagesIn).toEqual({ count: 1, size: 10 }) expect(webSocketCompleteEvents[0].trackingEndReason).toBe('close_event') - wsInstance = wsB - notifyMessageOut(35, 5) - notifyClosed(40, 1000, 'bye-b', true) + notifyMessageOut(socketB, 35, 5) + notifyClosed(socketB, 40, 1000, 'bye-b', true) expect(webSocketCompleteEvents.length).toBe(2) expect(webSocketCompleteEvents[1].url).toBe(urlB) @@ -238,12 +217,12 @@ describe('webSocketCollection', () => { const firstMessageOutAt = 17 startTracking() - notifyConnecting() - notifyOpen(openAt) - notifyMessageIn(firstMessageInAt, 1) - notifyMessageIn(25, 1) // not first; should not update - notifyMessageOut(firstMessageOutAt, 1) - notifyClosed(30, 1000, 'bye', true) + const socket = notifyConnecting() + notifyOpen(socket, openAt) + notifyMessageIn(socket, firstMessageInAt, 1) + notifyMessageIn(socket, 25, 1) // not first; should not update + notifyMessageOut(socket, firstMessageOutAt, 1) + notifyClosed(socket, 30, 1000, 'bye', true) const webSocket = webSocketCompleteEvents[0] expect(webSocket.firstMessageInOffset).toBe((firstMessageInAt - openAt) as Duration) @@ -252,24 +231,24 @@ describe('webSocketCollection', () => { it('tracks longestInboundSilence from consecutive message-in only', () => { startTracking() - notifyConnecting() - notifyOpen(10) - notifyMessageIn(20, 1) - notifyMessageIn(50, 1) // gap 30 - notifyMessageIn(75, 1) // gap 25 - notifyClosed(100, 1000, 'bye', true) + const socket = notifyConnecting() + notifyOpen(socket, 10) + notifyMessageIn(socket, 20, 1) + notifyMessageIn(socket, 50, 1) // gap 30 + notifyMessageIn(socket, 75, 1) // gap 25 + notifyClosed(socket, 100, 1000, 'bye', true) expect(webSocketCompleteEvents[0].longestInboundSilence).toBe(30 as Duration) }) it('ignores message-out when computing longestInboundSilence', () => { startTracking() - notifyConnecting() - notifyOpen(10) - notifyMessageIn(20, 1) - notifyMessageOut(100, 1) - notifyMessageIn(130, 1) // gap from last in (20) to 130 - notifyClosed(200, 1000, 'bye', true) + const socket = notifyConnecting() + notifyOpen(socket, 10) + notifyMessageIn(socket, 20, 1) + notifyMessageOut(socket, 100, 1) + notifyMessageIn(socket, 130, 1) // gap from last in (20) to 130 + notifyClosed(socket, 200, 1000, 'bye', true) expect(webSocketCompleteEvents[0].longestInboundSilence).toBe(110 as Duration) }) @@ -279,20 +258,20 @@ describe('webSocketCollection', () => { const closeAt = 50 startTracking() - notifyConnecting() - notifyOpen(10) - notifyMessageIn(lastMessageInAt, 1) - notifyClosed(closeAt, 1000, 'bye', true) + const socket = notifyConnecting() + notifyOpen(socket, 10) + notifyMessageIn(socket, lastMessageInAt, 1) + notifyClosed(socket, closeAt, 1000, 'bye', true) expect(webSocketCompleteEvents[0].inboundIdleDurationBeforeClose).toBe((closeAt - lastMessageInAt) as Duration) }) it('leaves inboundIdleDurationBeforeClose and lastMessageInAt undefined when no message was received', () => { startTracking() - notifyConnecting() - notifyOpen(10) - notifyMessageOut(30, 1) - notifyClosed(50, 1000, 'bye', true) + const socket = notifyConnecting() + notifyOpen(socket, 10) + notifyMessageOut(socket, 30, 1) + notifyClosed(socket, 50, 1000, 'bye', true) expect(webSocketCompleteEvents[0].lastMessageInAt).toBeUndefined() expect(webSocketCompleteEvents[0].inboundIdleDurationBeforeClose).toBeUndefined() @@ -306,9 +285,9 @@ describe('webSocketCollection', () => { const expectedSetupDuration = elapsed(startClocks.timeStamp, openClocks.timeStamp) startTracking() - notifyConnecting(startAt, 'wss://example.com/socket', startClocks) - notifyOpen(openAt, '', openClocks) - notifyClosed(40, 1000, 'bye', true) + const socket = notifyConnecting(startAt) + notifyOpen(socket, openAt) + notifyClosed(socket, 40, 1000, 'bye', true) expect(webSocketCompleteEvents[0].setupDuration).toBe(expectedSetupDuration) }) @@ -323,8 +302,8 @@ describe('webSocketCollection', () => { const expectedSetupDuration = elapsed(startClocks.timeStamp, closeClocks.timeStamp) startTracking() - notifyConnecting(startAt, 'wss://example.com/socket', startClocks) - notifyClosed(closeAt, closeCode, closeReason, false, closeClocks) + const socket = notifyConnecting(startAt) + notifyClosed(socket, closeAt, closeCode, closeReason, false) expect(webSocketCompleteEvents[0].setupDuration).toBe(expectedSetupDuration) }) @@ -342,12 +321,12 @@ describe('webSocketCollection', () => { const peakBufferedAmount = 100 startTracking() - notifyConnecting() - notifyOpen(10) - notifyMessageOut(20, 1, 10) - notifyMessageOut(30, 1, peakBufferedAmount) - notifyMessageOut(40, 1, 50) - notifyClosed(50, 1000, 'bye', true) + const socket = notifyConnecting() + notifyOpen(socket, 10) + notifyMessageOut(socket, 20, 1, 10) + notifyMessageOut(socket, 30, 1, peakBufferedAmount) + notifyMessageOut(socket, 40, 1, 50) + notifyClosed(socket, 50, 1000, 'bye', true) expect(webSocketCompleteEvents[0].bufferedAmountMax).toBe(peakBufferedAmount) }) @@ -366,8 +345,8 @@ describe('webSocketCollection', () => { ) startTracking(viewHistory) - notifyConnecting() - notifyClosed(startViewB, 1000, 'bye', true) + const socket = notifyConnecting() + notifyClosed(socket, startViewB, 1000, 'bye', true) const webSocket = webSocketCompleteEvents[0] expect(webSocket.startViewId).toBe('view-A') @@ -376,9 +355,9 @@ describe('webSocketCollection', () => { it('flushOpenConnections finalizes still-open connections with tracking_end_reason="session_end"', () => { const tracker = startTracking() - notifyConnecting() - notifyOpen(10) - notifyMessageIn(20, 1) + const socket = notifyConnecting() + notifyOpen(socket, 10) + notifyMessageIn(socket, 20, 1) tracker.flushOpenConnections('session_end') @@ -392,10 +371,10 @@ describe('webSocketCollection', () => { it('does not finalize twice when close arrives after flushOpenConnections', () => { const tracker = startTracking() - notifyConnecting() - notifyOpen(10) + const socket = notifyConnecting() + notifyOpen(socket, 10) tracker.flushOpenConnections('session_end') - notifyClosed(20, 1000, 'bye', true) + notifyClosed(socket, 20, 1000, 'bye', true) expect(webSocketCompleteEvents.length).toBe(1) expect(webSocketCompleteEvents[0].trackingEndReason).toBe('session_end') @@ -403,9 +382,9 @@ describe('webSocketCollection', () => { it('stop() unsubscribes from the observable and ignores further events', () => { const tracker = startTracking() - notifyConnecting() + const socket = notifyConnecting() tracker.stop() - notifyClosed(20, 1000, 'bye', true) + notifyClosed(socket, 20, 1000, 'bye', true) expect(webSocketCompleteEvents.length).toBe(0) }) @@ -428,8 +407,8 @@ describe('webSocketCollection', () => { it('uses a fresh UUID as the vital id and the connectionId in the context', () => { const addDurationVital = jasmine.createSpy<(vital: DurationVital) => void>() startTracking(mockViewHistory(), addDurationVital) - notifyConnecting() - notifyClosed(1, 1000, 'bye', true) + const socket = notifyConnecting() + notifyClosed(socket, 1, 1000, 'bye', true) const vital = addDurationVital.calls.first().args[0] expect(vital.id).not.toBe(webSocketCompleteEvents[0].connectionId) @@ -452,11 +431,11 @@ describe('webSocketCollection', () => { const addDurationVital = jasmine.createSpy<(vital: DurationVital) => void>() startTracking(viewHistory, addDurationVital) - notifyConnecting(startView, 'wss://example.com/socket?token=secret&tenant=acme', undefined, [ + const socket = notifyConnecting(startView, 'wss://example.com/socket?token=secret&tenant=acme', [ 'auth-token', 'chat.v1', ]) - notifyClosed(1, 1000, 'bye', true) + notifyClosed(socket, 1, 1000, 'bye', true) const vital = addDurationVital.calls.first().args[0] expect(vital.context).toEqual({ @@ -470,7 +449,7 @@ describe('webSocketCollection', () => { const addDurationVital = jasmine.createSpy<(vital: DurationVital) => void>() startTracking(mockViewHistory(), addDurationVital) - notifyConnecting(0, 'wss://example.com/socket', undefined, protocols) + notifyConnecting(0, 'wss://example.com/socket', protocols) const context = addDurationVital.calls.mostRecent().args[0].context expect('protocols' in context).toBeFalse() @@ -482,8 +461,8 @@ describe('webSocketCollection', () => { it('emits a duration-0 vital at close time on a close event', () => { const addDurationVital = jasmine.createSpy<(vital: DurationVital) => void>() startTracking(mockViewHistory(), addDurationVital) - notifyConnecting() - notifyClosed(40, 1000, 'bye', true) + const socket = notifyConnecting() + notifyClosed(socket, 40, 1000, 'bye', true) const connectingVital = addDurationVital.calls.argsFor(0)[0] const closedVital = addDurationVital.calls.argsFor(1)[0] @@ -523,55 +502,12 @@ describe('webSocketCollection', () => { }) describe('startWebSocketCollection', () => { - const wsInstance = {} as WebSocket - const wsUrl = 'wss://example.com/socket' - const singletonObservable = () => initWebSocketObservable() - function startCollection() { const collection = startWebSocketCollection(lifeCycle, mockViewHistory(), jasmine.createSpy()) registerCleanupTask(() => collection.stop()) return collection } - function notifyConnecting(offsetMs = 0) { - singletonObservable().notify({ - state: 'connecting', - instance: wsInstance, - url: wsUrl, - startClocks: relativeToClocks(clock.relative(offsetMs)), - }) - } - - function notifyOpen(openRelative = 10, protocol = '') { - singletonObservable().notify({ - state: 'open', - instance: wsInstance, - openClocks: relativeToClocks(clock.relative(openRelative)), - protocol, - }) - } - - function notifyMessageOut(at: number, size: number, bufferedAmountPreSend = 0) { - singletonObservable().notify({ - state: 'message-out', - instance: wsInstance, - size, - bufferedAmountPreSend, - at: relativeToClocks(clock.relative(at)), - }) - } - - function notifyClosed(at: number, code: number, reason: string, wasClean: boolean) { - singletonObservable().notify({ - state: 'closed', - instance: wsInstance, - code, - reason, - wasClean, - at: relativeToClocks(clock.relative(at)), - }) - } - it('finalizes open connections with tracking_end_reason="session_end" when the session expires', () => { const endClocks = relativeToClocks(clock.relative(40)) startCollection() @@ -590,21 +526,21 @@ describe('webSocketCollection', () => { it('ignores further WebSocket events from the same instance after stop()', () => { const collection = startCollection() - notifyConnecting() + const socket = notifyConnecting() collection.stop() const eventCountAfterStop = webSocketCompleteEvents.length - notifyClosed(1000, 1000, 'bye', true) + notifyClosed(socket, 1000, 1000, 'bye', true) expect(webSocketCompleteEvents.length).toBe(eventCountAfterStop) }) it('ignores further WebSocket events from the same instance after the session expires', () => { startCollection() - notifyConnecting() - notifyOpen(10) - notifyMessageOut(20, 10) + const socket = notifyConnecting() + notifyOpen(socket, 10) + notifyMessageOut(socket, 20, 10) expireSession() @@ -612,8 +548,8 @@ describe('webSocketCollection', () => { expect(webSocketCompleteEvents[0].trackingEndReason).toBe('session_end') expect(webSocketCompleteEvents[0].messagesOut).toEqual({ count: 1, size: 10 }) - notifyMessageOut(40, 7) - notifyClosed(50, 1000, 'bye', true) + notifyMessageOut(socket, 40, 7) + notifyClosed(socket, 50, 1000, 'bye', true) expect(webSocketCompleteEvents.length).toBe(1) expect(webSocketCompleteEvents[0].trackingEndReason).toBe('session_end') From ba4b6353e4623534fd51925633c01aa7f8937d6d Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 10:59:28 +0200 Subject: [PATCH 4/7] =?UTF-8?q?=E2=9C=A8=20Collect=20WebSocket=20connectio?= =?UTF-8?q?ns=20opened=20before=20init()?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WebSocket collection now consumes the buffered data observable instead of subscribing to the WebSocket observable directly, so connections opened before init() are replayed and reported as complete resource events. The opt-in gate (trackResources plus betaTrackWebSockets or the TRACK_WEBSOCKETS experimental feature) moves into the collection entry point, since it is not knowable until init(). RUM startup now calls it unconditionally; the returned stop handle is a no-op when the gate is closed. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/boot/startRum.spec.ts | 79 +------ .../browser-rum-core/src/boot/startRum.ts | 17 +- .../resource/webSocketCollection.spec.ts | 195 +++++++++++++++++- .../domain/resource/webSocketCollection.ts | 43 +++- 4 files changed, 236 insertions(+), 98 deletions(-) diff --git a/packages/browser-rum-core/src/boot/startRum.spec.ts b/packages/browser-rum-core/src/boot/startRum.spec.ts index d9a80328b6..dd3e8fc3f3 100644 --- a/packages/browser-rum-core/src/boot/startRum.spec.ts +++ b/packages/browser-rum-core/src/boot/startRum.spec.ts @@ -1,15 +1,7 @@ import { ONE_SECOND, toServerDuration, relativeNow, relativeToClocks } from '@datadog/js-core/time' import type { Duration } from '@datadog/js-core/time' import type { BufferedData, SessionManager } from '@datadog/browser-core' -import { - Observable, - findLast, - noop, - createIdentityEncoder, - BufferedObservable, - addExperimentalFeatures, - ExperimentalFeature, -} from '@datadog/browser-core' +import { Observable, findLast, noop, createIdentityEncoder, BufferedObservable } from '@datadog/browser-core' import type { Clock, SessionManagerMock } from '@datadog/browser-core/test' import { createNewEvent, @@ -268,72 +260,3 @@ describe('view events', () => { expect(lastRumViewEvent._dd.sdk_name).toBe('rum') }) }) - -describe('WebSocket resource collection activation', () => { - it('starts when betaTrackWebSockets is enabled and resources are tracked', () => { - const originalWebSocket = window.WebSocket - const { stop } = startRumStub( - new LifeCycle(), - mockRumConfiguration({ trackResources: true, betaTrackWebSockets: true }), - createSessionManagerMock(), - noop - ) - registerCleanupTask(stop) - - expect(window.WebSocket).not.toBe(originalWebSocket) - }) - - it('starts when the experimental flag is enabled and resources are tracked', () => { - addExperimentalFeatures([ExperimentalFeature.TRACK_WEBSOCKETS]) - const originalWebSocket = window.WebSocket - const { stop } = startRumStub( - new LifeCycle(), - mockRumConfiguration({ trackResources: true, betaTrackWebSockets: false }), - createSessionManagerMock(), - noop - ) - registerCleanupTask(stop) - - expect(window.WebSocket).not.toBe(originalWebSocket) - }) - - it('does not start when neither enablement mechanism is active', () => { - const originalWebSocket = window.WebSocket - const { stop } = startRumStub( - new LifeCycle(), - mockRumConfiguration({ trackResources: true, betaTrackWebSockets: false }), - createSessionManagerMock(), - noop - ) - registerCleanupTask(stop) - - expect(window.WebSocket).toBe(originalWebSocket) - }) - - it('does not start from the beta option when resource tracking is disabled', () => { - const originalWebSocket = window.WebSocket - const { stop } = startRumStub( - new LifeCycle(), - mockRumConfiguration({ trackResources: false, betaTrackWebSockets: true }), - createSessionManagerMock(), - noop - ) - registerCleanupTask(stop) - - expect(window.WebSocket).toBe(originalWebSocket) - }) - - it('does not start from the experimental flag when resource tracking is disabled', () => { - addExperimentalFeatures([ExperimentalFeature.TRACK_WEBSOCKETS]) - const originalWebSocket = window.WebSocket - const { stop } = startRumStub( - new LifeCycle(), - mockRumConfiguration({ trackResources: false, betaTrackWebSockets: false }), - createSessionManagerMock(), - noop - ) - registerCleanupTask(stop) - - expect(window.WebSocket).toBe(originalWebSocket) - }) -}) diff --git a/packages/browser-rum-core/src/boot/startRum.ts b/packages/browser-rum-core/src/boot/startRum.ts index 0ce5ec300c..d2df5ea19c 100644 --- a/packages/browser-rum-core/src/boot/startRum.ts +++ b/packages/browser-rum-core/src/boot/startRum.ts @@ -17,8 +17,6 @@ import { startUserContext, startTabContext, ErrorSource, - isExperimentalFeatureEnabled, - ExperimentalFeature, } from '@datadog/browser-core' import { clocksNow } from '@datadog/js-core/time' import { createDOMMutationObservable } from '../browser/domMutationObservable' @@ -226,13 +224,14 @@ export function startRumEventCollection( const vitalCollection = startVitalCollection(lifeCycle, pageStateHistory) - if ( - configuration.trackResources && - (configuration.betaTrackWebSockets || isExperimentalFeatureEnabled(ExperimentalFeature.TRACK_WEBSOCKETS)) - ) { - const webSocketCollection = startWebSocketCollection(lifeCycle, viewHistory, vitalCollection.addDurationVital) - cleanupTasks.push(webSocketCollection.stop) - } + const webSocketCollection = startWebSocketCollection( + lifeCycle, + configuration, + viewHistory, + vitalCollection.addDurationVital, + bufferedDataObservable + ) + cleanupTasks.push(webSocketCollection.stop) const internalContext = startInternalContext( configuration.applicationId, diff --git a/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts b/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts index c98214a4ed..48a2ba17e8 100644 --- a/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts +++ b/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts @@ -1,5 +1,16 @@ -import { initWebSocketObservable, resetAllowUntrustedEvents, setAllowUntrustedEvents } from '@datadog/browser-core' +import type { BufferedData } from '@datadog/browser-core' import { + addExperimentalFeatures, + BufferedDataType, + ExperimentalFeature, + initWebSocketObservable, + Observable, + resetAllowUntrustedEvents, + setAllowUntrustedEvents, + startBufferingData, +} from '@datadog/browser-core' +import { + collectAsyncCalls, createMockWebSocket, mockClock, mockWebSocket, @@ -9,7 +20,7 @@ import { } from '@datadog/browser-core/test' import type { Duration, RelativeTime } from '@datadog/js-core/time' import { elapsed, relativeToClocks } from '@datadog/js-core/time' -import { mockViewHistory } from '../../../test' +import { mockRumConfiguration, mockViewHistory } from '../../../test' import { VitalType } from '../../rawRumEvent.types' import type { ViewHistoryEntry } from '../contexts/viewHistory' import { LifeCycle, LifeCycleEventType } from '../lifeCycle' @@ -45,11 +56,23 @@ describe('webSocketCollection', () => { lifeCycle.notify(LifeCycleEventType.SESSION_EXPIRED, { endClocks }) } + // Feeds WebSocket instrumentation into a plain `BufferedData` observable, so that specs keep + // driving real sockets while observing the events synchronously. `startBufferingData` is used + // instead where the asynchronous buffer replay is what is under test. + function createWebSocketDataObservable() { + const observable = new Observable() + const subscription = initWebSocketObservable().subscribe((data) => + observable.notify({ type: BufferedDataType.WEB_SOCKET, data }) + ) + registerCleanupTask(() => subscription.unsubscribe()) + return observable + } + function startTracking( viewHistory = mockViewHistory(), addDurationVital: (vital: DurationVital) => void = jasmine.createSpy() ) { - const tracker = trackWebSocket(lifeCycle, initWebSocketObservable(), viewHistory, addDurationVital) + const tracker = trackWebSocket(lifeCycle, createWebSocketDataObservable(), viewHistory, addDurationVital) registerCleanupTask(tracker.stop) return tracker } @@ -502,12 +525,174 @@ describe('webSocketCollection', () => { }) describe('startWebSocketCollection', () => { - function startCollection() { - const collection = startWebSocketCollection(lifeCycle, mockViewHistory(), jasmine.createSpy()) + function startCollection( + configuration = mockRumConfiguration({ betaTrackWebSockets: true }), + bufferedDataObservable = createWebSocketDataObservable() + ) { + const collection = startWebSocketCollection( + lifeCycle, + configuration, + mockViewHistory(), + jasmine.createSpy(), + bufferedDataObservable + ) registerCleanupTask(() => collection.stop()) return collection } + describe('opt-in gate', () => { + ;( + [ + { trackResources: true, betaTrackWebSockets: true, experimentalFeature: false, collects: true }, + { trackResources: true, betaTrackWebSockets: false, experimentalFeature: true, collects: true }, + { trackResources: true, betaTrackWebSockets: false, experimentalFeature: false, collects: false }, + { trackResources: false, betaTrackWebSockets: true, experimentalFeature: false, collects: false }, + { trackResources: false, betaTrackWebSockets: false, experimentalFeature: true, collects: false }, + ] as const + ).forEach(({ trackResources, betaTrackWebSockets, experimentalFeature, collects }) => { + it(`${collects ? 'collects' : 'does not collect'} with trackResources=${trackResources}, betaTrackWebSockets=${betaTrackWebSockets}, TRACK_WEBSOCKETS=${experimentalFeature}`, () => { + if (experimentalFeature) { + addExperimentalFeatures([ExperimentalFeature.TRACK_WEBSOCKETS]) + } + + startCollection(mockRumConfiguration({ trackResources, betaTrackWebSockets })) + const socket = notifyConnecting() + notifyOpen(socket, 10) + notifyClosed(socket, 20, 1000, 'bye', true) + + expect(webSocketCompleteEvents.length).toBe(collects ? 1 : 0) + }) + }) + + it('does not subscribe to the buffered data observable when the gate is closed', () => { + const bufferedDataObservable = createWebSocketDataObservable() + const subscribeSpy = spyOn(bufferedDataObservable, 'subscribe').and.callThrough() + + startCollection( + mockRumConfiguration({ trackResources: true, betaTrackWebSockets: false }), + bufferedDataObservable + ) + + expect(subscribeSpy).not.toHaveBeenCalled() + }) + }) + + // WebSocket activity is instrumented and buffered from SDK load; collection only subscribes at + // init(), and receives everything that happened before as a replayed burst. + describe('connections started before collection subscribed', () => { + let bufferedDataObservable: Observable + let completeSpy: jasmine.Spy<(event: WebSocketCompleteEvent) => void> + + beforeEach(() => { + const buffering = startBufferingData() + bufferedDataObservable = buffering.observable + registerCleanupTask(buffering.stop) + + completeSpy = jasmine.createSpy() + lifeCycle.subscribe(LifeCycleEventType.WEBSOCKET_COMPLETED, completeSpy) + }) + + function subscribeCollection(addDurationVital: (vital: DurationVital) => void = jasmine.createSpy()) { + const collection = startWebSocketCollection( + lifeCycle, + mockRumConfiguration({ betaTrackWebSockets: true }), + mockViewHistory(), + addDurationVital, + bufferedDataObservable + ) + registerCleanupTask(() => collection.stop()) + } + + it('reports a connection that also completed before collection subscribed', async () => { + const socket = notifyConnecting(0) + notifyOpen(socket, 10) + notifyMessageIn(socket, 20, 30) + notifyClosed(socket, 40, 1000, 'bye', true) + + subscribeCollection() + await collectAsyncCalls(completeSpy) + + expect(webSocketCompleteEvents.length).toBe(1) + const webSocket = webSocketCompleteEvents[0] + expect(webSocket.messagesIn).toEqual({ count: 1, size: 30 }) + // measured from the real constructor call, not from the subscription + expect(webSocket.setupDuration).toBe(10 as Duration) + }) + + it('reports a connection spanning the subscription exactly once', async () => { + const socket = notifyConnecting(0) + notifyOpen(socket, 10) + notifyMessageIn(socket, 20, 30) + + // a connection that completed early gives a deterministic signal that the replay is over + const replayedSocket = notifyConnecting(21, 'wss://example.com/replayed') + notifyClosed(replayedSocket, 22, 1000, 'bye', true) + const replaySpy = jasmine.createSpy<(event: WebSocketCompleteEvent) => void>() + const replaySubscription = lifeCycle.subscribe(LifeCycleEventType.WEBSOCKET_COMPLETED, replaySpy) + + subscribeCollection() + await collectAsyncCalls(replaySpy) + replaySubscription.unsubscribe() + + notifyMessageIn(socket, 50, 5) + notifyClosed(socket, 60, 1000, 'bye', true) + + // a single event, merging what was replayed with what came in live + expect(webSocketCompleteEvents.length).toBe(2) + const webSocket = webSocketCompleteEvents[1] + expect(webSocket.messagesIn).toEqual({ count: 2, size: 35 }) + expect(webSocket.setupDuration).toBe(10 as Duration) + }) + + it('keeps a distinct connection id per early connection and emits both vitals', async () => { + const addDurationVital = jasmine.createSpy<(vital: DurationVital) => void>() + const firstSocket = notifyConnecting(0, 'wss://example.com/socket-a') + const secondSocket = notifyConnecting(5, 'wss://example.com/socket-b') + notifyClosed(firstSocket, 10, 1000, 'bye-a', true) + notifyClosed(secondSocket, 20, 1000, 'bye-b', true) + + subscribeCollection(addDurationVital) + await collectAsyncCalls(completeSpy, 2) + + const [first, second] = webSocketCompleteEvents + expect(first.connectionId).not.toBe(second.connectionId) + + const vitalNames = addDurationVital.calls.all().map((call) => call.args[0].name) + expect(vitalNames).toEqual([ + WEBSOCKET_CONNECTING_VITAL_NAME, + WEBSOCKET_CONNECTING_VITAL_NAME, + WEBSOCKET_CLOSED_VITAL_NAME, + WEBSOCKET_CLOSED_VITAL_NAME, + ]) + }) + }) + + it('leaves application-set handlers and exchanged payloads untouched', () => { + const openHandler = jasmine.createSpy<(event: Event) => void>() + const messageHandler = jasmine.createSpy<(event: MessageEvent) => void>() + const closeHandler = jasmine.createSpy<(event: CloseEvent) => void>() + // spied before instrumentation is installed, so that the instrumented `send` delegates to it + const sendSpy = spyOn(window.WebSocket.prototype, 'send').and.callThrough() + + startCollection() + const socket = notifyConnecting() + socket.onopen = openHandler + socket.onmessage = messageHandler + socket.onclose = closeHandler + + notifyOpen(socket, 10) + setClock(20) + socket.simulateMessage('hello') + setClock(30) + socket.send('world') + notifyClosed(socket, 40, 1000, 'bye', true) + + expect(openHandler).toHaveBeenCalledTimes(1) + expect(messageHandler.calls.mostRecent().args[0].data).toBe('hello') + expect(closeHandler.calls.mostRecent().args[0].code).toBe(1000) + expect(sendSpy).toHaveBeenCalledOnceWith('world') + }) + it('finalizes open connections with tracking_end_reason="session_end" when the session expires', () => { const endClocks = relativeToClocks(clock.relative(40)) startCollection() diff --git a/packages/browser-rum-core/src/domain/resource/webSocketCollection.ts b/packages/browser-rum-core/src/domain/resource/webSocketCollection.ts index 4a0cd37564..44e1ba4067 100644 --- a/packages/browser-rum-core/src/domain/resource/webSocketCollection.ts +++ b/packages/browser-rum-core/src/domain/resource/webSocketCollection.ts @@ -1,9 +1,17 @@ -import type { Observable, WebSocketContext } from '@datadog/browser-core' -import { generateUUID, initWebSocketObservable, sanitize } from '@datadog/browser-core' +import type { BufferedData, Observable } from '@datadog/browser-core' +import { + BufferedDataType, + ExperimentalFeature, + generateUUID, + isExperimentalFeatureEnabled, + noop, + sanitize, +} from '@datadog/browser-core' import type { ClocksState, Duration, TimeStamp } from '@datadog/js-core/time' import { clocksNow, elapsed } from '@datadog/js-core/time' import { buildUrl } from '@datadog/js-core/util' import { VitalType } from '../../rawRumEvent.types' +import type { RumConfiguration } from '../configuration' import type { ViewHistory } from '../contexts/viewHistory' import type { LifeCycle } from '../lifeCycle' import { LifeCycleEventType } from '../lifeCycle' @@ -62,12 +70,23 @@ export interface WebSocketConnectionTracker { stop: () => void } +/** + * The opt-in is enforced here rather than by withholding instrumentation, which happens from SDK + * load (see `startBufferingData`). When it is closed, nothing is subscribed nor allocated and the + * returned stop handle is a no-op. + */ export function startWebSocketCollection( lifeCycle: LifeCycle, + configuration: RumConfiguration, viewHistory: ViewHistory, - addDurationVital: (vital: DurationVital) => void + addDurationVital: (vital: DurationVital) => void, + bufferedDataObservable: Observable ) { - const tracker = trackWebSocket(lifeCycle, initWebSocketObservable(), viewHistory, addDurationVital) + if (!isWebSocketCollectionEnabled(configuration)) { + return { stop: noop } + } + + const tracker = trackWebSocket(lifeCycle, bufferedDataObservable, viewHistory, addDurationVital) // Session-boundary cleanup happens on SESSION_EXPIRED (fired before SESSION_RENEWED). Open // connections are finalized once with trackingEndReason "session_end"; later events on the same @@ -85,9 +104,16 @@ export function startWebSocketCollection( } } +function isWebSocketCollectionEnabled(configuration: RumConfiguration) { + return ( + configuration.trackResources && + (configuration.betaTrackWebSockets || isExperimentalFeatureEnabled(ExperimentalFeature.TRACK_WEBSOCKETS)) + ) +} + export function trackWebSocket( lifeCycle: LifeCycle, - webSocketContextObservable: Observable, + bufferedDataObservable: Observable, viewHistory: ViewHistory, addDurationVital: (vital: DurationVital) => void ): WebSocketConnectionTracker { @@ -114,7 +140,12 @@ export function trackWebSocket( }) } - const subscription = webSocketContextObservable.subscribe((context) => { + const subscription = bufferedDataObservable.subscribe((bufferedData) => { + if (bufferedData.type !== BufferedDataType.WEB_SOCKET) { + return + } + + const context = bufferedData.data switch (context.state) { case 'connecting': { const connectionId = generateUUID() From 0386e91cd4d75f44a48cc82ac1d763581f7303e4 Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 11:23:38 +0200 Subject: [PATCH 5/7] =?UTF-8?q?=E2=9C=85=20Test=20WebSocket=20passthrough?= =?UTF-8?q?=20at=20the=20instrumentation=20layer?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The collection layer only observes what webSocketObservable emits, so asserting that application handlers and payloads survive belongs with the instrumentation. Handler passthrough was already covered there; add the missing send-payload case and drop the collection-level duplicate. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/browser/webSocketObservable.spec.ts | 10 +++++++ .../test/emulate/mockWebSocket.ts | 7 +++-- .../resource/webSocketCollection.spec.ts | 26 ------------------- 3 files changed, 15 insertions(+), 28 deletions(-) diff --git a/packages/browser-core/src/browser/webSocketObservable.spec.ts b/packages/browser-core/src/browser/webSocketObservable.spec.ts index 019cc90e6b..5163d5e62c 100644 --- a/packages/browser-core/src/browser/webSocketObservable.spec.ts +++ b/packages/browser-core/src/browser/webSocketObservable.spec.ts @@ -111,6 +111,16 @@ describe('webSocketObservable', () => { expect(customerHandler).toHaveBeenCalledTimes(1) expect(getContexts('closed').length).toBe(1) }) + + it('forwards the sent payload to the native send unaltered', () => { + const ws = createMockWebSocket('wss://example.com/socket') + const payload = 'hello' + + ws.send(payload) + + expect(ws.sentData).toEqual([payload]) + expect(getContexts('message-out').length).toBe(1) + }) }) describe('open context', () => { diff --git a/packages/browser-core/test/emulate/mockWebSocket.ts b/packages/browser-core/test/emulate/mockWebSocket.ts index 1e7a02ac3e..7512363679 100644 --- a/packages/browser-core/test/emulate/mockWebSocket.ts +++ b/packages/browser-core/test/emulate/mockWebSocket.ts @@ -38,6 +38,9 @@ export class MockWebSocket extends EventTarget { onmessage: ((event: MessageEvent) => void) | null = null onopen: ((event: Event) => void) | null = null onclose: ((event: CloseEvent) => void) | null = null + // Payloads that reached the socket, in order, so that specs can check instrumentation forwards + // them unaltered. Tests set `bufferedAmount` before calling send to verify it is sampled. + sentData: Array = [] constructor(url: string | URL, protocols?: string | string[]) { super() @@ -47,8 +50,8 @@ export class MockWebSocket extends EventTarget { } } - send(_data: string | ArrayBufferLike | Blob | ArrayBufferView): void { - // no-op; tests will set `bufferedAmount` before calling send to verify it is sampled. + send(data: string | ArrayBufferLike | Blob | ArrayBufferView): void { + this.sentData.push(data) } close(_code?: number, _reason?: string): void { diff --git a/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts b/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts index 48a2ba17e8..ce41053ee8 100644 --- a/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts +++ b/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts @@ -667,32 +667,6 @@ describe('webSocketCollection', () => { }) }) - it('leaves application-set handlers and exchanged payloads untouched', () => { - const openHandler = jasmine.createSpy<(event: Event) => void>() - const messageHandler = jasmine.createSpy<(event: MessageEvent) => void>() - const closeHandler = jasmine.createSpy<(event: CloseEvent) => void>() - // spied before instrumentation is installed, so that the instrumented `send` delegates to it - const sendSpy = spyOn(window.WebSocket.prototype, 'send').and.callThrough() - - startCollection() - const socket = notifyConnecting() - socket.onopen = openHandler - socket.onmessage = messageHandler - socket.onclose = closeHandler - - notifyOpen(socket, 10) - setClock(20) - socket.simulateMessage('hello') - setClock(30) - socket.send('world') - notifyClosed(socket, 40, 1000, 'bye', true) - - expect(openHandler).toHaveBeenCalledTimes(1) - expect(messageHandler.calls.mostRecent().args[0].data).toBe('hello') - expect(closeHandler.calls.mostRecent().args[0].code).toBe(1000) - expect(sendSpy).toHaveBeenCalledOnceWith('world') - }) - it('finalizes open connections with tracking_end_reason="session_end" when the session expires', () => { const endClocks = relativeToClocks(clock.relative(40)) startCollection() From 179104c7c2445571fcc7bdedbde6f868b3e4f09b Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 14:20:44 +0200 Subject: [PATCH 6/7] =?UTF-8?q?=E2=9C=85=20Add=20a=20pre-init=20script=20i?= =?UTF-8?q?njection=20point=20to=20E2E=20page=20setups?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The page setup built the SDK bundle load and the init() call together, with no way to run application script in between. That window — SDK loaded and instrumenting, but not configured yet — is exactly what the buffered data sources serve, and it was untestable end to end. `createTest().withPreInitScript(js)` now emits that script as `window.DD_PRE_INIT`, called from the init site of the bundle, npm and async setups. The script may return a promise, in which case init() waits for it to settle, so an asynchronous exchange can complete entirely before the SDK starts. The bundle setup now emits every SDK script tag before any init call, so the pre-init script runs with all of them instrumenting the page. Setups that don't own their init call site reject the option rather than silently ignoring it. Co-Authored-By: Claude Opus 5 (1M context) --- test/apps/vanilla/app.ts | 13 ++- test/e2e/AGENTS.md | 3 + test/e2e/lib/framework/createTest.ts | 20 +++++ test/e2e/lib/framework/pageSetups.ts | 117 ++++++++++++++++++++++----- 4 files changed, 134 insertions(+), 19 deletions(-) diff --git a/test/apps/vanilla/app.ts b/test/apps/vanilla/app.ts index 0a67e56915..31a87e068b 100644 --- a/test/apps/vanilla/app.ts +++ b/test/apps/vanilla/app.ts @@ -7,10 +7,11 @@ declare global { LOGS_INIT?: () => void RUM_INIT?: () => void DEBUGGER_INIT?: () => void + DD_PRE_INIT?: () => unknown } } -if (typeof window !== 'undefined') { +function runInits() { if (window.LOGS_INIT) { window.LOGS_INIT() } @@ -22,6 +23,16 @@ if (typeof window !== 'undefined') { if (window.DEBUGGER_INIT) { window.DEBUGGER_INIT() } +} + +if (typeof window !== 'undefined') { + if (window.DD_PRE_INIT) { + // Application code running once the SDK modules are evaluated — and therefore instrumenting + // the page — but before init(). It may return a promise to hold init() back until it settles. + void Promise.resolve(window.DD_PRE_INIT()).then(runInits) + } else { + runInits() + } } else { // compat test datadogLogs.init({ clientToken: 'xxx', beforeSend: undefined }) diff --git a/test/e2e/AGENTS.md b/test/e2e/AGENTS.md index 7b8e555dd6..45a2bd4a77 100644 --- a/test/e2e/AGENTS.md +++ b/test/e2e/AGENTS.md @@ -125,6 +125,9 @@ test.describe('feature name', () => { - `.withHead(html)` - Add content to `` - `.withBody(html)` - Add content to `` +- `.withPreInitScript(js)` - Run application code after the SDK loads but before `init()`, to test + what the SDK buffers before being started. The script is a function body and may return a promise + that `init()` waits on. - `.withReactApp(name)` - Use a React test app - `.withExtension(ext)` - Test with browser extension diff --git a/test/e2e/lib/framework/createTest.ts b/test/e2e/lib/framework/createTest.ts index 21a501d1b6..2c2480f112 100644 --- a/test/e2e/lib/framework/createTest.ts +++ b/test/e2e/lib/framework/createTest.ts @@ -105,6 +105,7 @@ class TestBuilder { private remoteConfiguration?: RemoteConfiguration = undefined private head = '' private body = '' + private preInitScript = '' private baseUrlHooks: UrlHook[] = [] private eventBridge: EventBridgeOptions | undefined private setups: Array<{ factory: SetupFactory; name?: string }> = DEFAULT_SETUPS @@ -162,6 +163,24 @@ class TestBuilder { return this } + /** + * Runs application code in the window between the SDK load and `init()`: the SDK is already + * instrumenting the page, but is not configured yet. Useful to test what the SDK buffers before + * being started. + * + * `script` is raw JavaScript, evaluated as a function body. It may return a promise, in which + * case `init()` waits for it to settle — so an asynchronous exchange can complete entirely + * before the SDK starts. + * + * Supported by the bundle, npm and async setups. With the async setup the script runs from the + * first bundle to become ready, so a test configuring several products can't assume the others + * are loaded yet. + */ + withPreInitScript(script: string) { + this.preInitScript = script + return this + } + withEventBridge(options: EventBridgeOptions = {}) { this.eventBridge = options return this @@ -294,6 +313,7 @@ class TestBuilder { const setupOptions: SetupOptions = { body: this.body, head: this.head, + preInitScript: this.preInitScript, logs: this.logsConfiguration, rum: this.rumConfiguration, debugger: this.debuggerConfiguration, diff --git a/test/e2e/lib/framework/pageSetups.ts b/test/e2e/lib/framework/pageSetups.ts index d939baf4cf..1fbf4aed86 100644 --- a/test/e2e/lib/framework/pageSetups.ts +++ b/test/e2e/lib/framework/pageSetups.ts @@ -23,6 +23,7 @@ export interface SetupOptions { eventBridge?: EventBridgeOptions head?: string body?: string + preInitScript?: string baseUrlHooks: UrlHook[] context: { run_id: string @@ -70,6 +71,58 @@ export const DEFAULT_SETUPS = { name: 'bundle', factory: bundleSetup }, ] +/** + * Defines `window.DD_PRE_INIT`, the application script that runs once the SDK is loaded (and + * therefore instrumenting the page) but before `init()` is called. The script body may return a + * promise: `init()` then waits for it to settle, which lets a scenario complete asynchronous work + * — a request, a WebSocket exchange — entirely within the pre-init window. + * + * Setups call it from their init site through {@link afterPreInitScript}, so the definition itself + * can be emitted anywhere in the head. It runs at most once, on the first init site to reach it. + */ +function preInitDefinition(options: SetupOptions) { + if (!options.preInitScript) { + return '' + } + + return html`` +} + +/** + * Setups that don't control the init call site can't offer the pre-init window. Fail loudly rather + * than silently ignoring the script and letting the scenario pass for the wrong reason. + */ +function rejectPreInitScript(options: SetupOptions, setupName: string) { + if (options.preInitScript) { + throw new Error(`withPreInitScript() is not supported by the ${setupName} setup`) + } +} + +/** Wraps init code so that it runs after the pre-init script (see {@link preInitDefinition}). */ +function afterPreInitScript(options: SetupOptions, initCode: string) { + if (!options.preInitScript) { + return initCode + } + + return js`Promise.resolve(window.DD_PRE_INIT()).then(function () { + ${initCode} + })` +} + export function asyncSetup(options: SetupOptions, servers: Servers) { let header = options.head || '' let footer = '' @@ -82,6 +135,8 @@ export function asyncSetup(options: SetupOptions, servers: Servers) { header += setupExtension(options, servers) } + header += preInitDefinition(options) + function formatSnippet(url: string, globalName: string) { return `(function(h,o,u,n,d) { h=h[d]=h[d]||{q:[],onReady:function(c){h.q.push(c)}} @@ -96,8 +151,11 @@ n=o.getElementsByTagName(u)[0];n.parentNode.insertBefore(d,n) footer += html`` } @@ -106,8 +164,11 @@ n=o.getElementsByTagName(u)[0];n.parentNode.insertBefore(d,n) footer += html`` } @@ -141,31 +202,42 @@ export function bundleSetup(options: SetupOptions, servers: Servers) { const { logsScriptUrl, rumScriptUrl, debuggerScriptUrl } = createCrossOriginScriptUrls(servers, options) + // Every SDK bundle is loaded before any init() call, so the pre-init script runs with all of + // them instrumenting the page. + let sdkScripts = '' + let initScripts = '' + if (options.logs) { - header += html`` - header += html`` + initScripts += html`` } if (options.rum) { - header += html`` - header += html`` + initScripts += html`` } if (options.debugger) { - header += html` - - - ` + sdkScripts += html`` + initScripts += html`` } + header += preInitDefinition(options) + sdkScripts + initScripts + return basePage({ header, body: options.body, @@ -209,6 +281,9 @@ export function npmSetup(options: SetupOptions, servers: Servers) { ` } + // The app bundle imports the SDK and then calls the *_INIT globals, so it is the one calling + // DD_PRE_INIT in between. + header += preInitDefinition(options) header += html`` return basePage({ @@ -218,6 +293,8 @@ export function npmSetup(options: SetupOptions, servers: Servers) { } export function appSetup(options: SetupOptions, servers: Servers, appName: string) { + rejectPreInitScript(options, 'app') + let header = options.head || '' if (options.eventBridge) { @@ -276,6 +353,8 @@ export function workerSetup(setupOptions: SetupOptions, servers: Servers) { } export function microfrontendSetup(options: SetupOptions, servers: Servers) { + rejectPreInitScript(options, 'microfrontend') + let header = options.head || '' if (options.eventBridge) { @@ -307,6 +386,8 @@ export function microfrontendSetup(options: SetupOptions, servers: Servers) { // Salesforce apps don't serve a locally-generated page body; this factory only drives the // page-side setup needed to init RUM on the remote Salesforce page. export async function salesforceSetup(options: SetupOptions, servers: Servers, page: Page): Promise { + rejectPreInitScript(options, 'salesforce') + const salesforceAppDirectory = options.salesforceApp === 'experience-cloud' ? 'sf-experience-app' : 'sf-lwc-app' const salesforceBundlePath = resolve( __dirname, From 0595f6d7fffdf055faf6e14423d681721517ab8e Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 14:20:53 +0200 Subject: [PATCH 7/7] =?UTF-8?q?=E2=9C=85=20Test=20WebSocket=20collection?= =?UTF-8?q?=20before=20init()=20end=20to=20end?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WebSocket early data collection had unit coverage on both sides of the buffer boundary, but nothing exercised the pre-start → post-start strategy swap in a browser. The headline behaviour — a socket opened before DD_RUM.init() still reports — was unverified end to end. Two scenarios open a socket in the pre-init window: one exchanges a message and stays open until the test closes it, the other also closes before init() so the entire connection lifecycle exists only in the buffer. Both assert the message counts and that the timings are measured from the real constructor call rather than from init(). Verified red: both fail against the collection code from before the buffered source was consumed. Co-Authored-By: Claude Opus 5 (1M context) --- test/e2e/lib/pages/webSocketPage.ts | 63 ++++++++++++ test/e2e/scenario/rum/websockets.scenario.ts | 100 ++++++++++++++++++- 2 files changed, 161 insertions(+), 2 deletions(-) diff --git a/test/e2e/lib/pages/webSocketPage.ts b/test/e2e/lib/pages/webSocketPage.ts index 96c55a4308..d9116da6e6 100644 --- a/test/e2e/lib/pages/webSocketPage.ts +++ b/test/e2e/lib/pages/webSocketPage.ts @@ -17,6 +17,25 @@ export function expectedWsEchoMessage(out = DEFAULT_WS_OUT_MESSAGE) { return `echo: ${out}` } +/** In-page expression evaluating to the /ws-echo URL for the current origin. */ +const WS_ECHO_URL = `(function () { + var url = new URL('/ws-echo', location.href) + url.protocol = location.protocol === 'https:' ? 'wss:' : 'ws:' + return url.toString() +})()` + +/** Handle exposed by {@link preInitWebSocketScript} to drive the socket it opened. */ +interface PreInitWebSocketHandle { + closed: boolean + close: () => void +} + +declare global { + interface Window { + preInitWebSocket?: PreInitWebSocketHandle + } +} + export class WebSocketPage { readonly wsOpenButton: Locator readonly wsStatusParagraph: Locator @@ -114,3 +133,47 @@ export class WebSocketPage { await expect(this.wsLastMessageParagraph).toHaveText(text) } } + +/** + * Script for `createTest().withPreInitScript()`: opens a socket to /ws-echo and exchanges a + * message before `init()` runs. It returns a promise, so the exchange — and, with + * `closeBeforeInit`, the close as well — is guaranteed to complete while the SDK is only + * buffering. The socket is left open otherwise, and can be closed later with + * {@link closePreInitWebSocket}. + * + * It has no DOM dependency: it runs in the page ``, before the body exists. + */ +export function preInitWebSocketScript({ closeBeforeInit = false } = {}) { + return ` + var socket = new WebSocket(${WS_ECHO_URL}) + var handle = { + closed: false, + close: function () { + socket.close() + }, + } + window.preInitWebSocket = handle + return new Promise(function (resolve) { + socket.addEventListener('open', function () { + socket.send(${JSON.stringify(DEFAULT_WS_OUT_MESSAGE)}) + }) + socket.addEventListener('message', function () { + ${closeBeforeInit ? 'socket.close()' : 'resolve()'} + }) + // Also covers a failed handshake: an 'error' is always followed by a 'close', so init() is + // never held back forever. + socket.addEventListener('close', function () { + handle.closed = true + resolve() + }) + }) + ` +} + +export async function closePreInitWebSocket(page: Page) { + // page.goto() only waits for the load event, which can happen before the SDK bundle is + // evaluated (async setup) and therefore before the pre-init script has run. + await page.waitForFunction(() => window.preInitWebSocket !== undefined) + await page.evaluate(() => window.preInitWebSocket!.close()) + await page.waitForFunction(() => window.preInitWebSocket!.closed) +} diff --git a/test/e2e/scenario/rum/websockets.scenario.ts b/test/e2e/scenario/rum/websockets.scenario.ts index df1f4a7a8d..fefb1e9f98 100644 --- a/test/e2e/scenario/rum/websockets.scenario.ts +++ b/test/e2e/scenario/rum/websockets.scenario.ts @@ -1,9 +1,22 @@ import type { RumResourceEvent } from '@datadog/browser-rum' -import type { RawRumEvent } from '@datadog/browser-rum-core' +import type { RawRumEvent, RumInitConfiguration } from '@datadog/browser-rum-core' +import type { Page } from '@playwright/test' import { expect, test } from '@playwright/test' import { createTest } from '../../lib/framework' import { expireSession, renewSession } from '../../lib/helpers/session' -import { DEFAULT_WS_OUT_MESSAGE, expectedWsEchoMessage, WebSocketPage } from '../../lib/pages/webSocketPage' +import { + closePreInitWebSocket, + DEFAULT_WS_OUT_MESSAGE, + expectedWsEchoMessage, + preInitWebSocketScript, + WebSocketPage, +} from '../../lib/pages/webSocketPage' + +declare global { + interface Window { + RUM_INIT_TIME?: number + } +} type RawRumResource = Extract type WebSocketResourceProperties = NonNullable @@ -195,6 +208,70 @@ test.describe('rum websockets', () => { expect(wsResources).toHaveLength(0) }) + createTest('collects a websocket opened and used before init()') + .withRum({ enableExperimentalFeatures: ['track_websockets'] }) + .withRumInit(recordInitTime) + .withPreInitScript(preInitWebSocketScript()) + .run(async ({ intakeRegistry, flushEvents, page }) => { + await closePreInitWebSocket(page) + + // Read before flushing: flushEvents() navigates away and drops the page state. + const initTime = await getInitTime(page) + + await flushEvents() + + const rumEvent = getLastRumResourceEventWithWebSocket(intakeRegistry.rumResourceEvents) + expect(rumEvent).toBeDefined() + + const { websocket } = rumEvent!.resource + + // The message was exchanged before init(): it is only counted if the SDK buffered it. + expect(websocket.messages_out.count).toBe(1) + expect(websocket.messages_out.size).toBe(DEFAULT_WS_OUT_MESSAGE.length) + expect(websocket.messages_in.count).toBe(1) + expect(websocket.messages_in.size).toBe(expectedWsEchoMessage().length) + + // Timings come from the constructor call, not from init(). + expect(websocket.start_time).toBeLessThan(initTime) + expect(websocket.handshake_succeeded).toBe(true) + expect(websocket.start_time + websocket.setup_duration / NANOSECONDS_PER_MILLISECOND).toBeLessThanOrEqual( + initTime + ) + + const connectingVital = intakeRegistry.rumVitalEvents.find((e) => e.vital.name === 'websocket-connecting') + expect(connectingVital).toBeDefined() + expect(connectingVital!.context!.connection_id).toBe(websocket.connection_id) + expect(connectingVital!.date).toBeLessThan(initTime) + }) + + createTest('collects a websocket whose whole lifecycle happened before init()') + .withRum({ enableExperimentalFeatures: ['track_websockets'] }) + .withRumInit(recordInitTime) + .withPreInitScript(preInitWebSocketScript({ closeBeforeInit: true })) + .run(async ({ intakeRegistry, flushEvents, page }) => { + // Read before flushing: flushEvents() navigates away and drops the page state. + const initTime = await getInitTime(page) + + await flushEvents() + + const rumEvent = getLastRumResourceEventWithWebSocket(intakeRegistry.rumResourceEvents) + expect(rumEvent).toBeDefined() + + const { websocket } = rumEvent!.resource + + expect(websocket.tracking_end_reason).toBe('close_event') + expect(websocket.messages_out.count).toBe(1) + expect(websocket.messages_in.count).toBe(1) + + // The connection was opened, used and closed while the SDK was only buffering. + expect(websocket.start_time).toBeLessThan(initTime) + expect(websocket.end_time).toBeLessThanOrEqual(initTime) + + const closedVital = intakeRegistry.rumVitalEvents.find((e) => e.vital.name === 'websocket-closed') + expect(closedVital).toBeDefined() + expect(closedVital!.context!.connection_id).toBe(websocket.connection_id) + }) + createTest('websocket resource records different start and end views when it spanned multiple views') .withRum({ enableExperimentalFeatures: ['track_websockets'] }) .withBody(WebSocketPage.testBody()) @@ -223,6 +300,25 @@ test.describe('rum websockets', () => { }) }) +/** + * Serialized into the page, so it must stay self-contained. Marks the moment init() runs, to + * assert that pre-init timings are measured from the real WebSocket constructor call. + */ +function recordInitTime(configuration: RumInitConfiguration) { + window.RUM_INIT_TIME = Date.now() + window.DD_RUM!.init(configuration) +} + +/** + * init() is held back until the pre-init WebSocket exchange completes, which can outlast the load + * event page.goto() waits for — so wait for the marker rather than assuming it is already set. + */ +async function getInitTime(page: Page) { + await page.waitForFunction(() => window.RUM_INIT_TIME !== undefined) + const initTime = await page.evaluate(() => window.RUM_INIT_TIME) + return initTime! +} + function isWebSocketResource(event: RumResourceEvent): event is RumResourceEventWithWebSocket { // Public RumResourceEvent.resource omits `websocket` until rum-events-format is updated. const resource = event.resource as unknown as {