feat: publish Regex Tools 0.1.0
This commit is contained in:
223
src/regex/execution/WorkerSupervisor.ts
Normal file
223
src/regex/execution/WorkerSupervisor.ts
Normal file
@@ -0,0 +1,223 @@
|
||||
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: number;
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user