From 2c96c649c9128f6457b510251501381b2dc190b4 Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Mon, 10 Aug 2026 19:48:06 +0200 Subject: [PATCH 1/6] =?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 4467063a8f..7e56927853 100644 --- a/packages/browser-core/test/index.ts +++ b/packages/browser-core/test/index.ts @@ -16,6 +16,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 b7c153cf95814259765ad8299596e30b91962864 Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Mon, 10 Aug 2026 21:44:45 +0200 Subject: [PATCH 2/6] =?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 e70d01f3fb8be6745b609c1181e52da89b890a1f Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 09:56:22 +0200 Subject: [PATCH 3/6] =?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 cf01f7702b3dd6109f7d392b095f823f7a5bc6ce Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 10:59:28 +0200 Subject: [PATCH 4/6] =?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 557bee042272b78e536212f573f5a83522cbac4d Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 11:23:38 +0200 Subject: [PATCH 5/6] =?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 fa87695c3c8ef8deb35dc714fc73b1fed69773e5 Mon Sep 17 00:00:00 2001 From: Boris Dibon Date: Tue, 11 Aug 2026 18:23:40 +0200 Subject: [PATCH 6/6] =?UTF-8?q?=E2=9A=A1=EF=B8=8F=20Select=20buffered=20da?= =?UTF-8?q?ta=20sources=20per=20SDK?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/domain/bufferedData.spec.ts | 55 ++++++++++++++++++- .../browser-core/src/domain/bufferedData.ts | 8 +-- .../src/boot/logsPublicApi.spec.ts | 19 +++++++ .../browser-logs/src/boot/logsPublicApi.ts | 8 ++- .../src/boot/rumPublicApi.spec.ts | 28 +++++++++- .../browser-rum-core/src/boot/rumPublicApi.ts | 9 ++- 6 files changed, 119 insertions(+), 8 deletions(-) diff --git a/packages/browser-core/src/domain/bufferedData.spec.ts b/packages/browser-core/src/domain/bufferedData.spec.ts index bda8a5eb78..1a275f632e 100644 --- a/packages/browser-core/src/domain/bufferedData.spec.ts +++ b/packages/browser-core/src/domain/bufferedData.spec.ts @@ -1,5 +1,5 @@ import { clocksNow } from '@datadog/js-core/time' -import { ConsoleApiName } from '@datadog/js-core/util' +import { ConsoleApiName, globalObject } from '@datadog/js-core/util' import type { MockFetch } from '../../test' import { collectAsyncCalls, @@ -23,6 +23,59 @@ import { ErrorHandling, ErrorSource, type RawError } from './error/error.types' import { trackRuntimeError } from './error/trackRuntimeError' describe('startBufferingData', () => { + it('does not instrument unselected sources', () => { + mockWebSocket() + const originalWebSocket = globalObject.WebSocket + const originalOnError = globalObject.onerror + const originalOnUnhandledRejection = globalObject.onunhandledrejection + const { stop } = startBufferingData([]) + + registerCleanupTask(() => { + stop() + resetWebSocketObservable() + }) + + expect(globalObject.WebSocket).toBe(originalWebSocket) + expect(globalObject.onerror).toBe(originalOnError) + expect(globalObject.onunhandledrejection).toBe(originalOnUnhandledRejection) + }) + + it('collects only selected data sources', (done) => { + const runtimeErrorObservable = new Observable() + replaceMockable(trackRuntimeError, () => runtimeErrorObservable) + mockWebSocket() + const originalWebSocket = globalObject.WebSocket + const { observable, stop } = startBufferingData([BufferedDataType.RUNTIME_ERROR]) + + registerCleanupTask(() => { + stop() + resetWebSocketObservable() + }) + + expect(globalObject.WebSocket).toBe(originalWebSocket) + + const rawError = { + startClocks: clocksNow(), + source: ErrorSource.SOURCE, + type: 'Error', + stack: 'Error: error!', + handling: ErrorHandling.UNHANDLED, + causes: undefined, + fingerprint: undefined, + message: 'error!', + } + + observable.subscribe((data) => { + expect(data).toEqual({ + type: BufferedDataType.RUNTIME_ERROR, + data: rawError, + }) + done() + }) + + runtimeErrorObservable.notify(rawError) + }) + it('collects runtime errors', (done) => { const runtimeErrorObservable = new Observable() replaceMockable(trackRuntimeError, () => runtimeErrorObservable) diff --git a/packages/browser-core/src/domain/bufferedData.ts b/packages/browser-core/src/domain/bufferedData.ts index 6ba31b4ce9..86647ebf7e 100644 --- a/packages/browser-core/src/domain/bufferedData.ts +++ b/packages/browser-core/src/domain/bufferedData.ts @@ -31,7 +31,7 @@ export type BufferedData = | { type: BufferedDataType.CONSOLE; data: ConsoleLog } | { type: BufferedDataType.WEB_SOCKET; data: WebSocketContext } -export function startBufferingData() { +export function startBufferingData(sources?: BufferedDataType[]) { const observable = new BufferedObservable(BUFFER_LIMIT, (count) => { // monitor-until: 2026-10-14 addTelemetryDebug('Early data collection dropped data on unbuffer', { @@ -44,6 +44,9 @@ export function startBufferingData() { type: T, source: Observable['data']> ) { + if (sources && !sources.includes(type)) { + return + } subscriptions.push( source.subscribe((data) => { observable.notify({ type, data } as BufferedData) @@ -55,9 +58,6 @@ 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 { diff --git a/packages/browser-logs/src/boot/logsPublicApi.spec.ts b/packages/browser-logs/src/boot/logsPublicApi.spec.ts index dba7c4c1cc..6ec168a321 100644 --- a/packages/browser-logs/src/boot/logsPublicApi.spec.ts +++ b/packages/browser-logs/src/boot/logsPublicApi.spec.ts @@ -6,6 +6,8 @@ import { TrackingConsent, startTelemetry, startSessionManager, + BufferedDataType, + startBufferingData, } from '@datadog/browser-core' import { collectAsyncCalls, @@ -13,6 +15,7 @@ import { replaceMockable, replaceMockableWithSpy, createStartSessionManagerMock, + registerCleanupTask, } from '@datadog/browser-core/test' import { HandlerType } from '../domain/logger' import { StatusType } from '../domain/logger/isAuthorized' @@ -27,6 +30,22 @@ const mockSessionId = 'some-session-id' const getInternalContext = () => ({ session_id: mockSessionId }) describe('logs entry', () => { + it('buffers the data sources consumed by Logs', () => { + const startBufferingDataSpy = replaceMockableWithSpy(startBufferingData).and.callFake(startBufferingData) + + makeLogsPublicApi() + + if (startBufferingDataSpy.calls.count()) { + registerCleanupTask(startBufferingDataSpy.calls.mostRecent().returnValue.stop) + } + expect(startBufferingDataSpy).toHaveBeenCalledOnceWith([ + BufferedDataType.RUNTIME_ERROR, + BufferedDataType.FETCH, + BufferedDataType.XHR, + BufferedDataType.CONSOLE, + ]) + }) + it('should add a `_setDebug` that works', () => { const displaySpy = spyOn(display, 'error') const { logsPublicApi } = makeLogsPublicApiWithDefaults() diff --git a/packages/browser-logs/src/boot/logsPublicApi.ts b/packages/browser-logs/src/boot/logsPublicApi.ts index cb48393c3b..885a974abf 100644 --- a/packages/browser-logs/src/boot/logsPublicApi.ts +++ b/packages/browser-logs/src/boot/logsPublicApi.ts @@ -12,6 +12,7 @@ import { startBufferingData, callMonitored, mockable, + BufferedDataType, } from '@datadog/browser-core' import { deepClone } from '@datadog/js-core/util' import type { LogsInitConfiguration } from '../domain/configuration' @@ -273,7 +274,12 @@ export interface LogsPublicApiOptions { export function makeLogsPublicApi(options: LogsPublicApiOptions = {}): LogsPublicApi { const trackingConsentState = createTrackingConsentState() - const bufferedDataObservable = startBufferingData().observable + const bufferedDataObservable = mockable(startBufferingData)([ + BufferedDataType.RUNTIME_ERROR, + BufferedDataType.FETCH, + BufferedDataType.XHR, + BufferedDataType.CONSOLE, + ]).observable let strategy = createPreStartStrategy( buildCommonContext, diff --git a/packages/browser-rum-core/src/boot/rumPublicApi.spec.ts b/packages/browser-rum-core/src/boot/rumPublicApi.spec.ts index 1f353ecc2d..39ca6a99ec 100644 --- a/packages/browser-rum-core/src/boot/rumPublicApi.spec.ts +++ b/packages/browser-rum-core/src/boot/rumPublicApi.spec.ts @@ -1,7 +1,15 @@ import { ONE_SECOND, timeStampToClocks, toTimeStamp } from '@datadog/js-core/time' import type { TimeStamp, RelativeTime } from '@datadog/js-core/time' import type { DeflateWorker } from '@datadog/browser-core' -import { display, DefaultPrivacyLevel, ResourceType, startTelemetry, startSessionManager } from '@datadog/browser-core' +import { + display, + DefaultPrivacyLevel, + ResourceType, + startTelemetry, + startSessionManager, + BufferedDataType, + startBufferingData, +} from '@datadog/browser-core' import type { Clock } from '@datadog/browser-core/test' import { collectAsyncCalls, @@ -11,6 +19,7 @@ import { replaceMockable, replaceMockableWithSpy, createStartSessionManagerMock, + registerCleanupTask, } from '@datadog/browser-core/test' import { noopRecorderApi, noopProfilerApi } from '../../test' import { ActionType, VitalType } from '../rawRumEvent.types' @@ -55,6 +64,23 @@ const DEFAULT_INIT_CONFIGURATION = { applicationId: 'xxx', clientToken: 'xxx' } const FAKE_WORKER = {} as DeflateWorker describe('rum public api', () => { + it('buffers the data sources consumed by RUM', () => { + const startBufferingDataSpy = replaceMockableWithSpy(startBufferingData).and.callFake(startBufferingData) + + makeRumPublicApiWithDefaults() + + if (startBufferingDataSpy.calls.count()) { + registerCleanupTask(startBufferingDataSpy.calls.mostRecent().returnValue.stop) + } + expect(startBufferingDataSpy).toHaveBeenCalledOnceWith([ + BufferedDataType.RUNTIME_ERROR, + BufferedDataType.FETCH, + BufferedDataType.XHR, + BufferedDataType.CONSOLE, + BufferedDataType.WEB_SOCKET, + ]) + }) + describe('init', () => { describe('deflate worker', () => { let rumPublicApi: RumPublicApi diff --git a/packages/browser-rum-core/src/boot/rumPublicApi.ts b/packages/browser-rum-core/src/boot/rumPublicApi.ts index 82e2282479..daef044f7d 100644 --- a/packages/browser-rum-core/src/boot/rumPublicApi.ts +++ b/packages/browser-rum-core/src/boot/rumPublicApi.ts @@ -33,6 +33,7 @@ import { startBufferingData, mockable, generateUUID, + BufferedDataType, } from '@datadog/browser-core' import type { LifeCycle } from '../domain/lifeCycle' @@ -669,7 +670,13 @@ export function makeRumPublicApi( options: RumPublicApiOptions = {} ): RumPublicApi { const trackingConsentState = createTrackingConsentState() - const bufferedData = startBufferingData() + const bufferedData = mockable(startBufferingData)([ + BufferedDataType.RUNTIME_ERROR, + BufferedDataType.FETCH, + BufferedDataType.XHR, + BufferedDataType.CONSOLE, + BufferedDataType.WEB_SOCKET, + ]) let strategy = createPreStartStrategy( options,