228 lines
5.8 KiB
TypeScript
228 lines
5.8 KiB
TypeScript
import {
|
|
WORKER_PROTOCOL_VERSION,
|
|
type WorkerRequest,
|
|
type WorkerResponse,
|
|
} from "./worker-protocol";
|
|
|
|
export interface WorkerLike {
|
|
onmessage: ((event: MessageEvent<unknown>) => void) | null;
|
|
onerror: ((event: ErrorEvent) => void) | null;
|
|
onmessageerror: ((event: MessageEvent<unknown>) => void) | null;
|
|
postMessage(message: unknown): void;
|
|
terminate(): void;
|
|
}
|
|
|
|
export type WorkerFactory = () => WorkerLike;
|
|
|
|
export class WorkerRequestError extends Error {
|
|
readonly kind: "timeout" | "crash" | "cancelled" | "worker-error";
|
|
|
|
constructor(
|
|
kind: WorkerRequestError["kind"],
|
|
message: string,
|
|
options?: ErrorOptions,
|
|
) {
|
|
super(message, options);
|
|
this.name = "WorkerRequestError";
|
|
this.kind = kind;
|
|
}
|
|
}
|
|
|
|
interface ActiveRequest<TResult> {
|
|
readonly requestId: number;
|
|
readonly generation: number;
|
|
readonly resolve: (result: TResult) => void;
|
|
readonly reject: (error: Error) => void;
|
|
readonly timeout: ReturnType<typeof globalThis.setTimeout>;
|
|
}
|
|
|
|
export class WorkerSupervisor<TOperation, TResult> {
|
|
private readonly label: string;
|
|
private readonly workerFactory: WorkerFactory;
|
|
private worker: WorkerLike | null = null;
|
|
private generation = 0;
|
|
private nextRequestId = 1;
|
|
private active: ActiveRequest<TResult> | null = null;
|
|
private disposed = false;
|
|
|
|
constructor(label: string, workerFactory: WorkerFactory) {
|
|
this.label = label;
|
|
this.workerFactory = workerFactory;
|
|
}
|
|
|
|
get currentGeneration(): number {
|
|
return this.generation;
|
|
}
|
|
|
|
get isRunning(): boolean {
|
|
return this.active !== null;
|
|
}
|
|
|
|
get hasWorker(): boolean {
|
|
return this.worker !== null;
|
|
}
|
|
|
|
run(
|
|
operation: TOperation,
|
|
timeoutMs: number,
|
|
options: { readonly supersede?: boolean } = {},
|
|
): Promise<TResult> {
|
|
if (this.disposed) {
|
|
return Promise.reject(
|
|
new WorkerRequestError(
|
|
"worker-error",
|
|
`${this.label} supervisor has been disposed`,
|
|
),
|
|
);
|
|
}
|
|
if (this.active) {
|
|
if (!options.supersede) {
|
|
return Promise.reject(
|
|
new WorkerRequestError(
|
|
"worker-error",
|
|
`${this.label} already has an active request`,
|
|
),
|
|
);
|
|
}
|
|
this.restart(
|
|
new WorkerRequestError(
|
|
"cancelled",
|
|
`${this.label} request was superseded`,
|
|
),
|
|
);
|
|
}
|
|
const worker = this.ensureWorker();
|
|
const requestId = this.nextRequestId++;
|
|
const generation = this.generation;
|
|
return new Promise<TResult>((resolve, reject) => {
|
|
const timeout = globalThis.setTimeout(() => {
|
|
if (
|
|
this.active?.requestId !== requestId ||
|
|
this.active.generation !== generation
|
|
) {
|
|
return;
|
|
}
|
|
this.restart(
|
|
new WorkerRequestError(
|
|
"timeout",
|
|
`${this.label} exceeded its ${timeoutMs} ms time limit`,
|
|
),
|
|
);
|
|
}, timeoutMs);
|
|
this.active = {
|
|
requestId,
|
|
generation,
|
|
resolve,
|
|
reject,
|
|
timeout,
|
|
};
|
|
const message: WorkerRequest<TOperation> = {
|
|
protocolVersion: WORKER_PROTOCOL_VERSION,
|
|
requestId,
|
|
generation,
|
|
payload: operation,
|
|
};
|
|
worker.postMessage(message);
|
|
});
|
|
}
|
|
|
|
cancel(): void {
|
|
if (!this.active) return;
|
|
this.restart(
|
|
new WorkerRequestError(
|
|
"cancelled",
|
|
`${this.label} request was cancelled`,
|
|
),
|
|
);
|
|
}
|
|
|
|
dispose(): void {
|
|
if (this.disposed) return;
|
|
this.disposed = true;
|
|
this.restart(
|
|
new WorkerRequestError("cancelled", `${this.label} was disposed`),
|
|
);
|
|
}
|
|
|
|
private ensureWorker(): WorkerLike {
|
|
if (this.worker) return this.worker;
|
|
const worker = this.workerFactory();
|
|
this.generation += 1;
|
|
const generation = this.generation;
|
|
worker.onmessage = (event) => {
|
|
this.handleResponse(event.data, generation);
|
|
};
|
|
worker.onerror = (event) => {
|
|
this.handleCrash(
|
|
new WorkerRequestError(
|
|
"crash",
|
|
`${this.label} crashed: ${event.message || "unknown worker error"}`,
|
|
),
|
|
generation,
|
|
);
|
|
};
|
|
worker.onmessageerror = () => {
|
|
this.handleCrash(
|
|
new WorkerRequestError(
|
|
"crash",
|
|
`${this.label} returned an unreadable message`,
|
|
),
|
|
generation,
|
|
);
|
|
};
|
|
this.worker = worker;
|
|
return worker;
|
|
}
|
|
|
|
private handleResponse(value: unknown, generation: number): void {
|
|
const response = value as Partial<WorkerResponse<TResult>>;
|
|
if (
|
|
response.protocolVersion !== WORKER_PROTOCOL_VERSION ||
|
|
typeof response.requestId !== "number" ||
|
|
response.generation !== generation ||
|
|
!this.active ||
|
|
this.active.requestId !== response.requestId ||
|
|
this.active.generation !== generation
|
|
) {
|
|
return;
|
|
}
|
|
const active = this.active;
|
|
this.active = null;
|
|
globalThis.clearTimeout(active.timeout);
|
|
if (response.ok === true && "payload" in response) {
|
|
active.resolve(response.payload as TResult);
|
|
return;
|
|
}
|
|
const details =
|
|
"error" in response && response.error
|
|
? response.error
|
|
: { name: "Error", message: "Unknown worker failure" };
|
|
active.reject(
|
|
new WorkerRequestError(
|
|
"worker-error",
|
|
`${details.name}: ${details.message}`,
|
|
),
|
|
);
|
|
}
|
|
|
|
private handleCrash(error: WorkerRequestError, generation: number): void {
|
|
if (generation !== this.generation) return;
|
|
this.restart(error);
|
|
}
|
|
|
|
private restart(error: WorkerRequestError): void {
|
|
if (this.active) {
|
|
globalThis.clearTimeout(this.active.timeout);
|
|
this.active.reject(error);
|
|
this.active = null;
|
|
}
|
|
if (this.worker) {
|
|
this.worker.onmessage = null;
|
|
this.worker.onerror = null;
|
|
this.worker.onmessageerror = null;
|
|
this.worker.terminate();
|
|
this.worker = null;
|
|
}
|
|
}
|
|
}
|