diff --git a/packages/core/src/domain/session/sessionManager.spec.ts b/packages/core/src/domain/session/sessionManager.spec.ts index d709e3b54f..2cfd2d39c9 100644 --- a/packages/core/src/domain/session/sessionManager.spec.ts +++ b/packages/core/src/domain/session/sessionManager.spec.ts @@ -665,7 +665,7 @@ describe('startSessionManager', () => { const sessionManager = startSessionManagerWithDefaults() sessionManager.sessionStateUpdateObservable.subscribe(sessionStateUpdateSpy) - sessionManager.updateSessionState({ extra: 'extra' }) + sessionManager.updateSessionState(() => ({ extra: 'extra' })) expectSessionIdToBeDefined(sessionManager) expect(sessionStateUpdateSpy).toHaveBeenCalledTimes(1) diff --git a/packages/core/src/domain/session/sessionManager.ts b/packages/core/src/domain/session/sessionManager.ts index 03caf2f49c..fae90339a5 100644 --- a/packages/core/src/domain/session/sessionManager.ts +++ b/packages/core/src/domain/session/sessionManager.ts @@ -20,13 +20,19 @@ export interface SessionManager { expireObservable: Observable sessionStateUpdateObservable: Observable<{ previousState: SessionState; newState: SessionState }> expire: () => void - updateSessionState: (state: Partial) => void + updateSessionState: (update: (state: SessionState) => Partial | undefined) => void } export interface SessionContext extends Context { id: string trackingType: TrackingType isReplayForced: boolean + /** + * Whether an error has already been reported during this session. Persisted in the session store + * so it survives page navigation: an error session must not go back to withholding its replay + * just because the user moved to another page. + */ + hasError: boolean anonymousId: string | undefined } @@ -92,6 +98,7 @@ export function startSessionManager( id: sessionStore.getSession().id!, trackingType: sessionStore.getSession()[productKey] as TrackingType, isReplayForced: !!sessionStore.getSession().forcedReplay, + hasError: !!sessionStore.getSession().hasError, anonymousId: sessionStore.getSession().anonymousId, } } diff --git a/packages/core/src/domain/session/sessionStore.spec.ts b/packages/core/src/domain/session/sessionStore.spec.ts index 7d8105177c..2bcd9378f7 100644 --- a/packages/core/src/domain/session/sessionStore.spec.ts +++ b/packages/core/src/domain/session/sessionStore.spec.ts @@ -596,7 +596,7 @@ describe('session store', () => { sessionStoreManager = setupSessionStore(updateSpy) otherSessionStoreManager = setupSessionStore(otherUpdateSpy) - sessionStoreManager.updateSessionState({ extra: 'extra' }) + sessionStoreManager.updateSessionState(() => ({ extra: 'extra' })) expect(updateSpy).toHaveBeenCalledTimes(1) diff --git a/packages/core/src/domain/session/sessionStore.ts b/packages/core/src/domain/session/sessionStore.ts index c926e79410..cb2253a9e6 100644 --- a/packages/core/src/domain/session/sessionStore.ts +++ b/packages/core/src/domain/session/sessionStore.ts @@ -29,7 +29,12 @@ export interface SessionStore { sessionStateUpdateObservable: Observable<{ previousState: SessionState; newState: SessionState }> expire: () => void stop: () => void - updateSessionState: (state: Partial) => void + /** + * Applies a change to the stored session under the store lock. The producer sees the state the + * change would land on and returns `undefined` to make it a no-op - which is how a write meant for + * one session avoids landing on the one that replaced it while the write was waiting for the lock. + */ + updateSessionState: (update: (state: SessionState) => Partial | undefined) => void } /** @@ -216,10 +221,13 @@ export function startSessionStore( renewObservable.notify() } - function updateSessionState(partialSessionState: Partial) { + function updateSessionState(update: (state: SessionState) => Partial | undefined) { processSessionStoreOperations( { - process: (sessionState) => ({ ...sessionState, ...partialSessionState }), + process: (sessionState) => { + const partialSessionState = update(sessionState) + return partialSessionState && { ...sessionState, ...partialSessionState } + }, after: synchronizeSession, }, sessionStoreStrategy diff --git a/packages/core/src/domain/session/sessionStoreOperations.spec.ts b/packages/core/src/domain/session/sessionStoreOperations.spec.ts index 78d76b6efe..37a8ab374f 100644 --- a/packages/core/src/domain/session/sessionStoreOperations.spec.ts +++ b/packages/core/src/domain/session/sessionStoreOperations.spec.ts @@ -1,3 +1,4 @@ +import { startFakeTelemetry } from '../telemetry' import type { MockStorage } from '../../../test' import { mockClock, mockCookie, mockLocalStorage } from '../../../test' import type { CookieOptions } from '../../browser/cookie' @@ -232,6 +233,7 @@ const DEFAULT_INIT_CONFIGURATION = { trackAnonymousUser: true } as Configuration it('should abort after a max number of retry', () => { const clock = mockClock() + const telemetry = startFakeTelemetry() sessionStoreStrategy.persistSession(initialSession) storage.setSpy.calls.reset() @@ -246,6 +248,7 @@ const DEFAULT_INIT_CONFIGURATION = { trackAnonymousUser: true } as Configuration expect(processSpy).not.toHaveBeenCalled() expect(afterSpy).not.toHaveBeenCalled() expect(storage.setSpy).not.toHaveBeenCalled() + expect(telemetry).toContain(jasmine.objectContaining({ message: 'Session store lock retries exhausted' })) clock.cleanup() }) diff --git a/packages/core/src/domain/session/sessionStoreOperations.ts b/packages/core/src/domain/session/sessionStoreOperations.ts index 869347d0cc..5be4d4cc5f 100644 --- a/packages/core/src/domain/session/sessionStoreOperations.ts +++ b/packages/core/src/domain/session/sessionStoreOperations.ts @@ -1,3 +1,4 @@ +import { addTelemetryDebug } from '../telemetry' import { setTimeout } from '../../tools/timer' import { generateUUID } from '../../tools/utils/stringUtils' import type { SessionStoreStrategy } from './storeStrategies/sessionStoreStrategy' @@ -43,6 +44,7 @@ export function processSessionStoreOperations( return } if (isLockEnabled && numberOfRetries >= LOCK_MAX_TRIES) { + addTelemetryDebug('Session store lock retries exhausted', { retries: numberOfRetries }) next(sessionStoreStrategy) return } diff --git a/packages/core/src/domain/telemetry/telemetry.ts b/packages/core/src/domain/telemetry/telemetry.ts index 4b0a12c617..4651ebb253 100644 --- a/packages/core/src/domain/telemetry/telemetry.ts +++ b/packages/core/src/domain/telemetry/telemetry.ts @@ -4,7 +4,9 @@ import { NO_ERROR_STACK_PRESENT_MESSAGE, isError } from '../error/error' import { toStackTraceString } from '../../tools/stackTrace/handlingStack' import { getExperimentalFeatures } from '../../tools/experimentalFeatures' import type { Configuration } from '../configuration' -import { INTAKE_SITE_STAGING } from '../configuration' +// Import the constant without loading configuration construction, which uses session storage. +// eslint-disable-next-line local-rules/disallow-protected-directory-import +import { INTAKE_SITE_STAGING } from '../configuration/intakeSites' import { Observable } from '../../tools/observable' import { timeStampNow } from '../../tools/utils/timeUtils' import { displayIfDebugEnabled, startMonitorErrorCollection } from '../../tools/monitor' diff --git a/packages/rum-core/src/boot/startRum.ts b/packages/rum-core/src/boot/startRum.ts index 06299236f4..f874605b70 100644 --- a/packages/rum-core/src/boot/startRum.ts +++ b/packages/rum-core/src/boot/startRum.ts @@ -29,6 +29,7 @@ import { startErrorCollection } from '../domain/error/errorCollection' import { startResourceCollection } from '../domain/resource/resourceCollection' import { startViewCollection } from '../domain/view/viewCollection' import { startRumSessionManager, startRumSessionManagerStub } from '../domain/rumSessionManager' +import { startSessionErrorTracking } from '../domain/trackSessionError' import { startRumBatch } from '../transport/startRumBatch' import { startRumEventBridge } from '../transport/startRumEventBridge' import { startUrlContexts } from '../domain/contexts/urlContexts' @@ -121,6 +122,9 @@ export function startRum( : startRumSessionManager(configuration, lifeCycle, trackingConsentState) cleanupTasks.push(session.stop) + const sessionErrorTracking = startSessionErrorTracking(lifeCycle, session) + cleanupTasks.push(() => sessionErrorTracking.stop()) + if (!canUseEventBridge()) { // FLASHCAT FORK - keep the console's sampling rates fresh, at the rhythm the sessions read // them: once now and once per session renewal. It is skipped under an event bridge, where the diff --git a/packages/rum-core/src/domain/configuration/configuration.spec.ts b/packages/rum-core/src/domain/configuration/configuration.spec.ts index 64069cd281..20f89eb359 100644 --- a/packages/rum-core/src/domain/configuration/configuration.spec.ts +++ b/packages/rum-core/src/domain/configuration/configuration.spec.ts @@ -65,6 +65,81 @@ describe('validateAndBuildRumConfiguration', () => { }) }) + describe('sessionReplayOnError', () => { + it('is carried into the built configuration', () => { + expect( + validateAndBuildRumConfiguration({ ...DEFAULT_INIT_CONFIGURATION, sessionReplayOnError: true })! + .sessionReplayOnError + ).toBeTrue() + }) + + it('defaults to collecting no error replay at all', () => { + expect(validateAndBuildRumConfiguration(DEFAULT_INIT_CONFIGURATION)!.sessionReplayOnError).toBeFalse() + }) + + it('is read as a switch, whatever it was given', () => { + expect( + validateAndBuildRumConfiguration({ + ...DEFAULT_INIT_CONFIGURATION, + sessionReplayOnError: 1 as unknown as boolean, + })!.sessionReplayOnError + ).toBeTrue() + }) + + it('starts the recording on its own, since there is nothing to withhold otherwise', () => { + expect( + validateAndBuildRumConfiguration({ + ...DEFAULT_INIT_CONFIGURATION, + sessionReplaySampleRate: 0, + sessionReplayOnError: true, + })!.startSessionReplayRecordingManually + ).toBeFalse() + }) + + it('warns when the plain replay rate leaves it nothing to apply to', () => { + validateAndBuildRumConfiguration({ + ...DEFAULT_INIT_CONFIGURATION, + sessionReplaySampleRate: 100, + sessionReplayOnError: true, + }) + + expect(displayWarnSpy).toHaveBeenCalledTimes(1) + expect(displayWarnSpy.calls.argsFor(0)[0]).toContain('sessionReplaySampleRate did not draw') + }) + + it('warns when no session is tracked at all', () => { + validateAndBuildRumConfiguration({ + ...DEFAULT_INIT_CONFIGURATION, + sessionSampleRate: 0, + sessionReplayOnError: true, + }) + + expect(displayWarnSpy).toHaveBeenCalledTimes(1) + expect(displayWarnSpy.calls.argsFor(0)[0]).toContain('no session is tracked') + }) + + it('warns when the recording is left for the customer to start, since nothing would be held', () => { + validateAndBuildRumConfiguration({ + ...DEFAULT_INIT_CONFIGURATION, + sessionReplayOnError: true, + startSessionReplayRecordingManually: true, + }) + + expect(displayWarnSpy).toHaveBeenCalledTimes(1) + expect(displayWarnSpy.calls.argsFor(0)[0]).toContain('startSessionReplayRecordingManually') + }) + + it('says nothing about a switch that can apply', () => { + validateAndBuildRumConfiguration({ + ...DEFAULT_INIT_CONFIGURATION, + sessionReplaySampleRate: 20, + sessionReplayOnError: true, + }) + + expect(displayWarnSpy).not.toHaveBeenCalled() + }) + }) + describe('traceSampleRate', () => { it('defaults to 100 if the option is not provided', () => { expect(validateAndBuildRumConfiguration(DEFAULT_INIT_CONFIGURATION)!.traceSampleRate).toBe(100) @@ -283,6 +358,26 @@ describe('validateAndBuildRumConfiguration', () => { }) describe('startSessionReplayRecordingManually', () => { + it('keeps automatic recording available for remotely enabled replay', () => { + expect( + validateAndBuildRumConfiguration({ + ...DEFAULT_INIT_CONFIGURATION, + sessionReplaySampleRate: 0, + remoteConfigurationEnabled: true, + })!.startSessionReplayRecordingManually + ).toBeFalse() + }) + + it('respects explicit manual recording when remote configuration is enabled', () => { + expect( + validateAndBuildRumConfiguration({ + ...DEFAULT_INIT_CONFIGURATION, + remoteConfigurationEnabled: true, + startSessionReplayRecordingManually: true, + })!.startSessionReplayRecordingManually + ).toBeTrue() + }) + it('defaults to true if sessionReplaySampleRate is 0', () => { expect( validateAndBuildRumConfiguration({ ...DEFAULT_INIT_CONFIGURATION, sessionReplaySampleRate: 0 })! @@ -554,6 +649,7 @@ describe('serializeRumConfiguration', () => { enablePrivacyForActionName: false, subdomain: 'foo', sessionReplaySampleRate: 60, + sessionReplayOnError: true, startSessionReplayRecordingManually: true, sessionReplayDirectUpload: true, trackUserInteractions: true, @@ -587,6 +683,8 @@ describe('serializeRumConfiguration', () => { // FLASHCAT FORK: not reported to telemetry | 'sessionReplayDirectUpload' | 'beforeSampling' + // not reported yet: needs a rum-events-format schema change first + | 'sessionReplayOnError' ? never : CamelToSnakeCase // By specifying the type here, we can ensure that serializeConfiguration is returning an diff --git a/packages/rum-core/src/domain/configuration/configuration.ts b/packages/rum-core/src/domain/configuration/configuration.ts index 9ea4ca9ea7..0467f75a5f 100644 --- a/packages/rum-core/src/domain/configuration/configuration.ts +++ b/packages/rum-core/src/domain/configuration/configuration.ts @@ -200,6 +200,15 @@ export interface RumInitConfiguration extends InitConfiguration { * See [Configure Your Setup For Browser RUM and Browser RUM & Session Replay Sampling](https://docs.datadoghq.com/real_user_monitoring/guide/sampling-browser-plans) for further information. */ sessionReplaySampleRate?: number | undefined + /** + * Whether the tracked sessions that `sessionReplaySampleRate` did not draw still record a replay, + * uploaded only if the session reports an error. Default: false. + * + * Such a session records from the start and keeps at most the last minute of it in memory. If it + * never reports an error, nothing is uploaded and the session is not billed. On the first error, + * the withheld minute is uploaded and recording continues normally for the rest of the session. + */ + sessionReplayOnError?: boolean | undefined /** * If the session is sampled for Session Replay, only start the recording when `startSessionReplayRecording()` is called, instead of at the beginning of the session. Default: if startSessionReplayRecording is 0, true; otherwise, false. * See [Session Replay Usage](https://docs.datadoghq.com/real_user_monitoring/session_replay/browser/#usage) for further information. @@ -301,6 +310,7 @@ export interface RumConfiguration extends Configuration { defaultPrivacyLevel: DefaultPrivacyLevel enablePrivacyForActionName: boolean sessionReplaySampleRate: number + sessionReplayOnError: boolean startSessionReplayRecordingManually: boolean sessionReplayDirectUpload: boolean trackUserInteractions: boolean @@ -385,16 +395,40 @@ export function validateAndBuildRumConfiguration( const profilingEnabled = isExperimentalFeatureEnabled(ExperimentalFeature.PROFILING) const sessionReplaySampleRate = initConfiguration.sessionReplaySampleRate ?? 0 + const sessionReplayOnError = !!initConfiguration.sessionReplayOnError + + // Each of these is a combination the customer can set that cannot apply to a single session. It + // is valid, so validation lets it through - but silence would leave them waiting for data that is + // never coming. + if (sessionReplayOnError) { + if (sessionReplaySampleRate === 100) { + display.warn( + 'sessionReplayOnError only applies to sessions sessionReplaySampleRate did not draw, and that rate is 100: it will never apply.' + ) + } + if ((initConfiguration.sessionSampleRate ?? 100) === 0) { + display.warn('sessionReplayOnError has no effect while sessionSampleRate is 0: no session is tracked.') + } + if (initConfiguration.startSessionReplayRecordingManually) { + display.warn( + 'sessionReplayOnError needs the recording to already be running when the error happens, and startSessionReplayRecordingManually keeps it stopped until you start it: there would be nothing to release.' + ) + } + } return { applicationId: initConfiguration.applicationId, version: initConfiguration.version || undefined, actionNameAttribute: initConfiguration.actionNameAttribute, sessionReplaySampleRate, + sessionReplayOnError, startSessionReplayRecordingManually: initConfiguration.startSessionReplayRecordingManually !== undefined ? !!initConfiguration.startSessionReplayRecordingManually - : sessionReplaySampleRate === 0, + : // An error-sampled session has to be recording before the error happens, otherwise there is + // nothing to withhold and release. Remote configuration may enable replay on a later + // session, so keep the automatic start intent even when init disables replay. + sessionReplaySampleRate === 0 && !sessionReplayOnError && !initConfiguration.remoteConfigurationEnabled, sessionReplayDirectUpload: !!initConfiguration.sessionReplayDirectUpload, traceSampleRate: initConfiguration.traceSampleRate ?? 100, rulePsr: isNumber(initConfiguration.traceSampleRate) ? initConfiguration.traceSampleRate / 100 : undefined, @@ -485,6 +519,9 @@ export function serializeRumConfiguration(configuration: RumInitConfiguration) { return { session_replay_sample_rate: configuration.sessionReplaySampleRate, + // `session_replay_on_error` is deliberately not reported yet: the telemetry + // configuration type is generated from the rum-events-format schema, so adding it needs a schema + // change first, and that is a separate repository. start_session_replay_recording_manually: configuration.startSessionReplayRecordingManually, trace_sample_rate: configuration.traceSampleRate, trace_context_injection: configuration.traceContextInjection, diff --git a/packages/rum-core/src/domain/configuration/remoteConfiguration.spec.ts b/packages/rum-core/src/domain/configuration/remoteConfiguration.spec.ts index c3e409ab8c..1bc0b3fd24 100644 --- a/packages/rum-core/src/domain/configuration/remoteConfiguration.spec.ts +++ b/packages/rum-core/src/domain/configuration/remoteConfiguration.spec.ts @@ -121,6 +121,30 @@ describe('remoteConfiguration', () => { start(configurationWith()) }) + it('keeps the replay-on-error switch the server reports, either way it is set', (done) => { + interceptor.withMockXhr((xhr) => { + xhr.complete(200, body({ rum: { sessionReplaySampleRate: 10, sessionReplayOnError: false } })) + + expect(readRemoteConfig(setup)).toEqual({ + sessionReplaySampleRate: 10, + sessionReplayOnError: false, + version: 3, + }) + done() + }) + start(configurationWith()) + }) + + it('drops a switch that is not a boolean, so it reads as not delivered', (done) => { + interceptor.withMockXhr((xhr) => { + xhr.complete(200, body({ rum: { sessionSampleRate: 50, sessionReplayOnError: 'true' as unknown as boolean } })) + + expect(readRemoteConfig(setup)).toEqual({ sessionSampleRate: 50, version: 3 }) + done() + }) + start(configurationWith()) + }) + it('drops a privacy level it does not recognise rather than passing it on', (done) => { // A typo must not reach the recorders: an unknown value there falls through to recording // everything, which is the one outcome nobody asks for by accident. diff --git a/packages/rum-core/src/domain/configuration/remoteConfiguration.ts b/packages/rum-core/src/domain/configuration/remoteConfiguration.ts index 2c0ee781f9..24be706f18 100644 --- a/packages/rum-core/src/domain/configuration/remoteConfiguration.ts +++ b/packages/rum-core/src/domain/configuration/remoteConfiguration.ts @@ -102,6 +102,12 @@ export interface RemoteConfigValues { * fact. */ defaultPrivacyLevel?: DefaultPrivacyLevel + /** + * Whether the sessions `sessionReplaySampleRate` did not draw still record a replay, uploaded only + * if the session errors. Read at the draw like the rates, and for the same reason: a session + * either withholds its replay from the start or never does. + */ + sessionReplayOnError?: boolean /** * Which version of the settings these rates came from. Reported back on the next request so the * console can say how far a change has actually reached — a question the events cannot answer, @@ -261,6 +267,9 @@ function readStoredValues(parsed: unknown): RemoteConfigValues { if (isPrivacyLevel(stored.defaultPrivacyLevel)) { values.defaultPrivacyLevel = stored.defaultPrivacyLevel } + if (isSwitch(stored.sessionReplayOnError)) { + values.sessionReplayOnError = stored.sessionReplayOnError + } if (isBag(stored.custom)) { values.custom = stored.custom } @@ -484,6 +493,11 @@ function store(setup: RemoteConfigSetup, response: RemoteConfigurationResponse) if (isPrivacyLevel(response.rum.defaultPrivacyLevel)) { values.defaultPrivacyLevel = response.rum.defaultPrivacyLevel } + // A switch is a boolean or nothing. Anything else - a "true" string, a 1 - is dropped for the + // same reason a bad rate is: it must read as "not delivered", not as either position. + if (isSwitch(response.rum.sessionReplayOnError)) { + values.sessionReplayOnError = response.rum.sessionReplayOnError + } } // The custom bag rides along untouched — the platform's job is delivery, its meaning belongs to // the host application. Gone from the response (or the kill switch off) means gone from storage. @@ -732,6 +746,10 @@ export function isRate(value: unknown): value is number { return typeof value === 'number' && value >= 0 && value <= 100 } +export function isSwitch(value: unknown): value is boolean { + return typeof value === 'boolean' +} + /** * A version is a publish counter, so anything that is not a whole, non-negative number small enough * to survive a JSON round trip cannot be one. Checked on the way in and on the way out, because a diff --git a/packages/rum-core/src/domain/contexts/sessionContext.spec.ts b/packages/rum-core/src/domain/contexts/sessionContext.spec.ts index 263bb737c5..2d73f0d58f 100644 --- a/packages/rum-core/src/domain/contexts/sessionContext.spec.ts +++ b/packages/rum-core/src/domain/contexts/sessionContext.spec.ts @@ -94,6 +94,49 @@ describe('session context', () => { expect(eventWithoutHasReplay.session!.has_replay).toBeUndefined() }) + it('should tell a replay kept only because the session errored apart from an unconditional one', () => { + sessionManager.setTrackedWithErrorSessionReplay() + const errorReplayEvent = hooks.triggerHook(HookNames.Assemble, { + eventType: 'view', + startTime: 0 as RelativeTime, + }) as DefaultRumEventAttributes + + sessionManager.setTrackedWithSessionReplay() + const plainEvent = hooks.triggerHook(HookNames.Assemble, { + eventType: 'view', + startTime: 0 as RelativeTime, + }) as DefaultRumEventAttributes + + expect(errorReplayEvent.session!.sampled_for_error_replay).toBeTrue() + // absent rather than false, so it costs nothing on every ordinary session + expect(plainEvent.session!.sampled_for_error_replay).toBeUndefined() + }) + + it('should not set hasReplay when a dropped buffer left the view with nothing', () => { + // a withheld buffer that was dropped rolls back what it held, and a view left with an empty + // stats entry has no replay to offer + getReplayStatsSpy.and.returnValue({ segments_count: 0, records_count: 0, segments_total_raw_size: 0 }) + + const event = hooks.triggerHook(HookNames.Assemble, { + eventType: 'view', + startTime: 0 as RelativeTime, + }) as DefaultRumEventAttributes + + expect(event.session!.has_replay).toBeUndefined() + }) + + it('should set hasReplay when a host bridge took the records and no segment was built', () => { + // records go straight to the bridge, so nothing ever counts a segment for them + getReplayStatsSpy.and.returnValue({ ...fakeStats, segments_count: 0, records_count: 10 }) + + const event = hooks.triggerHook(HookNames.Assemble, { + eventType: 'view', + startTime: 0 as RelativeTime, + }) as DefaultRumEventAttributes + + expect(event.session!.has_replay).toBe(true) + }) + it('should set session.is_active when the session is active', () => { findViewSpy.and.returnValue({ ...fakeView, sessionIsActive: true }) const eventWithActiveSession = hooks.triggerHook(HookNames.Assemble, { diff --git a/packages/rum-core/src/domain/contexts/sessionContext.ts b/packages/rum-core/src/domain/contexts/sessionContext.ts index ced3bebbee..b928ac3acd 100644 --- a/packages/rum-core/src/domain/contexts/sessionContext.ts +++ b/packages/rum-core/src/domain/contexts/sessionContext.ts @@ -38,15 +38,28 @@ export function startSessionContext( return DISCARDED } + // A session withholding its replay is recording, but nothing has been uploaded and nothing may + // ever be. Reporting `has_replay` here would offer a replay that does not exist. + const isReplayWithheld = session.sessionReplay === SessionReplayState.BUFFERED_ON_ERROR + let hasReplay let sampledForReplay + let sampledForErrorReplay let isActive if (eventType === RumEventType.VIEW) { - hasReplay = recorderApi.getReplayStats(view.id) ? true : undefined + // Records rather than merely a stats entry: a withheld buffer that was dropped rolls back what + // it held, which leaves a view with an empty stats entry and no replay at all - and offering a + // replay that was never uploaded is worse than not offering one. Records, not segments, + // because a host bridge takes the records itself and no segment is ever built for them. + const replayStats = recorderApi.getReplayStats(view.id) + hasReplay = !isReplayWithheld && replayStats && replayStats.records_count > 0 ? true : undefined sampledForReplay = session.sessionReplay === SessionReplayState.SAMPLED + // Tells a replay collected only because the session errored apart from one collected + // unconditionally - the two cost differently and are answered by different questions. + sampledForErrorReplay = session.sampledOnErrorReplay || undefined isActive = view.sessionIsActive ? undefined : false } else { - hasReplay = recorderApi.isRecording() ? true : undefined + hasReplay = !isReplayWithheld && recorderApi.isRecording() ? true : undefined } return { @@ -56,6 +69,7 @@ export function startSessionContext( type: SessionType.USER, has_replay: hasReplay, sampled_for_replay: sampledForReplay, + sampled_for_error_replay: sampledForErrorReplay, is_active: isActive, }, // FLASHCAT FORK - overrides the init values reported by the default context with the rates diff --git a/packages/rum-core/src/domain/lifeCycle.ts b/packages/rum-core/src/domain/lifeCycle.ts index 78b6d9fccb..9f65158973 100644 --- a/packages/rum-core/src/domain/lifeCycle.ts +++ b/packages/rum-core/src/domain/lifeCycle.ts @@ -49,6 +49,8 @@ export const enum LifeCycleEventType { // at the end leaves upstream's numbering alone and keeps this file out of the way of the next // upstream merge. REMOTE_CONFIGURATION_STORED, + /** A local or shared session mark has released conditional collection. */ + SESSION_RELEASED, } // This is a workaround for an issue occurring when the Browser SDK is included in a TypeScript @@ -81,6 +83,7 @@ declare const LifeCycleEventTypeAsConst: { RAW_RUM_EVENT_COLLECTED: LifeCycleEventType.RAW_RUM_EVENT_COLLECTED RUM_EVENT_COLLECTED: LifeCycleEventType.RUM_EVENT_COLLECTED RAW_ERROR_COLLECTED: LifeCycleEventType.RAW_ERROR_COLLECTED + SESSION_RELEASED: LifeCycleEventType.SESSION_RELEASED REMOTE_CONFIGURATION_STORED: LifeCycleEventType.REMOTE_CONFIGURATION_STORED } @@ -106,6 +109,7 @@ export interface LifeCycleEventMap { error: RawError customerContext?: Context } + [LifeCycleEventTypeAsConst.SESSION_RELEASED]: { sessionId: string; reason: 'error' | 'force' } [LifeCycleEventTypeAsConst.REMOTE_CONFIGURATION_STORED]: void } diff --git a/packages/rum-core/src/domain/rumSessionManager.spec.ts b/packages/rum-core/src/domain/rumSessionManager.spec.ts index 40878e2db0..0fd2bf83f0 100644 --- a/packages/rum-core/src/domain/rumSessionManager.spec.ts +++ b/packages/rum-core/src/domain/rumSessionManager.spec.ts @@ -6,6 +6,7 @@ import { setCookie, stopSessionManager, ONE_SECOND, + isChromium, DOM_EVENT, createTrackingConsentState, TrackingConsent, @@ -228,6 +229,7 @@ describe('rum session manager', () => { sessionReplaySampleRate?: number traceSampleRate?: number defaultPrivacyLevel?: string + sessionReplayOnError?: boolean }) { localStorage.setItem(STORE_KEY, JSON.stringify(values)) registerCleanupTask(() => localStorage.removeItem(STORE_KEY)) @@ -255,6 +257,40 @@ describe('rum session manager', () => { expect(getSessionState(SESSION_STORE_KEY)[RUM_SESSION_KEY]).toBe(RumTrackingType.TRACKED_WITH_SESSION_REPLAY) }) + it('keeps a replay on error when the console says so, over what init said', () => { + storeRemoteConfigValues({ sessionReplaySampleRate: 0, sessionReplayOnError: true }) + + startRumSessionManagerWithDefaults({ + configuration: { + sessionSampleRate: 100, + sessionReplaySampleRate: 100, + sessionReplayOnError: false, + remoteConfig: REMOTE_SAMPLING_SETUP, + }, + }) + document.dispatchEvent(createNewEvent(DOM_EVENT.CLICK)) + + expect(getSessionState(SESSION_STORE_KEY)[RUM_SESSION_KEY]).toBe( + RumTrackingType.TRACKED_WITH_ERROR_SESSION_REPLAY + ) + }) + + it('turns the replay-on-error switch off when the console says so', () => { + storeRemoteConfigValues({ sessionReplayOnError: false }) + + startRumSessionManagerWithDefaults({ + configuration: { + sessionSampleRate: 100, + sessionReplaySampleRate: 0, + sessionReplayOnError: true, + remoteConfig: REMOTE_SAMPLING_SETUP, + }, + }) + document.dispatchEvent(createNewEvent(DOM_EVENT.CLICK)) + + expect(getSessionState(SESSION_STORE_KEY)[RUM_SESSION_KEY]).toBe(RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY) + }) + it('falls back to the rate passed to init for a knob the console did not set', () => { storeRemoteConfigValues({ sessionReplaySampleRate: 100 }) @@ -1195,6 +1231,194 @@ describe('rum session manager', () => { }) }) + describe('session replay on error', () => { + for (const mark of ['error', 'force'] as const) { + for (const replacement of [false, true]) { + it(`reconciles ${mark} after lock exhaustion only for its original session (replacement=${replacement})`, () => { + if (!isChromium()) { + pending('requires a cookie store lock') + } + const manager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + const id = manager.findTrackedSession()!.id + const state = `id=${id}&rum=3&created=${Date.now()}&expire=${Date.now() + DURATION}` + setCookie(SESSION_STORE_KEY, `lock=other-tab&${state}`, DURATION) + if (mark === 'error') { + manager.setSessionHasError(id) + } else { + manager.setForcedReplay() + } + clock.tick(1500) + setCookie(SESSION_STORE_KEY, replacement ? state.replace(id, 'replacement') : state, DURATION) + clock.tick(3000) + expect(getSessionState(SESSION_STORE_KEY)[mark === 'error' ? 'hasError' : 'forcedReplay']).toBe( + replacement ? undefined : '1' + ) + }) + } + } + + for (const force of ['setForcedReplay', 'setForcedSession'] as const) { + it(`${force} releases in memory before a locked store can persist it`, () => { + if (!isChromium()) { + pending('requires a cookie store lock') + } + const manager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + const id = manager.findTrackedSession()!.id + setCookie( + SESSION_STORE_KEY, + `lock=other-tab&id=${id}&rum=3&created=${Date.now()}&expire=${Date.now() + DURATION}`, + DURATION + ) + manager[force]() + expect(manager.findTrackedSession()!.sessionReplay).toBe(SessionReplayState.FORCED) + expect(getSessionState(SESSION_STORE_KEY).forcedReplay).toBeUndefined() + }) + + it(`${force} never writes its deferred mark into a replacement session`, () => { + if (!isChromium()) { + pending('requires a cookie store lock') + } + const manager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + const id = manager.findTrackedSession()!.id + setCookie( + SESSION_STORE_KEY, + `lock=other-tab&id=${id}&rum=3&created=${Date.now()}&expire=${Date.now() + DURATION}`, + DURATION + ) + manager[force]() + setCookie( + SESSION_STORE_KEY, + `id=replacement&rum=3&created=${Date.now()}&expire=${Date.now() + DURATION}`, + DURATION + ) + clock.tick(20) + expect(getSessionState(SESSION_STORE_KEY).forcedReplay).toBeUndefined() + }) + } + + it('applies the error-replay type only when the plain replay draw missed', () => { + startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 100, sessionReplayOnError: true }, + }) + + expect(getSessionState(SESSION_STORE_KEY)[RUM_SESSION_KEY]).toBe(RumTrackingType.TRACKED_WITH_SESSION_REPLAY) + }) + + it('stores the error-replay type when only the switch applies', () => { + startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + + expect(getSessionState(SESSION_STORE_KEY)[RUM_SESSION_KEY]).toBe( + RumTrackingType.TRACKED_WITH_ERROR_SESSION_REPLAY + ) + }) + + it('withholds the replay until the session reports an error', () => { + const sessionManager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + + expect(sessionManager.findTrackedSession()!.sessionReplay).toBe(SessionReplayState.BUFFERED_ON_ERROR) + + sessionManager.setSessionHasError(sessionManager.findTrackedSession()!.id) + + expect(sessionManager.findTrackedSession()!.sessionReplay).toBe(SessionReplayState.SAMPLED) + }) + + it('does not mark a session that has since been replaced by another one', () => { + const sessionManager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + + // another tab renewed the session while the mark was on its way to the store + setCookie(SESSION_STORE_KEY, 'id=other-session&rum=3', DURATION) + + sessionManager.setSessionHasError('a-session-that-is-gone') + + expect(getSessionState(SESSION_STORE_KEY).hasError).toBeUndefined() + }) + + it('releases the replay before the store write lands, since that write can be deferred', () => { + if (!isChromium()) { + pending('the store lock, and so a deferred write, only exists on Chromium') + } + const sessionManager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + const sessionId = sessionManager.findTrackedSession()!.id + + // another tab holds the store lock, so the write is deferred through retries + setCookie(SESSION_STORE_KEY, `lock=other-tab&id=${sessionId}&rum=3`, DURATION) + + sessionManager.setSessionHasError(sessionId) + + expect(getSessionState(SESSION_STORE_KEY).hasError).toBeUndefined() + // and yet the buffer must already see it as released: the page or the session may end before + // the write ever lands, and the buffer would otherwise be thrown away + expect(sessionManager.findTrackedSession()!.sessionReplay).toBe(SessionReplayState.SAMPLED) + }) + + it('keeps the released state across a page load, since it is persisted in the session store', () => { + setCookie( + SESSION_STORE_KEY, + `id=abcdef&rum=3&hasError=1&created=${Date.now()}&expire=${Date.now() + DURATION}`, + DURATION + ) + + const sessionManager = startRumSessionManagerWithDefaults() + + expect(sessionManager.findTrackedSession()!.sessionReplay).toBe(SessionReplayState.SAMPLED) + }) + + it('marks the session so a replay kept only because it errored can be told apart', () => { + const sessionManager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + + expect(sessionManager.findTrackedSession()!.sampledOnErrorReplay).toBeTrue() + + // still true once released, so what was stored can be told apart afterwards + sessionManager.setSessionHasError(sessionManager.findTrackedSession()!.id) + + expect(sessionManager.findTrackedSession()!.sampledOnErrorReplay).toBeTrue() + }) + + it('does not mark a session whose replay is collected unconditionally', () => { + const sessionManager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 100 }, + }) + + expect(sessionManager.findTrackedSession()!.sampledOnErrorReplay).toBeFalse() + }) + + it('releases the replay when it is forced, rather than waiting for an error that may never come', () => { + const sessionManager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: true }, + }) + expect(sessionManager.findTrackedSession()!.sessionReplay).toBe(SessionReplayState.BUFFERED_ON_ERROR) + + sessionManager.setForcedReplay() + + expect(sessionManager.findTrackedSession()!.sessionReplay).toBe(SessionReplayState.FORCED) + }) + + it('tracks the session even when neither the replay rate nor the switch applies', () => { + const sessionManager = startRumSessionManagerWithDefaults({ + configuration: { sessionSampleRate: 100, sessionReplaySampleRate: 0, sessionReplayOnError: false }, + }) + + expect(getSessionState(SESSION_STORE_KEY)[RUM_SESSION_KEY]).toBe(RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY) + expect(sessionManager.findTrackedSession()!.sessionReplay).toBe(SessionReplayState.OFF) + }) + }) + function startRumSessionManagerWithDefaults({ configuration, trackingConsentState = createTrackingConsentState(TrackingConsent.GRANTED), diff --git a/packages/rum-core/src/domain/rumSessionManager.ts b/packages/rum-core/src/domain/rumSessionManager.ts index 9e742043ee..898b4ac6c7 100644 --- a/packages/rum-core/src/domain/rumSessionManager.ts +++ b/packages/rum-core/src/domain/rumSessionManager.ts @@ -36,6 +36,12 @@ export interface RumSessionManager { expireObservable: Observable setForcedReplay: () => void setForcedSession: () => void + /** + * Marks the given session as having reported an error. For a session sampled by + * `sessionReplayOnError`, this is what releases the withheld replay. The id is required + * because the store write can be deferred by the lock, and it must not land on a later session. + */ + setSessionHasError: (sessionId: string) => void } /** @@ -80,6 +86,12 @@ export interface DrawnConfiguration { export type RumSession = { id: string sessionReplay: SessionReplayState + /** + * Whether the replay of this session is only kept if it reports an error. Unlike + * {@link sessionReplay} this stays true once the error has been reported, so a replay collected + * that way can be told apart from one collected unconditionally. + */ + sampledOnErrorReplay: boolean anonymousId?: string // FLASHCAT FORK - absent when the draw used exactly what init passed — nothing to override then, // the events already report those values — and when the record of the draw did not survive @@ -91,12 +103,19 @@ export const enum RumTrackingType { NOT_TRACKED = '0', TRACKED_WITH_SESSION_REPLAY = '1', TRACKED_WITHOUT_SESSION_REPLAY = '2', + TRACKED_WITH_ERROR_SESSION_REPLAY = '3', } export const enum SessionReplayState { OFF, SAMPLED, FORCED, + /** + * The session records, but every segment is withheld until it reports its first error. If no error + * ever happens, nothing is uploaded and the session is never billed. Once an error is reported the + * session moves to `SAMPLED` and the withheld buffer is released. + */ + BUFFERED_ON_ERROR, } export function startRumSessionManager( @@ -294,12 +313,45 @@ export function startRumSessionManager( endSessionIfSettingsAreDecisive ) - sessionManager.sessionStateUpdateObservable.subscribe(({ previousState, newState }) => { - if (!previousState.forcedReplay && newState.forcedReplay) { - const sessionEntity = sessionManager.findSession() - if (sessionEntity) { - sessionEntity.isReplayForced = true - } + function forceReplay() { + const session = sessionManager.findSession() + if (!session) { + return + } + const wasForced = session.isReplayForced + session.isReplayForced = true + if (!wasForced) { + lifeCycle.notify(LifeCycleEventType.SESSION_RELEASED, { sessionId: session.id, reason: 'force' }) + } + sessionManager.updateSessionState((state) => (state.id === session.id ? { forcedReplay: '1' } : undefined)) + } + + const sessionStateSubscription = sessionManager.sessionStateUpdateObservable.subscribe(({ newState }) => { + const session = sessionManager.findSession() + if (!session || session.id !== newState.id) { + return + } + const becameForced = !session.isReplayForced && newState.forcedReplay === '1' + const becameErrored = !session.hasError && newState.hasError === '1' + session.isReplayForced ||= becameForced + session.hasError ||= becameErrored + if (becameForced || becameErrored) { + lifeCycle.notify(LifeCycleEventType.SESSION_RELEASED, { + sessionId: session.id, + reason: becameForced ? 'force' : 'error', + }) + } + // A lock retry can be exhausted before a mark reaches storage. The existing poll is the + // next opportunity to reconcile it, and the session identity bounds how long it may live. + if ((session.hasError && newState.hasError !== '1') || (session.isReplayForced && newState.forcedReplay !== '1')) { + sessionManager.updateSessionState((state) => + state.id === session.id + ? { + ...(session.hasError ? { hasError: '1' } : {}), + ...(session.isReplayForced ? { forcedReplay: '1' } : {}), + } + : undefined + ) } }) return { @@ -310,12 +362,8 @@ export function startRumSessionManager( } return { id: session.id, - sessionReplay: - session.trackingType === RumTrackingType.TRACKED_WITH_SESSION_REPLAY - ? SessionReplayState.SAMPLED - : session.isReplayForced - ? SessionReplayState.FORCED - : SessionReplayState.OFF, + sessionReplay: computeSessionReplayState(session.trackingType, session.hasError, session.isReplayForced), + sampledOnErrorReplay: withholdsReplay(session.trackingType), anonymousId: session.anonymousId, // FLASHCAT FORK - looked up at the same time as the session itself, so an event that // belongs to a session already renewed still reports the draw that created it. @@ -325,29 +373,75 @@ export function startRumSessionManager( expire: sessionManager.expire, expireObservable: sessionManager.expireObservable, stop: () => { + sessionStateSubscription.unsubscribe() consentSubscription.unsubscribe() remoteConfigSubscription.unsubscribe() drawnHistory.stop() }, - setForcedReplay: () => sessionManager.updateSessionState({ forcedReplay: '1' }), + setForcedReplay: forceReplay, // FLASHCAT FORK - the escape hatch for "collect this visitor NOW": the host application knows // who needs debugging (its own allow-list, a support flow), the SDK only provides the switch. // A session keeps the decision it was drawn with, so forcing a visitor that was not being // collected means ending their current (empty) session; the next activity draws again with // `forcedSession` set and starts a collected session with replay. A session already collected - // only needs replay forced on, which is the existing forced-replay path. + // only needs replay forced on, which is the existing forced-replay path - and a session whose + // replay is withheld until it errors is released the same way, since the host asked for it now. setForcedSession: () => { forcedSession = true const session = sessionManager.findSession() if (!session || !isTypeTracked(session.trackingType)) { sessionManager.expire() - } else if (session.trackingType === RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY) { - sessionManager.updateSessionState({ forcedReplay: '1' }) + } else if ( + session.trackingType === RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY || + withholdsReplay(session.trackingType) + ) { + forceReplay() + } + }, + setSessionHasError: (sessionId) => { + const sessionEntity = sessionManager.findSession() + if (sessionEntity?.id === sessionId) { + // Marked in memory straight away, and not only once the store write lands: that write goes + // through a lock that can defer it by several retries, and until then the withheld buffer + // would still read the session as withholding - so an error followed closely by the page or + // the session ending would throw away the very buffer the error was meant to release. + const hadError = sessionEntity.hasError + sessionEntity.hasError = true + if (!hadError) { + lifeCycle.notify(LifeCycleEventType.SESSION_RELEASED, { sessionId, reason: 'error' }) + } } + sessionManager.updateSessionState((state) => (state.id === sessionId ? { hasError: '1' } : undefined)) }, } } +export function withholdsReplay(trackingType: RumTrackingType) { + return trackingType === RumTrackingType.TRACKED_WITH_ERROR_SESSION_REPLAY +} + +export function computeSessionReplayState( + trackingType: RumTrackingType, + hasError: boolean, + isReplayForced: boolean +): SessionReplayState { + if (trackingType === RumTrackingType.TRACKED_WITH_SESSION_REPLAY) { + return SessionReplayState.SAMPLED + } + if (withholdsReplay(trackingType) && hasError) { + return SessionReplayState.SAMPLED + } + // A forced replay wins over withholding: the host explicitly asked for this user's replay, so it + // must not keep waiting for an error that may never come. + if (isReplayForced) { + return SessionReplayState.FORCED + } + if (withholdsReplay(trackingType)) { + return SessionReplayState.BUFFERED_ON_ERROR + } + return SessionReplayState.OFF +} + /** * Session id used when the host application does not answer for one, because it was built against * an SDK that predates `getSessionId()`. Such a host is expected to override the session id of the @@ -439,6 +533,8 @@ export function startRumSessionManagerStub( return { id: sessionId ?? STUB_SESSION_ID, sessionReplay, + // The host records for us, or this page uploads what the plain rate drew: neither withholds. + sampledOnErrorReplay: false, anonymousId: bridge?.getAnonymousId(), } }, @@ -446,6 +542,7 @@ export function startRumSessionManagerStub( expireObservable, setForcedReplay: noop, setForcedSession: noop, + setSessionHasError: noop, stop: () => clearInterval(watchIntervalId), } } @@ -476,16 +573,22 @@ function computeSessionState( // the decision it was created with: settings arriving mid-session never start or stop // collecting for a visitor already on the site. const remote = readRemoteConfig(configuration.remoteConfig) - const { sessionSampleRate, sessionReplaySampleRate } = resolveSampleRates(configuration, remote) + const { sessionSampleRate, sessionReplaySampleRate, sessionReplayOnError } = resolveSampleRates( + configuration, + remote + ) reportDraw(configuration, remote, sessionSampleRate, sessionReplaySampleRate, onDraw) if (!performDraw(sessionSampleRate)) { trackingType = RumTrackingType.NOT_TRACKED - } else if (!performDraw(sessionReplaySampleRate)) { - trackingType = RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY - } else { + } else if (performDraw(sessionReplaySampleRate)) { trackingType = RumTrackingType.TRACKED_WITH_SESSION_REPLAY + } else if (sessionReplayOnError) { + // Only for sessions the plain replay draw missed, so a session is never counted by both. + trackingType = RumTrackingType.TRACKED_WITH_ERROR_SESSION_REPLAY + } else { + trackingType = RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY } } return { @@ -495,8 +598,10 @@ function computeSessionState( } /** - * FLASHCAT FORK - the rates a draw would use right now: what the console delivered, falling back to - * what the site passed to init, with the application's `beforeSampling` given the last word. This + * FLASHCAT FORK - the rates a draw would use right now, and the on-error switch beside them: what + * the console delivered, falling back to what the site passed to init, with the application's + * `beforeSampling` given the last word on the rates (the switch is not offered to it: it is a + * yes or a no the console already answered). This * is what turns the delivered custom values into sampling decisions without a wasted first draw or * a session restart: the console ships the data (an allow-list, a cohort rule), the application's * own code interprets it here. Its failure modes must never reach session creation, so a thrown @@ -530,7 +635,11 @@ function resolveSampleRates(configuration: RumConfiguration, remote: RemoteConfi } } - return { sessionSampleRate, sessionReplaySampleRate } + return { + sessionSampleRate, + sessionReplaySampleRate, + sessionReplayOnError: remote.sessionReplayOnError ?? configuration.sessionReplayOnError, + } } /** @@ -674,13 +783,15 @@ function hasValidRumSession(trackingType?: string): trackingType is RumTrackingT return ( trackingType === RumTrackingType.NOT_TRACKED || trackingType === RumTrackingType.TRACKED_WITH_SESSION_REPLAY || - trackingType === RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY + trackingType === RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY || + trackingType === RumTrackingType.TRACKED_WITH_ERROR_SESSION_REPLAY ) } function isTypeTracked(rumSessionType: RumTrackingType | undefined) { return ( rumSessionType === RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY || - rumSessionType === RumTrackingType.TRACKED_WITH_SESSION_REPLAY + rumSessionType === RumTrackingType.TRACKED_WITH_SESSION_REPLAY || + rumSessionType === RumTrackingType.TRACKED_WITH_ERROR_SESSION_REPLAY ) } diff --git a/packages/rum-core/src/domain/trackSessionError.spec.ts b/packages/rum-core/src/domain/trackSessionError.spec.ts new file mode 100644 index 0000000000..a0057d01f7 --- /dev/null +++ b/packages/rum-core/src/domain/trackSessionError.spec.ts @@ -0,0 +1,115 @@ +import type { Context } from '@flashcatcloud/browser-core' +import { registerCleanupTask } from '@flashcatcloud/browser-core/test' +import type { RumEvent } from '../rumEvent.types' +import { createRumSessionManagerMock } from '../../test' +import { LifeCycle, LifeCycleEventType } from './lifeCycle' +import { startSessionErrorTracking } from './trackSessionError' + +describe('startSessionErrorTracking', () => { + let lifeCycle: LifeCycle + let sessionManager: ReturnType + let setSessionHasErrorSpy: jasmine.Spy + + function collect(type: string, source = 'source') { + // only error events carry an `error` object; anything else that did would hide a guard that + // reads it before checking the type + const event = type === 'error' ? { type, session: { id: 'session-id' }, error: { source } } : { type } + lifeCycle.notify(LifeCycleEventType.RUM_EVENT_COLLECTED, event as unknown as RumEvent & Context) + } + + beforeEach(() => { + lifeCycle = new LifeCycle() + sessionManager = createRumSessionManagerMock().setTrackedWithErrorSessionReplay() + setSessionHasErrorSpy = spyOn(sessionManager, 'setSessionHasError').and.callThrough() + const { stop } = startSessionErrorTracking(lifeCycle, sessionManager) + registerCleanupTask(stop) + }) + + it('ignores an error from an earlier session without consuming the current session mark', () => { + lifeCycle.notify(LifeCycleEventType.RUM_EVENT_COLLECTED, { + type: 'error', + session: { id: 'previous-session' }, + error: { source: 'custom' }, + } as unknown as RumEvent & Context) + expect(setSessionHasErrorSpy).not.toHaveBeenCalled() + collect('error') + expect(setSessionHasErrorSpy).toHaveBeenCalledOnceWith('session-id') + }) + + it('does not attribute an error without a session id to the current session', () => { + lifeCycle.notify(LifeCycleEventType.RUM_EVENT_COLLECTED, { + type: 'error', + error: { source: 'custom' }, + } as unknown as RumEvent & Context) + expect(setSessionHasErrorSpy).not.toHaveBeenCalled() + }) + + it('marks the session on the first collected error', () => { + collect('error') + + // named, not just counted: the mark is refused if it does not name the session it belongs to + expect(setSessionHasErrorSpy).toHaveBeenCalledOnceWith('session-id') + }) + + it('leaves a session that withholds nothing alone, so an ordinary session store is never written', () => { + sessionManager.setTrackedWithSessionReplay() + + collect('error') + + expect(setSessionHasErrorSpy).not.toHaveBeenCalled() + }) + + it('leaves an untracked session alone', () => { + sessionManager.setNotTracked() + + collect('error') + + expect(setSessionHasErrorSpy).not.toHaveBeenCalled() + }) + + it('does not mark the session on other event types', () => { + collect('view') + collect('resource') + collect('action') + + expect(setSessionHasErrorSpy).not.toHaveBeenCalled() + }) + + it('ignores the SDK own failures, which are not the application reporting an error', () => { + collect('error', 'agent') + + expect(setSessionHasErrorSpy).not.toHaveBeenCalled() + }) + + it('still marks the session on a network error, which is the application reporting one', () => { + collect('error', 'network') + + expect(setSessionHasErrorSpy).toHaveBeenCalledTimes(1) + }) + + it('marks the session only once, however many errors follow', () => { + collect('error') + collect('error') + collect('error') + + expect(setSessionHasErrorSpy).toHaveBeenCalledTimes(1) + }) + + it('marks a renewed session again, since it is a different session', () => { + collect('error') + lifeCycle.notify(LifeCycleEventType.SESSION_RENEWED) + collect('error') + + expect(setSessionHasErrorSpy).toHaveBeenCalledTimes(2) + }) + + it('stops marking once stopped', () => { + const { stop } = startSessionErrorTracking(lifeCycle, sessionManager) + stop() + setSessionHasErrorSpy.calls.reset() + // the suite's own tracker is still running, so exactly one call is expected, not two + collect('error') + + expect(setSessionHasErrorSpy).toHaveBeenCalledTimes(1) + }) +}) diff --git a/packages/rum-core/src/domain/trackSessionError.ts b/packages/rum-core/src/domain/trackSessionError.ts new file mode 100644 index 0000000000..3414744bf9 --- /dev/null +++ b/packages/rum-core/src/domain/trackSessionError.ts @@ -0,0 +1,52 @@ +import { ErrorSource } from '@flashcatcloud/browser-core' +import { RumEventType } from '../rawRumEvent.types' +import type { LifeCycle } from './lifeCycle' +import { LifeCycleEventType } from './lifeCycle' +import type { RumSessionManager } from './rumSessionManager' + +/** + * Marks the session as having reported an error, which is what releases a replay withheld by + * `sessionReplayOnError`. + * + * It listens after assembly rather than on the raw error, so an error discarded by `beforeSend` or + * by a rate limiter does not release anything: a session billed for an error that cannot be found + * afterwards would be worse than no replay at all. + */ +export function startSessionErrorTracking(lifeCycle: LifeCycle, sessionManager: RumSessionManager) { + let hasReportedError = false + + const eventSubscription = lifeCycle.subscribe(LifeCycleEventType.RUM_EVENT_COLLECTED, (event) => { + if (hasReportedError || event.type !== RumEventType.ERROR) { + return + } + // The SDK's own failures — an intake request that could not be sent, for instance — are ours, + // not the application's. Counting them would turn every session into an error session for any + // customer whose network blocks our endpoint, billing them for replays of nothing. + if (event.error.source === ErrorSource.AGENT) { + return + } + // Only a session that is withholding something has any use for this mark. Setting it on any + // other session would write the session store for customers who enabled neither rate - and that + // write also pushes the session's expiry out (`processSessionStoreOperations` expands every + // state it persists), which would move where their sessions end. + const session = sessionManager.findTrackedSession() + if (!session?.sampledOnErrorReplay || event.session?.id !== session.id) { + return + } + hasReportedError = true + sessionManager.setSessionHasError(session.id) + }) + + // A renewed session is a different session: it draws its own sampling and starts out without an + // error, so anything withheld for it must stay withheld until it reports one of its own. + const renewSubscription = lifeCycle.subscribe(LifeCycleEventType.SESSION_RENEWED, () => { + hasReportedError = false + }) + + return { + stop: () => { + eventSubscription.unsubscribe() + renewSubscription.unsubscribe() + }, + } +} diff --git a/packages/rum-core/test/mockRumSessionManager.ts b/packages/rum-core/test/mockRumSessionManager.ts index 6bf0a1294a..80c9d90002 100644 --- a/packages/rum-core/test/mockRumSessionManager.ts +++ b/packages/rum-core/test/mockRumSessionManager.ts @@ -1,12 +1,20 @@ import { Observable } from '@flashcatcloud/browser-core' -import { SessionReplayState, type DrawnConfiguration, type RumSessionManager } from '../src/domain/rumSessionManager' +import { + RumTrackingType, + computeSessionReplayState, + withholdsReplay, + type DrawnConfiguration, + type RumSessionManager, +} from '../src/domain/rumSessionManager' export interface RumSessionManagerMock extends RumSessionManager { setId(id: string): RumSessionManagerMock setNotTracked(): RumSessionManagerMock setTrackedWithoutSessionReplay(): RumSessionManagerMock setTrackedWithSessionReplay(): RumSessionManagerMock + setTrackedWithErrorSessionReplay(): RumSessionManagerMock setForcedReplay(): RumSessionManagerMock + setSessionHasError(): RumSessionManagerMock setDrawnConfiguration(drawn: DrawnConfiguration): RumSessionManagerMock } @@ -14,31 +22,34 @@ const DEFAULT_ID = 'session-id' const enum SessionStatus { TRACKED_WITH_SESSION_REPLAY, TRACKED_WITHOUT_SESSION_REPLAY, + TRACKED_WITH_ERROR_SESSION_REPLAY, NOT_TRACKED, EXPIRED, } +const TRACKING_TYPES: { [key in SessionStatus]?: RumTrackingType } = { + [SessionStatus.TRACKED_WITH_SESSION_REPLAY]: RumTrackingType.TRACKED_WITH_SESSION_REPLAY, + [SessionStatus.TRACKED_WITHOUT_SESSION_REPLAY]: RumTrackingType.TRACKED_WITHOUT_SESSION_REPLAY, + [SessionStatus.TRACKED_WITH_ERROR_SESSION_REPLAY]: RumTrackingType.TRACKED_WITH_ERROR_SESSION_REPLAY, +} + export function createRumSessionManagerMock(): RumSessionManagerMock { let id = DEFAULT_ID let sessionStatus: SessionStatus = SessionStatus.TRACKED_WITH_SESSION_REPLAY let forcedReplay: boolean = false + let hasError: boolean = false let drawnConfiguration: DrawnConfiguration | undefined return { findTrackedSession() { - if ( - sessionStatus !== SessionStatus.TRACKED_WITH_SESSION_REPLAY && - sessionStatus !== SessionStatus.TRACKED_WITHOUT_SESSION_REPLAY - ) { + const trackingType = TRACKING_TYPES[sessionStatus] + if (!trackingType) { return undefined } return { id, - sessionReplay: - sessionStatus === SessionStatus.TRACKED_WITH_SESSION_REPLAY - ? SessionReplayState.SAMPLED - : forcedReplay - ? SessionReplayState.FORCED - : SessionReplayState.OFF, + // Derived the same way as in production, so the mock cannot drift from the real state machine + sessionReplay: computeSessionReplayState(trackingType, hasError, forcedReplay), + sampledOnErrorReplay: withholdsReplay(trackingType), anonymousId: 'device-123', drawnConfiguration, } @@ -64,10 +75,18 @@ export function createRumSessionManagerMock(): RumSessionManagerMock { sessionStatus = SessionStatus.TRACKED_WITH_SESSION_REPLAY return this }, + setTrackedWithErrorSessionReplay() { + sessionStatus = SessionStatus.TRACKED_WITH_ERROR_SESSION_REPLAY + return this + }, setForcedReplay() { forcedReplay = true return this }, + setSessionHasError() { + hasError = true + return this + }, setDrawnConfiguration(drawn) { drawnConfiguration = drawn return this diff --git a/packages/rum-legacy/src/domain/sessionStore.spec.ts b/packages/rum-legacy/src/domain/sessionStore.spec.ts index fdeddb47a4..9f88cdf96b 100644 --- a/packages/rum-legacy/src/domain/sessionStore.spec.ts +++ b/packages/rum-legacy/src/domain/sessionStore.spec.ts @@ -29,6 +29,31 @@ describe('session store', () => { deleteSessionCookie() }) + for (const [rum, flag, tracked] of [ + ['3', '', true], + ['4', '', false], + ['5', '', false], + ['4', '&hasError=1', true], + ['5', '&hasError=1', true], + ['4', '&forcedReplay=1', true], + ['5', '&forcedReplay=1', true], + ['0', '&hasError=1', false], + ] as const) { + it(`respects the shared tracking decision rum=${rum}${flag}`, () => { + document.cookie = `${SESSION_COOKIE_NAME}=id=shared-session&rum=${rum}${flag}&created=${Date.now()}&expire=${Date.now() + ONE_MINUTE};path=/` + expect(createSessionStore(100).getOrCreateSession().isTracked).toBe(tracked) + expect(toSessionState(readRawCookie()).rum).toBe(rum) + }) + } + + it('does not carry release marks into a renewed legacy session', () => { + document.cookie = `${SESSION_COOKIE_NAME}=id=old-session&rum=5&hasError=1&forcedReplay=1&created=${Date.now() - ONE_MINUTE}&expire=${Date.now() - 1};path=/` + createSessionStore(100).getOrCreateSession() + const stored = toSessionState(readRawCookie()) + expect(stored.hasError).toBeUndefined() + expect(stored.forcedReplay).toBeUndefined() + }) + it('creates a session with a lowercase uuid', () => { const session = createSessionStore(100).getOrCreateSession() diff --git a/packages/rum-legacy/src/domain/sessionStore.ts b/packages/rum-legacy/src/domain/sessionStore.ts index 60a870a3a7..28c49e701e 100644 --- a/packages/rum-legacy/src/domain/sessionStore.ts +++ b/packages/rum-legacy/src/domain/sessionStore.ts @@ -30,6 +30,9 @@ const EXPIRED = '1' const NOT_TRACKED = '0' const TRACKED_WITH_SESSION_REPLAY = '1' const TRACKED_WITHOUT_SESSION_REPLAY = '2' +const TRACKED_WITH_ERROR_SESSION_REPLAY = '3' +const TRACKED_ON_ERROR_WITHOUT_SESSION_REPLAY = '4' +const TRACKED_ON_ERROR_WITH_SESSION_REPLAY = '5' /** * How long a session may be reused without touching the cookie again. @@ -133,7 +136,7 @@ export function createSessionStore(sessionSampleRate: number) { } function toSession(state: SessionState): LegacySession { - // Both tracked values count. This build never writes '1' itself, but both builds share one cookie + // Honor collected and released decisions. This build writes only '0'/'2', but shares one cookie // jar per domain, and IE enterprise site lists routinely put some urls of a site in compatibility // mode and others not. Reading a session the modern bundle started as untracked would silence // this one for the rest of that session's lifetime. @@ -144,7 +147,13 @@ function toSession(state: SessionState): LegacySession { } function isTracked(state: SessionState): boolean { - return state.rum === TRACKED_WITHOUT_SESSION_REPLAY || state.rum === TRACKED_WITH_SESSION_REPLAY + return ( + state.rum === TRACKED_WITHOUT_SESSION_REPLAY || + state.rum === TRACKED_WITH_SESSION_REPLAY || + state.rum === TRACKED_WITH_ERROR_SESSION_REPLAY || + ((state.rum === TRACKED_ON_ERROR_WITHOUT_SESSION_REPLAY || state.rum === TRACKED_ON_ERROR_WITH_SESSION_REPLAY) && + (state.hasError === '1' || state.forcedReplay === '1')) + ) } /** @@ -193,7 +202,7 @@ function isExpired(state: SessionState, now: number): boolean { // `isExpired` belongs to the modern bundle's vocabulary, not to ours, but it has to be listed here // all the same: carried forward as an unknown field it would mark every session this build writes // as expired, and the modern bundle would start a new one on every page load. -const KNOWN_FIELDS = ['id', 'created', 'expire', 'rum', 'isExpired'] +const KNOWN_FIELDS = ['id', 'created', 'expire', 'rum', 'isExpired', 'hasError', 'forcedReplay'] function serialize(state: SessionState): string { const entries: string[] = [] diff --git a/packages/rum/README.md b/packages/rum/README.md index 11198510c8..e5b56ba897 100644 --- a/packages/rum/README.md +++ b/packages/rum/README.md @@ -32,3 +32,35 @@ flashcatRum.init({ [1]: https://docs.flashcat.cloud/zh/flashduty/rum/introduction [2]: https://www.npmjs.com/package/@flashcatcloud/browser-rum + +## Enabling error session collection across pages + +`sessionReplayOnError` needs the full `browser-rum` bundle. The slim and legacy +bundles do not contain a recorder. `sessionOnError` also requires a bundle with +conditional event buffering; the legacy bundle can only honor a shared session +that has already been released by a compatible modern page. + +Before enabling either option in initialization or remote configuration: + +1. Deploy compatible SDK bundles to every page sharing the session cookie, + including other applications and subdomains when cross-subdomain tracking is + enabled. Keep both error-collection options disabled during this deployment. +2. Account for already-open pages and cached application assets. Publishing a new + SDK does not replace JavaScript in those pages. Require those pages to reload, + or defer enablement until incompatible pages no longer share the session store. +3. Verify navigation and concurrent tabs using the deployed bundles. A session + must keep its identity and conditional decision until an error or explicit + force releases it. Verify that sessions without either trigger upload no + conditional data. +4. Enable the options only after that compatibility check. Before rolling back to + an incompatible bundle, disable conditional collection and end or drain the + existing conditional sessions across the affected pages. Disabling an option + alone does not rewrite every running session's decision. + +Older modern bundles recognize only session tracking values `0`, `1`, and `2`. +They can redraw conditional values `3`, `4`, or `5`, causing unexpected collection +or data loss. The compatible legacy reader recognizes `3` and released `4`/`5`, +but it cannot recover history it never recorded. A browser cannot guarantee +cross-page persistence if its shared store stays locked or becomes unavailable +until the page closes; the SDK retries missing marks through its existing session +poll while that same session remains active. diff --git a/packages/rum/src/boot/postStartStrategy.ts b/packages/rum/src/boot/postStartStrategy.ts index c9678c6f21..58e1854666 100644 --- a/packages/rum/src/boot/postStartStrategy.ts +++ b/packages/rum/src/boot/postStartStrategy.ts @@ -87,6 +87,13 @@ export function createPostStartStrategy( return } + if (shouldForceReplay(session!, options)) { + // Applied before the guard below, not after starting: a session that withholds its replay is + // already recording, so the guard would return without ever releasing it - and releasing what + // is held is the whole of what forcing means for such a session. + sessionManager.setForcedReplay() + } + if (isRecordingInProgress(status)) { return } @@ -95,10 +102,6 @@ export function createPostStartStrategy( // Intentionally not awaiting doStart() to keep it asynchronous doStart().catch(monitorError) - - if (shouldForceReplay(session!, options)) { - sessionManager.setForcedReplay() - } } function stop() { @@ -128,5 +131,11 @@ function isRecordingInProgress(status: RecorderStatus) { } function shouldForceReplay(session: RumSession, options?: StartRecordingOptions) { - return options && options.force && session.sessionReplay === SessionReplayState.OFF + return ( + options && + options.force && + // A withheld replay is as much in need of forcing as one that was never sampled: the host asked + // for this user's replay, so it must not go on waiting for an error that may never come. + (session.sessionReplay === SessionReplayState.OFF || session.sessionReplay === SessionReplayState.BUFFERED_ON_ERROR) + ) } diff --git a/packages/rum/src/boot/recorderApi.spec.ts b/packages/rum/src/boot/recorderApi.spec.ts index 2bfeba5311..f0db0dda63 100644 --- a/packages/rum/src/boot/recorderApi.spec.ts +++ b/packages/rum/src/boot/recorderApi.spec.ts @@ -10,6 +10,7 @@ import { mockRumConfiguration, mockViewHistory, } from '../../../rum-core/test' +import { validateAndBuildRumConfiguration } from '../../../rum-core/src/domain/configuration' import type { CreateDeflateWorker } from '../domain/deflate' import { MockWorker } from '../../test' import { resetDeflateWorkerState } from '../domain/deflate' @@ -73,6 +74,40 @@ describe('makeRecorderApi', () => { } describe('recorder boot', () => { + it('starts a remotely selected buffered replay with the built recording default', async () => { + const configuration = validateAndBuildRumConfiguration({ + applicationId: 'app', + clientToken: 'token', + remoteConfigurationEnabled: true, + })! + setupRecorderApi({ + sessionManager: createRumSessionManagerMock().setTrackedWithErrorSessionReplay(), + startSessionReplayRecordingManually: configuration.startSessionReplayRecordingManually, + }) + rumInit() + expect(loadRecorderSpy).toHaveBeenCalledTimes(1) + await collectAsyncCalls(startRecordingSpy, 1) + }) + + it('keeps automatic start intent until a later session enables buffered replay', async () => { + const configuration = validateAndBuildRumConfiguration({ + applicationId: 'app', + clientToken: 'token', + remoteConfigurationEnabled: true, + })! + const sessionManager = createRumSessionManagerMock().setNotTracked() + setupRecorderApi({ + sessionManager, + startSessionReplayRecordingManually: configuration.startSessionReplayRecordingManually, + }) + rumInit() + expect(loadRecorderSpy).not.toHaveBeenCalled() + sessionManager.setTrackedWithErrorSessionReplay() + lifeCycle.notify(LifeCycleEventType.SESSION_RENEWED) + expect(loadRecorderSpy).toHaveBeenCalledTimes(1) + await collectAsyncCalls(startRecordingSpy, 1) + }) + describe('with automatic start', () => { it('starts recording when init() is called', async () => { setupRecorderApi() @@ -186,6 +221,27 @@ describe('makeRecorderApi', () => { expect(setForcedReplaySpy).toHaveBeenCalledTimes(1) }) + it('releases a withheld replay when forced, although it is already recording', async () => { + const setForcedReplaySpy = jasmine.createSpy() + + setupRecorderApi({ + sessionManager: { + ...createRumSessionManagerMock().setTrackedWithErrorSessionReplay(), + setForcedReplay: setForcedReplaySpy, + }, + startSessionReplayRecordingManually: false, + }) + + rumInit() + await collectAsyncCalls(startRecordingSpy, 1) + + // the recording is already running - what forcing asks for here is that what it holds stops + // waiting for an error + recorderApi.start({ force: true }) + + expect(setForcedReplaySpy).toHaveBeenCalledTimes(1) + }) + it('uses the previously created worker if available', async () => { setupRecorderApi({ startSessionReplayRecordingManually: true }) rumInit({ worker: mockWorker }) diff --git a/packages/rum/src/boot/startRecording.spec.ts b/packages/rum/src/boot/startRecording.spec.ts index 1f0cd6f471..ef8412c261 100644 --- a/packages/rum/src/boot/startRecording.spec.ts +++ b/packages/rum/src/boot/startRecording.spec.ts @@ -157,6 +157,47 @@ describe('startRecording', () => { expect(requests[0].metadata.records_count).toBe(1 + recordsPerFullSnapshot()) }) + it('sends the withheld replay once its session reports an error', async () => { + sessionManager.setTrackedWithErrorSessionReplay() + setupStartRecording() + + document.body.dispatchEvent(createNewEvent('click', { clientX: 1, clientY: 2 })) + // a page exit while the session is still waiting for an error keeps the buffer rather than + // sending it, so what follows joins the same segment + flushSegment(lifeCycle) + document.body.dispatchEvent(createNewEvent('click', { clientX: 3, clientY: 4 })) + + sessionManager.setSessionHasError() + flushSegment(lifeCycle) + + const requests = await readSentRequests(1) + expect(requestSendSpy).toHaveBeenCalledTimes(1) + // one segment, held since the recording started, carrying everything from before the error + expect(requests[0].metadata.creation_reason).toBe('init') + expect(requests[0].metadata.records_count).toBe(2 + recordsPerFullSnapshot()) + }) + + it('drops a withheld replay when the session stops withholding without having errored', async () => { + sessionManager.setTrackedWithErrorSessionReplay() + setupStartRecording() + + document.body.dispatchEvent(createNewEvent('click', { clientX: 1, clientY: 2 })) + + // an older SDK sharing the same session store does not know this tracking type and redraws it. + // The session stops withholding, but it never reported an error, so what it held is not owed a + // trip to the intake. + sessionManager.setTrackedWithSessionReplay() + changeView(lifeCycle) + + document.body.dispatchEvent(createNewEvent('click', { clientX: 3, clientY: 4 })) + flushSegment(lifeCycle) + + const requests = await readSentRequests(1) + // 'init' would be the withheld segment; the first one to reach the intake is the one created + // after the session stopped withholding + expect(requests[0].metadata.creation_reason).toBe('view_change') + }) + it('restarts sending segments when the session is renewed', async () => { sessionManager.setNotTracked() setupStartRecording() diff --git a/packages/rum/src/boot/startRecording.ts b/packages/rum/src/boot/startRecording.ts index 086852c49d..b00d4e4e94 100644 --- a/packages/rum/src/boot/startRecording.ts +++ b/packages/rum/src/boot/startRecording.ts @@ -1,7 +1,7 @@ import type { RawError, HttpRequest, DeflateEncoder } from '@flashcatcloud/browser-core' -import { createHttpRequest, addTelemetryDebug, canUseEventBridge } from '@flashcatcloud/browser-core' +import { createHttpRequest, addTelemetryDebug, canUseEventBridge, noop } from '@flashcatcloud/browser-core' import type { LifeCycle, ViewHistory, RumConfiguration, RumSessionManager } from '@flashcatcloud/browser-rum-core' -import { LifeCycleEventType } from '@flashcatcloud/browser-rum-core' +import { LifeCycleEventType, SessionReplayState } from '@flashcatcloud/browser-rum-core' import { record } from '../domain/record' import { startSegmentCollection, SEGMENT_BYTES_LIMIT } from '../domain/segmentCollection' @@ -28,6 +28,10 @@ export function startRecording( let addRecord: (record: BrowserRecord) => void + // Assigned once recording has started. Segment collection is created first because `record()` + // emits into it, so the buffer reaches for the snapshot through this holder rather than directly. + let takeSubsequentFullSnapshot: () => void = noop + // FLASHCAT FORK (2/4) - see `sessionReplayDirectUpload` in RumInitConfiguration. // Without the option, records are handed over to the host application through the bridge. With // it, they go through the regular segment collection and are uploaded from this page. @@ -38,7 +42,27 @@ export function startRecording( sessionManager, viewHistory, replayRequest, - encoder + encoder, + { + getWithholdingSessionId: () => { + const session = sessionManager.findTrackedSession() + return session?.sessionReplay === SessionReplayState.BUFFERED_ON_ERROR ? session.id : undefined + }, + isReleased: (sessionId) => { + const session = sessionManager.findTrackedSession() + // Still the same session, still one whose replay is kept only on an error, and no longer + // withholding. The middle condition matters: a session can stop withholding without ever + // erroring - an older SDK sharing the same store does not know these tracking types and + // redraws them - and that is a session ending, not a replay earning its way out. + return ( + !!session && + session.id === sessionId && + session.sampledOnErrorReplay && + session.sessionReplay !== SessionReplayState.BUFFERED_ON_ERROR + ) + }, + restartFromFullSnapshot: () => takeSubsequentFullSnapshot(), + } ) addRecord = segmentCollection.addRecord cleanupTasks.push(segmentCollection.stop) @@ -57,13 +81,14 @@ export function startRecording( sessionManager.findTrackedSession()?.drawnConfiguration?.defaultPrivacyLevel ?? configuration.defaultPrivacyLevel, } - const { stop: stopRecording } = record({ + const recording = record({ emit: addRecord, configuration: recordConfiguration, lifeCycle, viewHistory, }) - cleanupTasks.push(stopRecording) + takeSubsequentFullSnapshot = recording.takeSubsequentFullSnapshot + cleanupTasks.push(recording.stop) return { stop: () => { diff --git a/packages/rum/src/domain/getSessionReplayLink.ts b/packages/rum/src/domain/getSessionReplayLink.ts index 1bb7c38ea5..e8df168276 100644 --- a/packages/rum/src/domain/getSessionReplayLink.ts +++ b/packages/rum/src/domain/getSessionReplayLink.ts @@ -34,6 +34,11 @@ function getErrorType(session: RumSession | undefined, isRecordingStarted: boole // - replay sampled out return 'incorrect-session-plan' } + if (session.sessionReplay === SessionReplayState.BUFFERED_ON_ERROR) { + // the session records, but nothing has been uploaded yet and nothing may ever be: there is no + // replay to link to until the session reports an error + return 'replay-not-started' + } if (!isRecordingStarted) { return 'replay-not-started' } diff --git a/packages/rum/src/domain/record/record.ts b/packages/rum/src/domain/record/record.ts index 82e41cb1d6..c7187f1a44 100644 --- a/packages/rum/src/domain/record/record.ts +++ b/packages/rum/src/domain/record/record.ts @@ -33,6 +33,11 @@ export interface RecordOptions { export interface RecordAPI { stop: () => void flushMutations: () => void + /** + * Re-serializes the document so that the records that follow are replayable on their own. Needed + * when a withheld replay buffer is dropped, since it takes its full snapshot with it. + */ + takeSubsequentFullSnapshot: () => void shadowRootsController: ShadowRootsController } @@ -54,7 +59,7 @@ export function record(options: RecordOptions): RecordAPI { const shadowRootsController = initShadowRootsController(configuration, emitAndComputeStats, elementsScrollPositions) - const { stop: stopFullSnapshots } = startFullSnapshots( + const { stop: stopFullSnapshots, takeSubsequentFullSnapshot } = startFullSnapshots( elementsScrollPositions, shadowRootsController, lifeCycle, @@ -95,6 +100,7 @@ export function record(options: RecordOptions): RecordAPI { stopFullSnapshots() }, flushMutations, + takeSubsequentFullSnapshot, shadowRootsController, } } diff --git a/packages/rum/src/domain/record/startFullSnapshots.ts b/packages/rum/src/domain/record/startFullSnapshots.ts index 0e57d03ade..8fc0a64228 100644 --- a/packages/rum/src/domain/record/startFullSnapshots.ts +++ b/packages/rum/src/domain/record/startFullSnapshots.ts @@ -97,5 +97,19 @@ export function startFullSnapshots( unsubscribeViewCreated() unsubscribeReactivated() }, + /** + * Re-serializes the document so that what follows is replayable on its own. Used when a withheld + * replay buffer is dropped: the records kept afterwards need a full snapshot to start from. + */ + takeSubsequentFullSnapshot: () => { + flushMutations() + fullSnapshotCallback( + takeFullSnapshot(timeStampNow(), { + shadowRootsController, + status: SerializationContextStatus.SUBSEQUENT_FULL_SNAPSHOT, + elementsScrollPositions, + }) + ) + }, } } diff --git a/packages/rum/src/domain/replayStats.ts b/packages/rum/src/domain/replayStats.ts index 76c5273f6c..1dc2370f60 100644 --- a/packages/rum/src/domain/replayStats.ts +++ b/packages/rum/src/domain/replayStats.ts @@ -19,6 +19,33 @@ export function addWroteData(viewId: string, additionalBytesCount: number) { getOrCreateReplayStats(viewId).segments_total_raw_size += additionalBytesCount } +/** + * Gives back the segment count {@link addSegment} took, and with it the `index_in_view` the segment + * was holding. Segment collection serializes encoder operations, so a dropped segment returns its + * reservation after the release decision and before the next segment is created. + */ +export function removeSegment(viewId: string) { + const replayStats = statsPerView?.get(viewId) + if (!replayStats) { + return + } + replayStats.segments_count = Math.max(0, replayStats.segments_count - 1) +} + +/** + * Rolls back what a dropped segment's records contributed. These are the counters reported on view + * events, and a withheld segment that is dropped never reached the intake, so it must leave no + * trace in them. + */ +export function discardSegmentData(viewId: string, rawBytesCount: number, recordsCount: number) { + const replayStats = statsPerView?.get(viewId) + if (!replayStats) { + return + } + replayStats.records_count = Math.max(0, replayStats.records_count - recordsCount) + replayStats.segments_total_raw_size = Math.max(0, replayStats.segments_total_raw_size - rawBytesCount) +} + export function getReplayStats(viewId: string) { return statsPerView?.get(viewId) } diff --git a/packages/rum/src/domain/segmentCollection/segmentCollection.spec.ts b/packages/rum/src/domain/segmentCollection/segmentCollection.spec.ts index fe528a3661..dc7d4d1fea 100644 --- a/packages/rum/src/domain/segmentCollection/segmentCollection.spec.ts +++ b/packages/rum/src/domain/segmentCollection/segmentCollection.spec.ts @@ -1,5 +1,5 @@ import type { ClocksState, HttpRequest, TimeStamp } from '@flashcatcloud/browser-core' -import { DeflateEncoderStreamId, PageExitReason } from '@flashcatcloud/browser-core' +import { DeflateEncoderStreamId, noop, PageExitReason } from '@flashcatcloud/browser-core' import type { ViewHistory, ViewHistoryEntry, RumConfiguration } from '@flashcatcloud/browser-rum-core' import { LifeCycle, LifeCycleEventType } from '@flashcatcloud/browser-rum-core' import type { Clock } from '@flashcatcloud/browser-core/test' @@ -9,7 +9,9 @@ import type { BrowserRecord, SegmentContext } from '../../types' import { RecordType } from '../../types' import { MockWorker, readMetadataFromReplayPayload } from '../../../test' import { createDeflateEncoder } from '../deflate' +import * as replayStats from '../replayStats' import { + BUFFER_CHECKOUT_TIME, computeSegmentContext, doStartSegmentCollection, SEGMENT_BYTES_LIMIT, @@ -70,7 +72,8 @@ describe('startSegmentCollection', () => { lifeCycle, () => context, httpRequestSpy, - createDeflateEncoder(configuration, worker, DeflateEncoderStreamId.REPLAY) + createDeflateEncoder(configuration, worker, DeflateEncoderStreamId.REPLAY), + { getWithholdingSessionId: () => undefined, isReleased: () => false, restartFromFullSnapshot: noop } )) registerCleanupTask(() => { @@ -329,3 +332,488 @@ describe('computeSegmentContext', () => { } as any } }) + +describe('startSegmentCollection withholding (error session replay)', () => { + let clock: Clock + let lifeCycle: LifeCycle + let worker: MockWorker + let httpRequestSpy: { + sendOnExit: jasmine.Spy + send: jasmine.Spy + } + let addRecord: (record: BrowserRecord) => void + let withholdingSessionId: string | undefined + let releasedSessionId: string | undefined + let restartFromFullSnapshotSpy: jasmine.Spy<() => void> + let stopCollection: () => void + + function reportError() { + releasedSessionId = withholdingSessionId + withholdingSessionId = undefined + } + + beforeEach(() => { + clock = mockClock() + lifeCycle = new LifeCycle() + worker = new MockWorker() + httpRequestSpy = { sendOnExit: jasmine.createSpy(), send: jasmine.createSpy() } + withholdingSessionId = CONTEXT.session.id + releasedSessionId = undefined + restartFromFullSnapshotSpy = jasmine.createSpy() + replayStats.resetReplayStats() + + const { stop, addRecord: add } = doStartSegmentCollection( + lifeCycle, + () => CONTEXT, + httpRequestSpy, + createDeflateEncoder({} as RumConfiguration, worker, DeflateEncoderStreamId.REPLAY), + { + getWithholdingSessionId: () => withholdingSessionId, + isReleased: (sessionId) => releasedSessionId === sessionId, + restartFromFullSnapshot: restartFromFullSnapshotSpy, + } + ) + addRecord = add + stopCollection = stop + + registerCleanupTask(() => { + stop() + clock.cleanup() + replayStats.resetReplayStats() + }) + }) + + it('releases a checkout still being encoded without reusing its segment index', async () => { + addRecord({ ...RECORD, type: RecordType.FullSnapshot, data: {} } as BrowserRecord) + worker.processAllMessages() + clock.tick(BUFFER_CHECKOUT_TIME) + reportError() + lifeCycle.notify(LifeCycleEventType.RUM_EVENT_COLLECTED, { type: 'error' } as any) + addRecord(RECORD) + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + const metadata = await Promise.all( + httpRequestSpy.send.calls.allArgs().map(([payload]) => readMetadataFromReplayPayload(payload)) + ) + expect(metadata.map((segment) => segment.index_in_view)).toEqual([0, 1]) + expect(metadata[0]?.has_full_snapshot).toBeTrue() + }) + + it('remembers a release if recording ends before the worker answers', () => { + addRecord(RECORD) + worker.processAllMessages() + clock.tick(BUFFER_CHECKOUT_TIME) + reportError() + lifeCycle.notify(LifeCycleEventType.RUM_EVENT_COLLECTED, { type: 'error' } as any) + stopCollection() + releasedSessionId = undefined + worker.processAllMessages() + expect(httpRequestSpy.send).toHaveBeenCalledTimes(1) + }) + + it('drains records and a stop queued behind a released flush', async () => { + addRecord(RECORD) + worker.processAllMessages() + clock.tick(BUFFER_CHECKOUT_TIME) + addRecord(RECORD) + reportError() + stopCollection() + releasedSessionId = undefined + worker.processAllMessages() + const metadata = await Promise.all( + httpRequestSpy.send.calls.allArgs().map(([payload]) => readMetadataFromReplayPayload(payload)) + ) + expect(metadata.map((segment) => segment.index_in_view)).toEqual([0, 1]) + expect(metadata.map((segment) => segment.records_count)).toEqual([1, 1]) + }) + + it('preserves encoder ordering when a new recording starts before the old flush completes', async () => { + const sharedWorker = new MockWorker() + const sharedEncoder = createDeflateEncoder({} as RumConfiguration, sharedWorker, DeflateEncoderStreamId.REPLAY) + const sent: Array[0]> = [] + let released = false + const request = { send: (payload: Parameters[0]) => sent.push(payload), sendOnExit: noop } + const first = doStartSegmentCollection(lifeCycle, () => CONTEXT, request, sharedEncoder, { + getWithholdingSessionId: () => (released ? undefined : CONTEXT.session.id), + isReleased: () => released, + restartFromFullSnapshot: noop, + }) + first.addRecord(RECORD) + clock.tick(BUFFER_CHECKOUT_TIME) + first.addRecord(RECORD) + released = true + first.stop() + const second = doStartSegmentCollection( + new LifeCycle(), + () => ({ ...CONTEXT, session: { id: 'next-session' }, view: { id: 'next-view' } }), + request, + sharedEncoder, + { + getWithholdingSessionId: () => undefined, + isReleased: () => false, + restartFromFullSnapshot: noop, + } + ) + second.addRecord(RECORD) + second.stop() + sharedWorker.processAllMessages() + const segments = await Promise.all( + sent.map( + async (payload) => + JSON.parse(await ((payload.data as FormData).get('segment') as Blob).text()) as { + session: { id: string } + records: BrowserRecord[] + index_in_view: number + } + ) + ) + expect(segments.map((segment) => segment.session.id)).toEqual([ + CONTEXT.session.id, + CONTEXT.session.id, + 'next-session', + ]) + expect(segments.map((segment) => segment.index_in_view)).toEqual([0, 1, 0]) + expect(segments.map((segment) => segment.records.length)).toEqual([1, 1, 1]) + }) + + it('never releases an unfinished flush for a different session', () => { + addRecord(RECORD) + clock.tick(BUFFER_CHECKOUT_TIME) + releasedSessionId = 'different-session' + stopCollection() + worker.processAllMessages() + expect(httpRequestSpy.send).not.toHaveBeenCalled() + }) + + it('does not send anything while the session has not reported an error', () => { + addRecord(RECORD) + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + + expect(httpRequestSpy.send).not.toHaveBeenCalled() + expect(httpRequestSpy.sendOnExit).not.toHaveBeenCalled() + }) + + it('keeps buffering across several duration limits instead of cutting the segment', () => { + addRecord(RECORD) + clock.tick(SEGMENT_DURATION_LIMIT * 3) + addRecord(RECORD) + worker.processAllMessages() + + expect(httpRequestSpy.send).not.toHaveBeenCalled() + // still the same buffer: dropping it would have asked for a fresh full snapshot + expect(restartFromFullSnapshotSpy).not.toHaveBeenCalled() + }) + + it('sends the withheld buffer once the session reports an error', async () => { + addRecord(RECORD) + addRecord(RECORD) + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + expect(httpRequestSpy.send).not.toHaveBeenCalled() + + reportError() + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + + expect(httpRequestSpy.send).toHaveBeenCalledTimes(1) + // the records collected before the error are part of what is sent + expect((await readMetadataFromReplayPayload(httpRequestSpy.send.calls.mostRecent().args[0])).records_count).toBe(2) + }) + + it('drops the buffer and restarts from a full snapshot once it spans the checkout time', () => { + addRecord(RECORD) + clock.tick(BUFFER_CHECKOUT_TIME) + worker.processAllMessages() + + expect(httpRequestSpy.send).not.toHaveBeenCalled() + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(1) + }) + + it('drops the buffer and restarts from a full snapshot when it grows past the bytes limit', () => { + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + + expect(httpRequestSpy.send).not.toHaveBeenCalled() + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(1) + }) + + it('does not restart in a hot loop when the full snapshot alone exceeds the bytes limit', () => { + // every restart would blow the limit again straight away on such a document + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(1) + + clock.tick(SEGMENT_DURATION_LIMIT) + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(2) + }) + + it('restores a full snapshot after consecutive oversized snapshots and an error', async () => { + restartFromFullSnapshotSpy.and.callFake(() => addRecord(VERY_BIG_RECORD)) + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + expect(httpRequestSpy.send).not.toHaveBeenCalled() + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(1) + + reportError() + addRecord(RECORD) + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + + expect(httpRequestSpy.send).toHaveBeenCalled() + expect( + (await readMetadataFromReplayPayload(httpRequestSpy.send.calls.first().args[0])).has_full_snapshot + ).toBeTrue() + }) + + it('restores a missing snapshot before an errored page exits during the restart delay', async () => { + restartFromFullSnapshotSpy.and.callFake(() => addRecord(VERY_BIG_RECORD)) + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + reportError() + addRecord(RECORD) + lifeCycle.notify(LifeCycleEventType.PAGE_MAY_EXIT, { reason: PageExitReason.UNLOADING }) + worker.processAllMessages() + + expect(httpRequestSpy.sendOnExit).toHaveBeenCalled() + expect( + (await readMetadataFromReplayPayload(httpRequestSpy.sendOnExit.calls.first().args[0])).has_full_snapshot + ).toBeTrue() + }) + + it('restores the missing snapshot as soon as an error releases the session', () => { + restartFromFullSnapshotSpy.and.callFake(() => addRecord(VERY_BIG_RECORD)) + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + reportError() + lifeCycle.notify(LifeCycleEventType.RUM_EVENT_COLLECTED, { type: 'error' } as any) + + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(2) + worker.processAllMessages() + expect(httpRequestSpy.send).toHaveBeenCalled() + }) + + it('cancels the delayed replacement when a new view supplies a snapshot', () => { + restartFromFullSnapshotSpy.and.callFake(() => addRecord(VERY_BIG_RECORD)) + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + lifeCycle.notify(LifeCycleEventType.VIEW_CREATED, {} as any) + addRecord({ ...VERY_BIG_RECORD, data: {} } as BrowserRecord) + worker.processAllMessages() + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(1) + expect(httpRequestSpy.send).not.toHaveBeenCalled() + }) + + it('does not repeatedly serialize an oversized document while waiting for an error', () => { + restartFromFullSnapshotSpy.and.callFake(() => addRecord(VERY_BIG_RECORD)) + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + for (let i = 0; i < 4; i++) { + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + } + + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(1) + expect(httpRequestSpy.send).not.toHaveBeenCalled() + }) + + it('cancels a delayed snapshot when recording stops', () => { + restartFromFullSnapshotSpy.and.callFake(() => addRecord(VERY_BIG_RECORD)) + addRecord(VERY_BIG_RECORD) + worker.processAllMessages() + stopCollection() + clock.tick(SEGMENT_DURATION_LIMIT * 2) + worker.processAllMessages() + + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(1) + expect(httpRequestSpy.send).not.toHaveBeenCalled() + }) + + it('keeps the buffer when the page is only hidden, so the replay can still start from its snapshot', () => { + // switching tabs is ordinary; dropping here would take the only full snapshot with it + addRecord(RECORD) + addRecord(RECORD) + lifeCycle.notify(LifeCycleEventType.PAGE_MAY_EXIT, { reason: PageExitReason.HIDDEN }) + worker.processAllMessages() + + expect(httpRequestSpy.sendOnExit).not.toHaveBeenCalled() + expect(restartFromFullSnapshotSpy).not.toHaveBeenCalled() + + reportError() + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + + expect(httpRequestSpy.send).toHaveBeenCalledTimes(1) + }) + + it('sends nothing on page exit for a session that never errored', () => { + addRecord(RECORD) + lifeCycle.notify(LifeCycleEventType.PAGE_MAY_EXIT, { reason: PageExitReason.UNLOADING }) + worker.processAllMessages() + + expect(httpRequestSpy.sendOnExit).not.toHaveBeenCalled() + }) + + it('does not let a dropped buffer leave its index_in_view behind for the next one to collide with', async () => { + // the restart emits records, exactly as taking a fresh full snapshot does in production + restartFromFullSnapshotSpy.and.callFake(() => addRecord(RECORD)) + + addRecord(RECORD) + clock.tick(BUFFER_CHECKOUT_TIME) + worker.processAllMessages() + expect(restartFromFullSnapshotSpy).toHaveBeenCalledTimes(1) + + reportError() + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + + // the dropped buffer never reached the intake, so the first segment that does is index 0 + const metadata = await readMetadataFromReplayPayload(httpRequestSpy.send.calls.mostRecent().args[0]) + expect(metadata.index_in_view).toBe(0) + // and it carries a reason the segment schema knows, not the internal one that dropped the buffer + expect(metadata.creation_reason).toBe('segment_duration_limit') + }) + + it('does not restart the buffer when collection was stopped while the flush was in flight', () => { + addRecord(RECORD) + // the checkout flush is posted to the worker, and recording is stopped before it answers + clock.tick(BUFFER_CHECKOUT_TIME) + stopCollection() + worker.processAllMessages() + + expect(restartFromFullSnapshotSpy).not.toHaveBeenCalled() + }) + + it('does not hand the next segment an index the dropped one still holds when a record lands mid-flush', async () => { + restartFromFullSnapshotSpy.and.callFake(() => addRecord(RECORD)) + + addRecord(RECORD) + // The flush is posted to the worker but not answered yet - in production that round trip always + // happens, because flushing writes the trailer before finishing. A record arriving now creates + // the next segment, which reads its index while the dropped one is still counted. + clock.tick(BUFFER_CHECKOUT_TIME) + addRecord(RECORD) + worker.processAllMessages() + + reportError() + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + + expect((await readMetadataFromReplayPayload(httpRequestSpy.send.calls.mostRecent().args[0])).index_in_view).toBe(0) + }) + + it('leaves no trace of a dropped buffer in the replay stats', () => { + addRecord(RECORD) + clock.tick(BUFFER_CHECKOUT_TIME) + worker.processAllMessages() + + const stats = replayStats.getReplayStats(CONTEXT.view.id) + expect(stats?.segments_count ?? 0).toBe(0) + expect(stats?.segments_total_raw_size ?? 0).toBe(0) + }) + + it('sends normally once released, without withholding the following segments', () => { + reportError() + addRecord(RECORD) + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + addRecord(RECORD) + clock.tick(SEGMENT_DURATION_LIMIT) + worker.processAllMessages() + + expect(httpRequestSpy.send).toHaveBeenCalledTimes(2) + }) +}) + +describe('startSegmentCollection withholding, session lifecycle', () => { + let clock: Clock + let lifeCycle: LifeCycle + let worker: MockWorker + let httpRequestSpy: { + sendOnExit: jasmine.Spy + send: jasmine.Spy + } + let addRecord: (record: BrowserRecord) => void + let stopSegmentCollection: () => void + let withholdingSessionId: string | undefined + let releasedSessionId: string | undefined + + beforeEach(() => { + clock = mockClock() + lifeCycle = new LifeCycle() + worker = new MockWorker() + httpRequestSpy = { sendOnExit: jasmine.createSpy(), send: jasmine.createSpy() } + withholdingSessionId = CONTEXT.session.id + releasedSessionId = undefined + + const { stop, addRecord: add } = doStartSegmentCollection( + lifeCycle, + () => CONTEXT, + httpRequestSpy, + createDeflateEncoder({} as RumConfiguration, worker, DeflateEncoderStreamId.REPLAY), + { + getWithholdingSessionId: () => withholdingSessionId, + isReleased: (sessionId) => releasedSessionId === sessionId, + restartFromFullSnapshot: () => undefined, + } + ) + addRecord = add + stopSegmentCollection = stop + + registerCleanupTask(() => { + stopSegmentCollection() + clock.cleanup() + }) + }) + + it('drops the buffer when the session expires without ever reporting an error', () => { + addRecord(RECORD) + // the session is gone, so nothing answers for these records any more + withholdingSessionId = undefined + releasedSessionId = undefined + + stopSegmentCollection() + worker.processAllMessages() + + expect(httpRequestSpy.send).not.toHaveBeenCalled() + expect(httpRequestSpy.sendOnExit).not.toHaveBeenCalled() + }) + + it('drops the buffer when the session is renewed into a different one', () => { + addRecord(RECORD) + withholdingSessionId = undefined + releasedSessionId = 'a-different-session' + + stopSegmentCollection() + worker.processAllMessages() + + expect(httpRequestSpy.send).not.toHaveBeenCalled() + expect(httpRequestSpy.sendOnExit).not.toHaveBeenCalled() + }) + + it('sends the buffer when its own session reports the error', () => { + addRecord(RECORD) + releasedSessionId = withholdingSessionId + withholdingSessionId = undefined + + stopSegmentCollection() + worker.processAllMessages() + + expect(httpRequestSpy.send).toHaveBeenCalledTimes(1) + }) +}) diff --git a/packages/rum/src/domain/segmentCollection/segmentCollection.ts b/packages/rum/src/domain/segmentCollection/segmentCollection.ts index ac2dbfabb9..7c7ed72d98 100644 --- a/packages/rum/src/domain/segmentCollection/segmentCollection.ts +++ b/packages/rum/src/domain/segmentCollection/segmentCollection.ts @@ -1,13 +1,29 @@ -import type { DeflateEncoder, HttpRequest, TimeoutId } from '@flashcatcloud/browser-core' -import { isPageExitReason, ONE_SECOND, clearTimeout, setTimeout } from '@flashcatcloud/browser-core' +import type { DeflateEncoder, HttpRequest, RelativeTime, TimeoutId } from '@flashcatcloud/browser-core' +import { + addTelemetryDebug, + isPageExitReason, + ONE_SECOND, + clearTimeout, + relativeNow, + setTimeout, +} from '@flashcatcloud/browser-core' import type { LifeCycle, ViewHistory, RumSessionManager, RumConfiguration } from '@flashcatcloud/browser-rum-core' import { LifeCycleEventType } from '@flashcatcloud/browser-rum-core' import type { BrowserRecord, CreationReason, SegmentContext } from '../../types' +import { RecordType } from '../../types' +import { discardSegmentData, removeSegment } from '../replayStats' import { buildReplayPayload } from './buildReplayPayload' import type { FlushReason, Segment } from './segment' import { createSegment } from './segment' export const SEGMENT_DURATION_LIMIT = 5 * ONE_SECOND + +/** + * How much history a withheld buffer may span before it is dropped and restarted from a fresh full + * snapshot. This bounds two things at once: the memory a session that never errors holds on to, and + * how far back an error session can show once its buffer is released. + */ +export const BUFFER_CHECKOUT_TIME = 60 * ONE_SECOND /** * beacon payload max queue size implementation is 64kb * ensure that we leave room for logs, rum and potential other users @@ -39,19 +55,43 @@ export let SEGMENT_BYTES_LIMIT = 60_000 // To help investigate session replays issues, each segment is created with a "creation reason", // indicating why the session has been created. +/** + * Lets a session record without uploading anything until it reports an error. Sessions drawn by + * `sessionReplayOnError` record from the start, but every segment is withheld: dropped on + * checkout while no error has happened, sent normally from the moment one has. + */ +export interface SegmentBuffering { + /** + * The id of the current session if it is withholding its replay, `undefined` otherwise. A segment + * remembers this at creation, so that what happens to it later is decided by the session that + * actually produced its records. + */ + getWithholdingSessionId: () => string | undefined + /** + * Whether that same session has since reported its error. Anything else — the session expired, or + * was renewed into a different one — means the records were never released and must be dropped: + * uploading them would bill a session for a replay nobody asked for and nobody can explain. + */ + isReleased: (sessionId: string) => boolean + /** Restarts the buffer from a fresh full snapshot, after the previous one was dropped. */ + restartFromFullSnapshot: () => void +} + export function startSegmentCollection( lifeCycle: LifeCycle, configuration: RumConfiguration, sessionManager: RumSessionManager, viewHistory: ViewHistory, httpRequest: HttpRequest, - encoder: DeflateEncoder + encoder: DeflateEncoder, + buffering: SegmentBuffering ) { return doStartSegmentCollection( lifeCycle, () => computeSegmentContext(configuration.applicationId, sessionManager, viewHistory), httpRequest, - encoder + encoder, + buffering ) } @@ -69,30 +109,86 @@ type SegmentCollectionState = status: SegmentCollectionStatus.SegmentPending segment: Segment expirationTimeoutId: TimeoutId + /** Only armed while the segment is withheld: bounds how much history the buffer may span. */ + bufferCheckoutTimeoutId: TimeoutId | undefined + /** Set when the segment was created while its session was withholding its replay. */ + withheldForSessionId: string | undefined } | { status: SegmentCollectionStatus.Stopped } +/** + * These two are internal and never reach the intake, so they are mapped back to a schema value where + * the next segment records why it was created. `buffer_checkout` drops a withheld buffer that has + * grown past {@link BUFFER_CHECKOUT_TIME}; `page_reactivated` cuts a segment when the page is + * switched back to, so the next one starts from the fresh full snapshot taken on the same event. + */ +type InternalFlushReason = FlushReason | 'buffer_checkout' | 'page_reactivated' + +// Recordings can stop and restart while the same encoder is still finishing a segment. +// Serialize at the encoder boundary so their metadata and index reservations cannot overlap. +let encodingQueues: WeakMap void> }> | undefined + export function doStartSegmentCollection( lifeCycle: LifeCycle, getSegmentContext: () => SegmentContext | undefined, httpRequest: HttpRequest, - encoder: DeflateEncoder + encoder: DeflateEncoder, + buffering: SegmentBuffering ) { let state: SegmentCollectionState = { status: SegmentCollectionStatus.WaitingForInitialRecord, nextSegmentCreationReason: 'init', } + // How many buffers were dropped before one was finally released. Without this, "the replay goes + // back up to a minute" is a promise nobody can check. + let droppedBufferCount = 0 + let lastBufferRestartAt: RelativeTime | undefined + let bufferRestartTimeoutId: TimeoutId | undefined + encodingQueues ||= new WeakMap() + const encodingQueue = encodingQueues.get(encoder) || { flushing: false, operations: [] } + encodingQueues.set(encoder, encodingQueue) + let stopped = false + const withholdingSessionIds = new Set() + const releasedSessionIds = new Set() + + function rememberReleases() { + withholdingSessionIds.forEach((sessionId) => { + if (buffering.isReleased(sessionId)) { + releasedSessionIds.add(sessionId) + } + }) + } + + function runWhenReady(operation: () => void) { + encodingQueue.operations.push(operation) + drainPendingOperations() + } + + function drainPendingOperations() { + while (!encodingQueue.flushing && encodingQueue.operations.length) { + encodingQueue.operations.shift()!() + } + } + + function requestFlush(reason: InternalFlushReason) { + rememberReleases() + if (reason !== 'view_change' && reason !== 'page_reactivated') { + restoreReleasedSnapshot() + } + runWhenReady(() => flushSegment(reason)) + } + const { unsubscribe: unsubscribeViewCreated } = lifeCycle.subscribe(LifeCycleEventType.VIEW_CREATED, () => { - flushSegment('view_change') + requestFlush('view_change') }) const { unsubscribe: unsubscribePageMayExit } = lifeCycle.subscribe( LifeCycleEventType.PAGE_MAY_EXIT, (pageMayExitEvent) => { - flushSegment(pageMayExitEvent.reason as FlushReason) + requestFlush(pageMayExitEvent.reason as FlushReason) } ) @@ -100,12 +196,103 @@ export function doStartSegmentCollection( // next one starts fresh with the full snapshot taken by startFullSnapshots on the same event. // Reuses the 'view_change' creation reason to avoid a schema change. const { unsubscribe: unsubscribeReactivated } = lifeCycle.subscribe(LifeCycleEventType.PAGE_REACTIVATED, () => { - flushSegment('view_change') + requestFlush('page_reactivated') }) - function flushSegment(flushReason: FlushReason) { + const { unsubscribe: unsubscribeRumEvent } = lifeCycle.subscribe( + LifeCycleEventType.RUM_EVENT_COLLECTED, + restoreReleasedSnapshot + ) + + const { unsubscribe: unsubscribeSessionReleased } = lifeCycle.subscribe( + LifeCycleEventType.SESSION_RELEASED, + ({ sessionId }) => { + if (withholdingSessionIds.has(sessionId)) { + releasedSessionIds.add(sessionId) + } + restoreReleasedSnapshot() + } + ) + + function restoreReleasedSnapshot() { + rememberReleases() + if (bufferRestartTimeoutId === undefined) { + return + } + const context = getSegmentContext() + if (context && buffering.isReleased(context.session.id)) { + // The error tracker marks the session before this listener runs. Restore the missing + // baseline now, before a view change or page exit can flush an unplayable segment. + clearTimeout(bufferRestartTimeoutId) + bufferRestartTimeoutId = undefined + lastBufferRestartAt = relativeNow() + buffering.restartFromFullSnapshot() + } else { + // The same oversized snapshot would be discarded again. Poll only for a release, without + // repeatedly serializing the document when neither an error nor new activity has arrived. + clearTimeout(bufferRestartTimeoutId) + bufferRestartTimeoutId = setTimeout(restoreReleasedSnapshot, SEGMENT_DURATION_LIMIT) + } + } + + function flushSegment(flushReason: InternalFlushReason) { + // Keep the encoder and index reservation owned by this segment until its asynchronous + // decision settles. Later records retain their emission context while waiting in FIFO order. + const withheldForSessionId = + state.status === SegmentCollectionStatus.SegmentPending ? state.withheldForSessionId : undefined + const isWithheld = withheldForSessionId !== undefined && !releasedSessionIds.has(withheldForSessionId) + if (state.status === SegmentCollectionStatus.SegmentPending) { + if (isWithheld && flushReason === 'page_reactivated') { + // The fresh full snapshot taken on the same event lands inside the withheld buffer, which + // stays replayable from it. Cutting here would only throw away what came before the switch. + return + } + + if (isWithheld && (flushReason === 'segment_duration_limit' || isPageExitReason(flushReason))) { + // Nothing can be sent while withheld, so these rotations would only throw the buffer away - + // and with it the full snapshot a released replay has to start from, leaving the rest of the + // session as incremental records nothing can be played from. A page that is merely hidden or + // frozen comes back and goes on recording; one that is really unloading takes the buffer with + // it either way. Keeping it is never worse than dropping it. + if (flushReason === 'segment_duration_limit') { + // Re-armed, so the buffer is flushed normally within one rotation of the session erroring. + // An expiring session does not lose it: the session history entry is still open when the + // recorder is stopped (`sessionManager.ts` notifies before closing it), so the stop flush + // still sees the session as released and sends. Only losing the page outright loses it. + state.expirationTimeoutId = setTimeout(() => requestFlush('segment_duration_limit'), SEGMENT_DURATION_LIMIT) + } + return + } + + encodingQueue.flushing = true state.segment.flush((metadata, encoderResult) => { + rememberReleases() + if (withheldForSessionId !== undefined && !releasedSessionIds.has(withheldForSessionId)) { + removeSegment(metadata.view.id) + // No error was reported, so this buffer is dropped rather than sent. Rolling back what its + // records contributed keeps `has_replay` and the counters on view events honest. + discardSegmentData(metadata.view.id, encoderResult.rawBytesCount, metadata.records_count) + droppedBufferCount += 1 + // Restarted from here rather than synchronously below, so the fresh full snapshot lands in + // the segment that follows this one rather than in the one being thrown away. + restartBuffer(flushReason) + encodingQueue.flushing = false + drainPendingOperations() + return + } + + if (withheldForSessionId !== undefined) { + // The first segment released by an error: report how much history it actually carried, so + // the window we promise can be compared against the one users get. + addTelemetryDebug('Error session replay buffer released', { + 'buffer.duration': metadata.end - metadata.start, + 'buffer.records_count': metadata.records_count, + 'buffer.dropped_count': droppedBufferCount, + }) + droppedBufferCount = 0 + } + const payload = buildReplayPayload(encoderResult.output, metadata, encoderResult.rawBytesCount) if (isPageExitReason(flushReason)) { @@ -113,14 +300,22 @@ export function doStartSegmentCollection( } else { httpRequest.send(payload) } + encodingQueue.flushing = false + drainPendingOperations() }) clearTimeout(state.expirationTimeoutId) + clearTimeout(state.bufferCheckoutTimeoutId) } if (flushReason !== 'stop') { state = { status: SegmentCollectionStatus.WaitingForInitialRecord, - nextSegmentCreationReason: flushReason, + nextSegmentCreationReason: + flushReason === 'buffer_checkout' + ? 'segment_duration_limit' + : flushReason === 'page_reactivated' + ? 'view_change' + : flushReason, } } else { state = { @@ -129,39 +324,103 @@ export function doStartSegmentCollection( } } - return { - addRecord: (record: BrowserRecord) => { - if (state.status === SegmentCollectionStatus.Stopped) { + /** + * A dropped buffer leaves no full snapshot behind, so the next one would not be replayable on its + * own. A view change does not need this: the new view emits its own full snapshot. + */ + function restartBuffer(flushReason: InternalFlushReason) { + if (flushReason !== 'buffer_checkout' && flushReason !== 'segment_bytes_limit') { + return + } + if (stopped || state.status === SegmentCollectionStatus.Stopped) { + // The flush that got here waited on the deflate worker, and recording was stopped in the + // meantime. Re-serializing the document now would cost a full snapshot on a page that asked + // to stop, and count records into the replay stats that no segment will ever hold. + return + } + // A snapshot can itself exceed the budget. After a rapid second discard, wait for release + // before replacing it: ordinary flushes no longer restart buffers once the session errors. + clearTimeout(bufferRestartTimeoutId) + bufferRestartTimeoutId = undefined + const now = relativeNow() + const delay = lastBufferRestartAt === undefined ? 0 : SEGMENT_DURATION_LIMIT - (now - lastBufferRestartAt) + if (delay > 0) { + bufferRestartTimeoutId = setTimeout(restoreReleasedSnapshot, delay) + } else { + lastBufferRestartAt = now + buffering.restartFromFullSnapshot() + } + } + + function addRecord( + record: BrowserRecord, + context: SegmentContext | undefined, + withheldForSessionId: string | undefined + ) { + if (state.status === SegmentCollectionStatus.Stopped) { + return + } + + if (record.type === RecordType.FullSnapshot) { + // A view change or page reactivation can supply the replacement before the timer does. + clearTimeout(bufferRestartTimeoutId) + bufferRestartTimeoutId = undefined + } + + if (state.status === SegmentCollectionStatus.WaitingForInitialRecord) { + if (!context) { return } - if (state.status === SegmentCollectionStatus.WaitingForInitialRecord) { - const context = getSegmentContext() - if (!context) { - return - } + state = { + status: SegmentCollectionStatus.SegmentPending, + segment: createSegment({ encoder, context, creationReason: state.nextSegmentCreationReason }), + expirationTimeoutId: setTimeout(() => { + requestFlush('segment_duration_limit') + }, SEGMENT_DURATION_LIMIT), + bufferCheckoutTimeoutId: + withheldForSessionId !== undefined + ? setTimeout(() => { + requestFlush('buffer_checkout') + }, BUFFER_CHECKOUT_TIME) + : undefined, + withheldForSessionId, + } + } - state = { - status: SegmentCollectionStatus.SegmentPending, - segment: createSegment({ encoder, context, creationReason: state.nextSegmentCreationReason }), - expirationTimeoutId: setTimeout(() => { - flushSegment('segment_duration_limit') - }, SEGMENT_DURATION_LIMIT), - } + state.segment.addRecord(record, (encodedBytesCount) => { + if (encodedBytesCount > SEGMENT_BYTES_LIMIT) { + requestFlush('segment_bytes_limit') } + }) + } - state.segment.addRecord(record, (encodedBytesCount) => { - if (encodedBytesCount > SEGMENT_BYTES_LIMIT) { - flushSegment('segment_bytes_limit') - } - }) + return { + addRecord: (record: BrowserRecord) => { + if (stopped) { + return + } + const context = getSegmentContext() + const withheldForSessionId = buffering.getWithholdingSessionId() + if (withheldForSessionId !== undefined) { + withholdingSessionIds.add(withheldForSessionId) + } + rememberReleases() + runWhenReady(() => addRecord(record, context, withheldForSessionId)) }, - stop: () => { - flushSegment('stop') + if (stopped) { + return + } + requestFlush('stop') + stopped = true + clearTimeout(bufferRestartTimeoutId) + bufferRestartTimeoutId = undefined unsubscribeViewCreated() unsubscribePageMayExit() unsubscribeReactivated() + unsubscribeRumEvent() + unsubscribeSessionReleased() }, } }