mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-08-26 00:52:18 +03:00
Validated on a 17-PR combined board: compression-worker + colocate-standalone-esm-scope within the board's 287/287, typecheck:core clean, env-doc-sync clean. Offloads eligible sync compression engines into a bounded worker_threads pool with a strict serializable DTO boundary and fail-open on spawn/worker/timeout failure. Closes #11023. Thank you @RaviTharuma!
41 lines
1.2 KiB
TypeScript
41 lines
1.2 KiB
TypeScript
import { parentPort } from "node:worker_threads";
|
|
import {
|
|
applyCompression,
|
|
applyStackedCompression,
|
|
type StackedCompressionStep,
|
|
} from "./strategySelector.ts";
|
|
import type {
|
|
CompressionWorkerJob,
|
|
CompressionWorkerMessage,
|
|
} from "./compressionWorkerProtocol.ts";
|
|
|
|
if (!parentPort) throw new Error("compressionWorker must run in a worker thread");
|
|
parentPort.on("message", (job: CompressionWorkerJob) => {
|
|
try {
|
|
const onEngineStep = (step: StackedCompressionStep) =>
|
|
parentPort.postMessage({
|
|
id: job.id,
|
|
type: "step",
|
|
step,
|
|
} satisfies CompressionWorkerMessage);
|
|
const result =
|
|
job.mode === "stacked"
|
|
? applyStackedCompression(job.body, job.options?.config?.stackedPipeline, {
|
|
...job.options,
|
|
onEngineStep,
|
|
})
|
|
: applyCompression(job.body, job.mode, job.options);
|
|
parentPort.postMessage({
|
|
id: job.id,
|
|
type: "result",
|
|
result,
|
|
} satisfies CompressionWorkerMessage);
|
|
} catch (error) {
|
|
parentPort.postMessage({
|
|
id: job.id,
|
|
type: "error",
|
|
error: error instanceof Error ? error.message : String(error),
|
|
} satisfies CompressionWorkerMessage);
|
|
}
|
|
});
|