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
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type { AddEventResult, EventBuffer, EventBufferType, RecordingEvent } fro
import { debug } from '../util/logger';
import { EventBufferArray } from './EventBufferArray';
import { EventBufferCompressionWorker } from './EventBufferCompressionWorker';
import { WorkerDestroyedError } from './error';

/**
* This proxy will try to use the compression worker, and fall back to use the simple buffer if an error occurs there.
Expand Down Expand Up @@ -130,6 +131,11 @@ export class EventBufferProxy implements EventBuffer {
// Can now clear fallback buffer as it's no longer necessary
this._fallback.clear();
} catch (error) {
// Destroying the worker (e.g. when the session expires) rejects the
// in-flight requests. This is expected teardown, not a failure.
if (error instanceof WorkerDestroyedError) {
return;
}
DEBUG_BUILD && debug.exception(error, 'Failed to add events when switching buffers.');
}
}
Expand Down
3 changes: 2 additions & 1 deletion packages/replay-internal/src/eventBuffer/WorkerHandler.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { DEBUG_BUILD } from '../debug-build';
import type { WorkerRequest, WorkerResponse } from '../types';
import { debug } from '../util/logger';
import { WorkerDestroyedError } from './error';

interface PendingRequest {
method: WorkerRequest['method'];
Expand Down Expand Up @@ -75,7 +76,7 @@ export class WorkerHandler {
public destroy(): void {
DEBUG_BUILD && debug.log('Destroying compression worker');
this._worker.removeEventListener('message', this._onMessage);
this._pending.forEach(pending => pending.reject(new Error('Worker destroyed')));
this._pending.forEach(pending => pending.reject(new WorkerDestroyedError()));
this._pending.clear();
this._worker.terminate();
}
Expand Down
7 changes: 7 additions & 0 deletions packages/replay-internal/src/eventBuffer/error.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,3 +6,10 @@ export class EventBufferSizeExceededError extends Error {
super(`Event buffer exceeded maximum size of ${REPLAY_MAX_EVENT_BUFFER_SIZE}.`);
}
}

/** This error indicates that the compression worker was intentionally destroyed (e.g. on session expiry). */
export class WorkerDestroyedError extends Error {
public constructor() {
super('Worker destroyed');
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -6,23 +6,55 @@ import 'jsdom-worker';
import type { MockInstance } from 'vitest';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { EventBufferProxy } from '../../../src/eventBuffer/EventBufferProxy';
import { debug } from '../../../src/util/logger';
import { BASE_TIMESTAMP } from '../..';
import { decompress } from '../../utils/compression';
import { getTestEventIncremental } from '../../utils/getTestEvent';
import { createEventBuffer } from './../../../src/eventBuffer';

const TEST_EVENT = getTestEventIncremental({ timestamp: BASE_TIMESTAMP });

/**
* Worker stub that only answers when the test tells it to, so the buffer can be
* destroyed while the switch to the compression worker is still in flight.
*/
class ControlledWorker extends EventTarget {
public posted: Array<{ id: number; method: string }> = [];

public postMessage(data: unknown): void {
this.posted.push(data as { id: number; method: string });
}

public terminate(): void {
// noop
}

/** Emit the message the worker sends once its script has loaded. */
public sendReady(): void {
this.dispatchEvent(new MessageEvent('message', { data: { success: true } }));
}

/** Answer all posted requests with an unsuccessful response. */
public failAll(): void {
this.posted.forEach(({ id, method }) => {
this.dispatchEvent(new MessageEvent('message', { data: { id, method, success: false } }));
});
}
}

describe('Unit | eventBuffer | EventBufferProxy', () => {
let consoleErrorSpy: MockInstance<any>;
let exceptionSpy: MockInstance<any>;

beforeEach(() => {
// Avoid logging errors to console
consoleErrorSpy = vi.spyOn(console, 'error').mockImplementation(() => {});
exceptionSpy = vi.spyOn(debug, 'exception').mockImplementation(() => {});
});

afterEach(() => {
consoleErrorSpy.mockRestore();
exceptionSpy.mockRestore();
});

it('waits for the worker to be loaded when calling finish', async function () {
Expand Down Expand Up @@ -67,4 +99,34 @@ describe('Unit | eventBuffer | EventBufferProxy', () => {
expect(typeof result2).toBe('string');
expect(result2).toEqual(JSON.stringify([TEST_EVENT, TEST_EVENT, TEST_EVENT]));
});

it('does not report an error if the worker is destroyed while switching buffers', async function () {
const worker = new ControlledWorker();
const buffer = new EventBufferProxy(worker as unknown as Worker);

await buffer.addEvent(TEST_EVENT);

worker.sendReady();
await vi.waitFor(() => expect(worker.posted).toHaveLength(1));

buffer.destroy();

await buffer.ensureWorkerIsLoaded();
expect(exceptionSpy).not.toHaveBeenCalled();
});

it('reports an error if adding events fails while switching buffers', async function () {
const worker = new ControlledWorker();
const buffer = new EventBufferProxy(worker as unknown as Worker);

await buffer.addEvent(TEST_EVENT);

worker.sendReady();
await vi.waitFor(() => expect(worker.posted).toHaveLength(1));

worker.failAll();

await buffer.ensureWorkerIsLoaded();
expect(exceptionSpy).toHaveBeenCalledWith(expect.any(Error), 'Failed to add events when switching buffers.');
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
*/

import { describe, expect, it } from 'vitest';
import { WorkerDestroyedError } from '../../../src/eventBuffer/error';
import { WorkerHandler } from '../../../src/eventBuffer/WorkerHandler';
import type { WorkerResponse } from '../../../src/types';

Expand Down Expand Up @@ -166,8 +167,8 @@ describe('Unit | eventBuffer | WorkerHandler', () => {

handler.destroy();

await expect(p1).rejects.toThrow('Worker destroyed');
await expect(p2).rejects.toThrow('Worker destroyed');
await expect(p1).rejects.toThrow(WorkerDestroyedError);
await expect(p2).rejects.toThrow(WorkerDestroyedError);
expect(worker.terminated).toBe(true);
expect(worker.listenerCount).toBe(0);
});
Expand Down
Loading