diff --git a/packages/browser-core/src/browser/webSocketObservable.spec.ts b/packages/browser-core/src/browser/webSocketObservable.spec.ts index 05ef1ccf96..5163d5e62c 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 @@ -184,11 +111,21 @@ 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', () => { 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 +138,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 +149,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 +160,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 +170,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 +179,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 +188,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 +199,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 +213,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 +221,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 +229,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 +239,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 +259,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/src/domain/bufferedData.spec.ts b/packages/browser-core/src/domain/bufferedData.spec.ts index d7238297a6..1a275f632e 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 { ConsoleApiName, globalObject } 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' @@ -13,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) @@ -136,6 +199,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..86647ebf7e 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,8 +29,9 @@ 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() { +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', { @@ -40,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) @@ -51,6 +58,7 @@ export function startBufferingData() { subscribe(BufferedDataType.FETCH, initFetchObservable()) subscribe(BufferedDataType.XHR, initXhrObservable()) subscribe(BufferedDataType.CONSOLE, initConsoleObservable(Object.values(ConsoleApiName))) + subscribe(BufferedDataType.WEB_SOCKET, initWebSocketObservable()) return { observable, diff --git a/packages/browser-core/test/emulate/mockWebSocket.ts b/packages/browser-core/test/emulate/mockWebSocket.ts new file mode 100644 index 0000000000..7512363679 --- /dev/null +++ b/packages/browser-core/test/emulate/mockWebSocket.ts @@ -0,0 +1,93 @@ +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 + // 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() + this.url = resolveWebSocketUrl(String(url)) + if (typeof protocols === 'string') { + this.protocol = protocols + } + } + + send(data: string | ArrayBufferLike | Blob | ArrayBufferView): void { + this.sentData.push(data) + } + + 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' 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, 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 c73773d9fc..ce41053ee8 100644 --- a/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts +++ b/packages/browser-rum-core/src/domain/resource/webSocketCollection.spec.ts @@ -1,9 +1,26 @@ -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 type { BufferedData } from '@datadog/browser-core' +import { + addExperimentalFeatures, + BufferedDataType, + ExperimentalFeature, + initWebSocketObservable, + Observable, + resetAllowUntrustedEvents, + setAllowUntrustedEvents, + startBufferingData, +} from '@datadog/browser-core' +import { + collectAsyncCalls, + 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 { mockRumConfiguration, mockViewHistory } from '../../../test' import { VitalType } from '../../rawRumEvent.types' import type { ViewHistoryEntry } from '../contexts/viewHistory' import { LifeCycle, LifeCycleEventType } from '../lifeCycle' @@ -20,17 +37,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) }) @@ -40,81 +56,77 @@ 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() ) { - return trackWebSocket(lifeCycle, wsObservable, viewHistory, addDurationVital) + const tracker = trackWebSocket(lifeCycle, createWebSocketDataObservable(), 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 +145,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 +169,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 +240,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 +254,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 +281,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 +308,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 +325,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 +344,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 +368,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 +378,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 +394,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 +405,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 +430,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 +454,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 +472,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 +484,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,54 +525,147 @@ 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()) + function startCollection( + configuration = mockRumConfiguration({ betaTrackWebSockets: true }), + bufferedDataObservable = createWebSocketDataObservable() + ) { + const collection = startWebSocketCollection( + lifeCycle, + configuration, + mockViewHistory(), + jasmine.createSpy(), + bufferedDataObservable + ) registerCleanupTask(() => collection.stop()) return collection } - function notifyConnecting(offsetMs = 0) { - singletonObservable().notify({ - state: 'connecting', - instance: wsInstance, - url: wsUrl, - startClocks: relativeToClocks(clock.relative(offsetMs)), + 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) + }) }) - } - function notifyOpen(openRelative = 10, protocol = '') { - singletonObservable().notify({ - state: 'open', - instance: wsInstance, - openClocks: relativeToClocks(clock.relative(openRelative)), - protocol, + 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() }) - } + }) - function notifyMessageOut(at: number, size: number, bufferedAmountPreSend = 0) { - singletonObservable().notify({ - state: 'message-out', - instance: wsInstance, - size, - bufferedAmountPreSend, - at: relativeToClocks(clock.relative(at)), + // 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 notifyClosed(at: number, code: number, reason: string, wasClean: boolean) { - singletonObservable().notify({ - state: 'closed', - instance: wsInstance, - code, - reason, - wasClean, - at: relativeToClocks(clock.relative(at)), + 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('finalizes open connections with tracking_end_reason="session_end" when the session expires', () => { const endClocks = relativeToClocks(clock.relative(40)) @@ -590,21 +685,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 +707,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') 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()