Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions docs/browser-sdk.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ Every field accepted by `MapleBrowser.init`:
| `errors` | `ErrorFilterOptions` | see [Filtering errors](#filtering-errors) | Drop captured errors by message, script URL, or a `beforeCapture` hook. |
| `replay.enabled` | `boolean` | `true` | Enable rrweb session recording. |
| `replay.sampleRate` | `number` | `1` | Fraction of sessions to record, `0` to `1`. Out-of-range values are clamped with a warning. See [Sampling](#sampling). |
| `replay.onErrorSampleRate` | `number` | `0` | Fraction of the sessions not recorded that buffer the last minute in memory and keep it only if an error happens. See [Sampling](#sampling). |
| `privacy.maskAllInputs` | `boolean` | `true` | Mask all `<input>` values in the recording. |
| `privacy.maskAllText` | `boolean` | `false` | Mask all text in the rrweb recording and omit captured click-target text from session events. |
| `privacy.persistVisitorId` | `boolean` | `true` | Store a persistent visitor id (localStorage + cookie) so unique visitors and new-vs-returning are measurable. Turning it off also purges any id already stored. |
Expand Down Expand Up @@ -456,6 +457,26 @@ privacy: {

To record only a fraction of sessions, set `replay.sampleRate` between `0` and `1`. For example, `0.1` records ~10% of sessions. Tracing is unaffected by this setting.

### Replay on error

`replay.onErrorSampleRate` covers the sessions `replay.sampleRate` leaves out. Those sessions run
the recorder into memory only, keeping roughly the last minute (the segments since the
second-to-last full snapshot, taken every 30s). Nothing is uploaded. When an error is recorded (an
uncaught error, an unhandled rejection or `captureException`, after [filters](#filtering-errors)),
the buffered minute is uploaded and the rest of the session is recorded normally, including its
later page loads. The session is marked `maple.session.replay_trigger: "error"`, and its replay starts
up to a minute before the error rather than at the start of the session.

```ts
MapleBrowser.init({
// ...
replay: { sampleRate: 0.05, onErrorSampleRate: 1 }, // 5% of sessions, plus every session with an error
})
```

Buffered sessions download the replay chunk like recorded ones, and keep up to 4 MB of events in
memory.

`tracing.sampleRate` does the same for traces. The decision is made once per session (a hash of
`session.id`), so a sampled session keeps every one of its traces and its replay never links to a
dropped one. Spans that record an error are always exported, whatever the rate.
Expand Down
7 changes: 7 additions & 0 deletions packages/browser-session/src/events/meta-row.ts
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,8 @@ export interface SessionMetaRowInput {
* rendering a player with nothing to play.
*/
readonly recorded: boolean
/** Why a recording exists when it did not start with the session: an error triggered it. */
readonly replayTrigger?: "error" | undefined
/**
* Whether this session owns its visit claim. The gateway bills
* `billable_start == 1 && version == 1`, and falls back to `version == 1`
Expand Down Expand Up @@ -119,6 +121,11 @@ export function buildSessionMetaRow(input: SessionMetaRowInput): SessionMetaRow
// first, and treats an absent key as "rrweb" — every session recorded
// before this marker existed is a browser recording.
"maple.session.replay_format": "rrweb",
// Replay buffered in memory and uploaded because an error happened: the
// recording starts up to a minute before the error, not at session start.
...(input.recorded && input.replayTrigger
? { "maple.session.replay_trigger": input.replayTrigger }
: undefined),
// Storage-blocked visitors get an in-memory id, so their sessions each
// look like a distinct visitor. Flag it rather than inflate silently.
...(input.visitorId && input.visitorIdPersisted === false
Expand Down
9 changes: 8 additions & 1 deletion packages/browser-session/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,14 @@ export { startMetadataSession } from "./session/metadata-session"
// package-internal: `startSessionLifecycle` owns those invariants, and an SDK
// reaching past it would write counts the lifecycle then overwrites.
export type { SessionRecord } from "./session/session"
export { claimReplaySample, getSession, getSessionId, rotateSession } from "./session/session"
export type { ReplayMode } from "./session/session"
export {
claimReplayMode,
claimReplaySample,
getSession,
getSessionId,
rotateSession,
} from "./session/session"
export type { MapleBrowserSessionSink } from "./session/sink"
export { clearSessionSink } from "./session/sink"
export { getObservedTraceIds, publishSessionSink, readSessionSink, recordTraceId } from "./session/sink"
Expand Down
58 changes: 57 additions & 1 deletion packages/browser-session/src/replay/record.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ vi.mock("../platform/transport", () => ({
}),
}))

const { startRecording } = await import("./record")
const { startBufferedRecording, startRecording } = await import("./record")

const CONFIG = {
endpoint: "https://ingest.example",
Expand Down Expand Up @@ -206,3 +206,59 @@ describe("startRecording", () => {
}
})
})

const META = 4
const meta = (timestamp: number) => ({
type: META,
timestamp,
data: { href: "https://app.example/?token=abc" },
})
/** One rrweb snapshot: a Meta event, then the FullSnapshot. */
const snapshot = (timestamp: number) => {
emitRef!(meta(timestamp), true)
emitRef!(fullSnapshot(timestamp + 1), true)
}

describe("startBufferedRecording", () => {
beforeEach(() => {
posted.length = 0
outcomes.length = 0
stopFn.mockClear()
emitRef = undefined
})

it("uploads nothing until drained, then the last two snapshots' segments as checkpoints", async () => {
const recorder = startBufferedRecording(CONFIG, "session-1")
emitRef!(incremental(500))
snapshot(1_000)
emitRef!(incremental(1_500))
snapshot(31_000)
emitRef!(incremental(31_500))
snapshot(61_000)
emitRef!(incremental(61_500))
emitRef!(incremental(62_000))
expect(posted).toEqual([])

await recorder.drain()
expect(posted.map((chunk) => chunk.meta)).toEqual([
{ sessionId: "session-1", chunkSeq: 1, isCheckpoint: true, eventCount: 3, durationMs: 500 },
{ sessionId: "session-1", chunkSeq: 1, isCheckpoint: true, eventCount: 4, durationMs: 1_000 },
])
const first = JSON.parse(posted[0]!.body) as Array<{ timestamp: number; data: { href?: string } }>
expect(first[0]?.timestamp).toBe(31_000)
expect(first[0]?.data.href).toBe("https://app.example/?token=REDACTED")

await recorder.drain()
expect(posted).toHaveLength(2)
})

it("discards the buffer on stop", async () => {
const recorder = startBufferedRecording(CONFIG, "session-1")
snapshot(1_000)
emitRef!(incremental(1_500))
recorder.stop()
await recorder.drain()
expect(posted).toEqual([])
expect(stopFn).toHaveBeenCalled()
})
})
165 changes: 142 additions & 23 deletions packages/browser-session/src/replay/record.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,13 @@ import { record } from "rrweb"
import { BLOCK_SELECTOR } from "../privacy-markers"
import { markActivity, nextChunkSeq } from "../session/session"
import type { IngestConfig } from "../platform/transport"
import { gzip, postSessionBlob, warnDropped, type ChunkMeta } from "../platform/transport"
import {
type BlobPostOutcome,
gzip,
postSessionBlob,
warnDropped,
type ChunkMeta,
} from "../platform/transport"
import { scrubUrl } from "../platform/url-privacy"

// rrweb event shape — typed loosely to avoid coupling to @rrweb/types across
Expand Down Expand Up @@ -54,6 +60,33 @@ export interface Recorder {
getClickCount: () => number
}

/**
* gzip and POST one chunk. `seq` is claimed here, monotonic across reloads
* (persisted on the session record), so a refresh continues the sequence
* instead of overwriting the previous load's blobs.
*/
async function uploadChunk(
config: IngestConfig,
sessionId: string,
body: string,
chunk: Omit<ChunkMeta, "sessionId" | "chunkSeq">,
keepalive: boolean,
): Promise<BlobPostOutcome | undefined> {
const chunkSeq = nextChunkSeq()
// `gzip` rejects rather than returning a truncated stream, and callers run
// this as a floating promise: an escaping rejection would surface in the
// host app's console as ours. Dropping the chunk is the same outcome ingest
// produced by refusing it, minus the wasted POST.
let gzipped: Uint8Array
try {
gzipped = await gzip(new TextEncoder().encode(body))
} catch (error) {
warnDropped("chunk compression", error)
return undefined
}
return postSessionBlob(config, { sessionId, chunkSeq, ...chunk }, gzipped, keepalive)
}

export function startRecording(config: IngestConfig, sessionId: string): Recorder {
// Events are serialized once at emit time and buffered as JSON strings, so
// flushing is a cheap `join` instead of re-stringifying the whole buffer
Expand Down Expand Up @@ -92,30 +125,14 @@ export function startRecording(config: IngestConfig, sessionId: string): Recorde
const isCheckpoint = bufferHasCheckpoint
const eventCount = parts.length
const durationMs = Math.max(0, lastTimestamp - firstTimestamp)
// Monotonic across reloads (persisted on the session record), so a refresh
// continues the sequence instead of overwriting the previous load's blobs.
const seq = nextChunkSeq()
resetBuffer()

// `gzip` now rejects rather than returning a truncated stream, and `flush`
// is called as a floating promise — an escaping rejection would surface in
// the host app's console as ours. Dropping the chunk is the same outcome
// ingest produced by refusing it, minus the wasted POST.
let gzipped: Uint8Array
try {
gzipped = await gzip(new TextEncoder().encode(body))
} catch (error) {
warnDropped("chunk compression", error)
return
}
const meta: ChunkMeta = {
const outcome = await uploadChunk(
config,
sessionId,
chunkSeq: seq,
isCheckpoint,
eventCount,
durationMs,
}
const outcome = await postSessionBlob(config, meta, gzipped, keepalive)
body,
{ isCheckpoint, eventCount, durationMs },
keepalive,
)
if (outcome === "exhausted" && !exhausted) {
exhausted = true
warnExhausted(sessionId)
Expand Down Expand Up @@ -243,3 +260,105 @@ export function startRecording(config: IngestConfig, sessionId: string): Recorde
getClickCount: () => clickCount,
}
}

/** Buffer mode checks out often, so the retained window stays near a minute. */
const BUFFER_CHECKOUT_MS = 30_000

interface Segment {
parts: string[]
bytes: number
first: number
last: number
}

export interface BufferedRecorder {
/** Upload what is buffered, oldest first, each segment a checkpoint chunk. */
drain: (keepalive?: boolean) => Promise<void>
stop: () => void
getClickCount: () => number
}

/**
* Record into memory only: the segments since the second-to-last checkout, so
* 30-60s of replay, and nothing is uploaded until `drain()`. For sessions that
* keep a replay only when an error happens.
*/
export function startBufferedRecording(config: IngestConfig, sessionId: string): BufferedRecorder {
let segments: Segment[] = []
let bytes = 0
let clickCount = 0
let stopped = false

const stopRecord = record({
emit: (event: unknown) => {
const active = markActivity()
if (stopped || (active && active.id !== sessionId)) return
const e = event as RrwebEvent
if (e.type === META && e.data && typeof e.data.href === "string")
e.data.href = scrubUrl(e.data.href)
if (
e.type === INCREMENTAL &&
e.data?.source === SOURCE_MOUSE_INTERACTION &&
e.data.type === MOUSE_CLICK
) {
clickCount++
}
let json: string
try {
json = JSON.stringify(e)
} catch {
return
}
// Every snapshot, first or checkout, is a Meta event then a FullSnapshot.
if (e.type === META) {
segments.push({ parts: [], bytes: 0, first: e.timestamp, last: e.timestamp })
while (segments.length > 2) bytes -= segments.shift()?.bytes ?? 0
}
const segment = segments.at(-1)
// Nothing to play back before the first snapshot.
if (!segment) return
segment.parts.push(json)
segment.bytes += json.length
segment.last = e.timestamp
bytes += json.length
while (bytes > MAX_BUFFER_BYTES && segments.length > 0) bytes -= segments.shift()?.bytes ?? 0
},
maskAllInputs: config.maskAllInputs,
blockSelector: BLOCK_SELECTOR,
...(config.maskAllText ? { maskTextSelector: "*" } : undefined),
checkoutEveryNms: BUFFER_CHECKOUT_MS,
})

return {
drain: async (keepalive = false) => {
const pending = segments
segments = []
bytes = 0
if (stopped) return
// Each call claims its chunk seq synchronously, so these stay ahead of
// whatever the streaming recorder that follows uploads.
await Promise.all(
pending.map((segment) =>
uploadChunk(
config,
sessionId,
`[${segment.parts.join(",")}]`,
{
isCheckpoint: true,
eventCount: segment.parts.length,
durationMs: Math.max(0, segment.last - segment.first),
},
keepalive,
),
),
)
},
stop: () => {
stopped = true
segments = []
bytes = 0
stopRecord?.()
},
getClickCount: () => clickCount,
}
}
19 changes: 15 additions & 4 deletions packages/browser-session/src/session/lifecycle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -59,8 +59,11 @@ export interface SessionSuspendOptions {

/** The parts of the lifecycle each capture mode owns. */
export interface SessionLifecycleHooks {
/** Whether the rows this owner posts are accompanied by an rrweb recording. */
readonly recorded: boolean
/**
* Whether the rows this owner posts are accompanied by an rrweb recording.
* A function when it can change mid-run: a buffered session becomes recorded when an error happens.
*/
readonly recorded: boolean | (() => boolean)
/** POST one metadata row. Best-effort — must never throw. */
readonly post: (row: Record<string, unknown>, keepalive: boolean) => void
/**
Expand All @@ -83,6 +86,8 @@ export interface SessionLifecycleHooks {

export interface SessionLifecycleHandle {
readonly sessionId: string
/** Post a fresh `active` row now, e.g. after the session became recorded. */
readonly announce: () => void
readonly shutdown: (options?: { readonly flush?: boolean }) => Promise<void>
}

Expand Down Expand Up @@ -112,6 +117,8 @@ export function startSessionLifecycle(
let errorCountBase = 0
let sinkClickCountAtStart = 0
let sinkErrorCountAtStart = 0
const isRecorded = (): boolean =>
typeof hooks.recorded === "function" ? hooks.recorded() : hooks.recorded

/**
* The persisted record of the session this lifecycle owns.
Expand Down Expand Up @@ -196,7 +203,8 @@ export function startSessionLifecycle(
pageViews: counts.pageViews,
errorCount: counts.errorCount,
traceIds: status === "ended" ? options.getTraceIds?.(record.id) : undefined,
recorded: hooks.recorded,
recorded: isRecorded(),
replayTrigger: record.replayTrigger,
billableStart,
}),
keepalive,
Expand All @@ -215,7 +223,7 @@ export function startSessionLifecycle(
const record = liveRecord()
// A session minted by idle rotation mid-page has no sampling decision yet;
// it takes this page's mode so its later loads agree with it.
adoptReplayDecision(record.id, hooks.recorded)
adoptReplayDecision(record.id, isRecorded())
rebaseCounts(record)
hooks.onStart?.(record)
post("active", false)
Expand Down Expand Up @@ -301,6 +309,9 @@ export function startSessionLifecycle(
get sessionId() {
return current.id
},
announce: () => {
if (running) post("active", false)
},
shutdown: async (shutdownOptions) => {
if (stopped) return
stopped = true
Expand Down
Loading
Loading