212 lines
6.9 KiB
TypeScript
212 lines
6.9 KiB
TypeScript
import { describe, expect, it, vi } from 'vitest';
|
|
import {
|
|
aggregatePcmPeaks,
|
|
createWaveformRuntime,
|
|
WaveformCancellationError,
|
|
WaveformRuntimeBusyError,
|
|
WaveformRuntimeTerminatedError,
|
|
type WaveformWorkerLike,
|
|
} from '../../src/waveform';
|
|
import type { WaveformWorkerAggregateRequest } from '../../src/waveform/waveform.types';
|
|
|
|
describe('waveform worker runtime', () => {
|
|
it('offloads aggregation, reports progress, and preserves caller bytes', async () => {
|
|
const worker = new FakeWorker();
|
|
const input = new Int16Array([-32_768, 0, 16_384, 32_767]);
|
|
const originalBytes = [...new Uint8Array(input.buffer)];
|
|
const progress = vi.fn();
|
|
const runtime = createWaveformRuntime({ workerFactory: () => worker });
|
|
|
|
const resultPromise = runtime.aggregate(
|
|
{
|
|
container: 'raw',
|
|
data: input,
|
|
encoding: 's16le',
|
|
sampleRate: 8_000,
|
|
},
|
|
{ bucketCount: 2, onProgress: progress }
|
|
);
|
|
const request = worker.lastRequest;
|
|
expect(request).toBeDefined();
|
|
expect(request?.input.data).not.toBe(input.buffer);
|
|
expect([...new Uint8Array(input.buffer)]).toEqual(originalBytes);
|
|
|
|
worker.emitMessage({
|
|
type: 'progress',
|
|
requestId: request?.requestId,
|
|
progress: 0.5,
|
|
});
|
|
const peaks = aggregatePcmPeaks(
|
|
{
|
|
container: 'raw',
|
|
data: request?.input.data ?? new ArrayBuffer(0),
|
|
encoding: 's16le',
|
|
sampleRate: 8_000,
|
|
},
|
|
{ bucketCount: 2 }
|
|
);
|
|
worker.emitMessage({
|
|
type: 'result',
|
|
requestId: request?.requestId,
|
|
peaks,
|
|
});
|
|
|
|
await expect(resultPromise).resolves.toMatchObject({ bucketCount: 2 });
|
|
expect(progress).toHaveBeenCalledWith(0.5);
|
|
expect(progress).toHaveBeenLastCalledWith(1);
|
|
expect(worker.terminated).toBe(true);
|
|
expect(runtime.isActive).toBe(false);
|
|
});
|
|
|
|
it('terminates the worker on AbortSignal cancellation and remains reusable', async () => {
|
|
const workers: FakeWorker[] = [];
|
|
const runtime = createWaveformRuntime({
|
|
workerFactory: () => {
|
|
const worker = new FakeWorker();
|
|
workers.push(worker);
|
|
return worker;
|
|
},
|
|
});
|
|
const controller = new AbortController();
|
|
const first = runtime.aggregate(rawInput(), {
|
|
bucketCount: 2,
|
|
signal: controller.signal,
|
|
});
|
|
controller.abort();
|
|
|
|
await expect(first).rejects.toBeInstanceOf(WaveformCancellationError);
|
|
expect(workers[0]?.terminated).toBe(true);
|
|
expect(runtime.isActive).toBe(false);
|
|
|
|
const second = runtime.aggregate(rawInput(), { bucketCount: 2 });
|
|
const request = workers[1]?.lastRequest;
|
|
const peaks = aggregatePcmPeaks(rawInput(), { bucketCount: 2 });
|
|
workers[1]?.emitMessage({
|
|
type: 'result',
|
|
requestId: request?.requestId,
|
|
peaks,
|
|
});
|
|
await expect(second).resolves.toMatchObject({ bucketCount: 2 });
|
|
});
|
|
|
|
it('supports explicit cancellation and rejects overlapping work', async () => {
|
|
const worker = new FakeWorker();
|
|
const runtime = createWaveformRuntime({ workerFactory: () => worker });
|
|
const active = runtime.aggregate(rawInput(), { bucketCount: 2 });
|
|
|
|
await expect(
|
|
runtime.aggregate(rawInput(), { bucketCount: 2 })
|
|
).rejects.toBeInstanceOf(WaveformRuntimeBusyError);
|
|
expect(runtime.cancelActive()).toBe(true);
|
|
await expect(active).rejects.toBeInstanceOf(WaveformCancellationError);
|
|
expect(runtime.cancelActive()).toBe(false);
|
|
});
|
|
|
|
it('falls back synchronously when workers are unavailable or blocked', async () => {
|
|
const runtime = createWaveformRuntime({ workerFactory: null });
|
|
await expect(
|
|
runtime.aggregate(rawInput(), { bucketCount: 2 })
|
|
).resolves.toMatchObject({ bucketCount: 2 });
|
|
|
|
const blockedRuntime = createWaveformRuntime({
|
|
workerFactory: () => {
|
|
throw new Error('CSP blocked worker');
|
|
},
|
|
});
|
|
await expect(
|
|
blockedRuntime.aggregate(rawInput(), { bucketCount: 2 })
|
|
).resolves.toMatchObject({ bucketCount: 2 });
|
|
expect(blockedRuntime.mode).toBe('synchronous');
|
|
});
|
|
|
|
it('tracks and can cancel synchronous fallback before it starts', async () => {
|
|
const runtime = createWaveformRuntime({ workerFactory: null });
|
|
const active = runtime.aggregate(rawInput(), { bucketCount: 2 });
|
|
|
|
expect(runtime.isActive).toBe(true);
|
|
const overlapping = runtime.aggregate(rawInput(), { bucketCount: 2 });
|
|
expect(runtime.cancelActive()).toBe(true);
|
|
await expect(overlapping).rejects.toBeInstanceOf(WaveformRuntimeBusyError);
|
|
await expect(active).rejects.toBeInstanceOf(WaveformCancellationError);
|
|
expect(runtime.isActive).toBe(false);
|
|
});
|
|
|
|
it('makes terminate permanent and cancels active work', async () => {
|
|
const worker = new FakeWorker();
|
|
const runtime = createWaveformRuntime({ workerFactory: () => worker });
|
|
const active = runtime.aggregate(rawInput(), { bucketCount: 2 });
|
|
runtime.terminate();
|
|
|
|
await expect(active).rejects.toBeInstanceOf(WaveformCancellationError);
|
|
await expect(
|
|
runtime.aggregate(rawInput(), { bucketCount: 2 })
|
|
).rejects.toBeInstanceOf(WaveformRuntimeTerminatedError);
|
|
});
|
|
});
|
|
|
|
class FakeWorker implements WaveformWorkerLike {
|
|
readonly #messageListeners = new Set<
|
|
(event: MessageEvent<unknown>) => void
|
|
>();
|
|
readonly #errorListeners = new Set<(event: Event) => void>();
|
|
readonly #messageErrorListeners = new Set<(event: Event) => void>();
|
|
lastRequest: WaveformWorkerAggregateRequest | undefined;
|
|
terminated = false;
|
|
|
|
postMessage(message: unknown): void {
|
|
this.lastRequest = message as WaveformWorkerAggregateRequest;
|
|
}
|
|
|
|
addEventListener(
|
|
type: 'message' | 'error' | 'messageerror',
|
|
listener:
|
|
((event: MessageEvent<unknown>) => void) | ((event: Event) => void)
|
|
): void {
|
|
if (type === 'message') {
|
|
this.#messageListeners.add(
|
|
listener as (event: MessageEvent<unknown>) => void
|
|
);
|
|
} else if (type === 'error') {
|
|
this.#errorListeners.add(listener as (event: Event) => void);
|
|
} else {
|
|
this.#messageErrorListeners.add(listener as (event: Event) => void);
|
|
}
|
|
}
|
|
|
|
removeEventListener(
|
|
type: 'message' | 'error' | 'messageerror',
|
|
listener:
|
|
((event: MessageEvent<unknown>) => void) | ((event: Event) => void)
|
|
): void {
|
|
if (type === 'message') {
|
|
this.#messageListeners.delete(
|
|
listener as (event: MessageEvent<unknown>) => void
|
|
);
|
|
} else if (type === 'error') {
|
|
this.#errorListeners.delete(listener as (event: Event) => void);
|
|
} else {
|
|
this.#messageErrorListeners.delete(listener as (event: Event) => void);
|
|
}
|
|
}
|
|
|
|
terminate(): void {
|
|
this.terminated = true;
|
|
}
|
|
|
|
emitMessage(data: unknown): void {
|
|
const event = { data } as MessageEvent<unknown>;
|
|
for (const listener of this.#messageListeners) {
|
|
listener(event);
|
|
}
|
|
}
|
|
}
|
|
|
|
function rawInput() {
|
|
return {
|
|
container: 'raw' as const,
|
|
data: new Int16Array([-100, 100, -200, 200]),
|
|
encoding: 's16le' as const,
|
|
sampleRate: 8_000,
|
|
};
|
|
}
|