mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-11 03:38:07 +00:00
feat(workers): wire protocol.ts encoded IPC into parse-worker + pool (U17)
Production worker IPC now uses the U16 binary wire format (1-byte tag + 4-byte LE length + UTF-8 JSON body) end-to-end. The pool encodes every outgoing `sub-batch` / `flush` dispatch via `encodeMessage`; the worker decodes incoming frames via `decodeMessage` and encodes its `ready`, `starting-file`, `progress`, `sub-batch-done`, `result`, `warning`, and `error` outputs the same way. The load-bearing correctness fix is making `decodeMessage` accept `Uint8Array` rather than only `Buffer`: Node's `worker_threads` `postMessage` structured-clones the payload, which strips the `Buffer` prototype on the receive side. A frame sent as `Buffer` arrives as a plain `Uint8Array`, and `Buffer.isBuffer(raw)` returns false — so the first attempt at U17 (gating decode on `Buffer.isBuffer`) silently treated every incoming frame as POJO and the worker never responded. The fix adopts the underlying memory zero-copy via `Buffer.from(view.buffer, view.byteOffset, view.byteLength)` and uses `raw instanceof Uint8Array` at every call site (parse-worker decode, pool dispatch handler, pool ready-handshake handler, FakeWorker test mocks, and the integration-test worker preamble). The pool stays tolerant of POJO incoming so unit-test FakeWorkers don't need rewriting — only the new outgoing encoded dispatches require the test scaffolding to decode on receive, which the test FakeWorkers and the integration test's inline `parentPort.on` wrapper now do. The slot-drop integration test was rewritten from a shared-counter-file race (which pre-U17 timing happened to land on the assertion-friendly counter==2 endpoint, but post-U17 protocol decoding latency shifted to counter==1 and produced 3 quarantines instead of 2) to a deterministic path-based crash trigger: slot 0 crashes on a.ts, respawns, crashes on the requeued b.ts, slot is dropped after budget exhausted; slot 1 handles [c.ts, d.ts] normally. Outcome no longer depends on inter-worker file-write ordering. Protocol coverage adds two regression tests pinning the Uint8Array decode path: structured-clone-stripped frames decode identically to their Buffer originals, and Uint8Array views with non-zero byteOffset into a wider ArrayBuffer also decode correctly (catches `Buffer.from(uint8)` copying semantics if a future refactor loses the zero-copy adoption). All 94 worker-pool tests (9 files, unit + integration) pass; the full unit suite (6128 tests across 268 files) passes unchanged.
This commit is contained in:
parent
832d97678f
commit
7744e75b77
8 changed files with 266 additions and 56 deletions
|
|
@ -1,5 +1,6 @@
|
|||
import { parentPort } from 'node:worker_threads';
|
||||
import Parser from 'tree-sitter';
|
||||
import { encodeMessage, decodeMessage, MessageTag } from './protocol.js';
|
||||
import JavaScript from 'tree-sitter-javascript';
|
||||
import TypeScript from 'tree-sitter-typescript';
|
||||
import Python from 'tree-sitter-python';
|
||||
|
|
@ -1390,7 +1391,7 @@ const processFileGroup = (
|
|||
} catch (err) {
|
||||
const message = `Query compilation failed for ${language}: ${err instanceof Error ? err.message : String(err)}`;
|
||||
if (parentPort) {
|
||||
parentPort.postMessage({ type: 'warning', message });
|
||||
parentPort.postMessage(encodeMessage(MessageTag.Warning, { type: 'warning', message }));
|
||||
} else {
|
||||
logger.warn(message);
|
||||
}
|
||||
|
|
@ -1406,7 +1407,11 @@ const processFileGroup = (
|
|||
// guessing from `items[lastProgress]` (which the language-grouped order
|
||||
// here would defeat). The pool gracefully ignores this when running an
|
||||
// older worker build that doesn't emit it.
|
||||
if (parentPort) parentPort.postMessage({ type: 'starting-file', path: file.path });
|
||||
if (parentPort) {
|
||||
parentPort.postMessage(
|
||||
encodeMessage(MessageTag.StartingFile, { type: 'starting-file', path: file.path }),
|
||||
);
|
||||
}
|
||||
|
||||
// Vue SFC preprocessing: extract <script> block content
|
||||
let parseContent = file.content;
|
||||
|
|
@ -1465,8 +1470,11 @@ const processFileGroup = (
|
|||
parseContent,
|
||||
file.path,
|
||||
(message) => {
|
||||
if (parentPort) parentPort.postMessage({ type: 'warning', message });
|
||||
else logger.warn(message);
|
||||
if (parentPort) {
|
||||
parentPort.postMessage(encodeMessage(MessageTag.Warning, { type: 'warning', message }));
|
||||
} else {
|
||||
logger.warn(message);
|
||||
}
|
||||
},
|
||||
tree,
|
||||
);
|
||||
|
|
@ -2453,37 +2461,69 @@ const mergeResult = (target: ParseWorkerResult, src: ParseWorkerResult) => {
|
|||
// before the script body runs) and the pool only notices via the first
|
||||
// dispatch's idle timeout (~30s). Emit once; the dispatch handler treats
|
||||
// any subsequent `ready` message as a benign no-op.
|
||||
parentPort!.postMessage({ type: 'ready' });
|
||||
//
|
||||
// Post-U17: every postMessage call uses `encodeMessage` from protocol.ts
|
||||
// so the bytes on the wire carry an explicit tag + length header. The pool
|
||||
// tolerates both encoded Buffer messages and raw POJO (for backward compat
|
||||
// with FakeWorkers in the test suite), but production parse-worker.ts is
|
||||
// strict — every outgoing message is encoded.
|
||||
parentPort!.postMessage(encodeMessage(MessageTag.Ready, { type: 'ready' }));
|
||||
|
||||
parentPort!.on('message', (msg: WorkerIncomingMessage) => {
|
||||
// Decode a single incoming message. The pool always sends Buffer-encoded
|
||||
// frames post-U17, but the parameter type is `unknown` because Node's
|
||||
// worker_threads typings declare `on('message', (value: any) => void)`.
|
||||
//
|
||||
// `instanceof Uint8Array` (not `Buffer.isBuffer`) is load-bearing: Node's
|
||||
// worker_threads `postMessage` structured-clones the payload, which
|
||||
// strips the Buffer prototype — a frame sent as a Buffer arrives here
|
||||
// as a plain Uint8Array, where `Buffer.isBuffer` returns false.
|
||||
// `decodeMessage` already adopts a Uint8Array view zero-copy.
|
||||
//
|
||||
// Falling back to the raw value for non-Uint8Array inputs preserves the
|
||||
// legacy POJO path some tests might still exercise.
|
||||
function decodeIncomingMessage(raw: unknown): WorkerIncomingMessage {
|
||||
if (raw instanceof Uint8Array) {
|
||||
return decodeMessage(raw).payload as WorkerIncomingMessage;
|
||||
}
|
||||
return raw as WorkerIncomingMessage;
|
||||
}
|
||||
|
||||
parentPort!.on('message', (raw: unknown) => {
|
||||
try {
|
||||
const msg = decodeIncomingMessage(raw);
|
||||
// Legacy single-message mode (backward compat): array of files
|
||||
if (Array.isArray(msg)) {
|
||||
const result = processBatch(msg, (filesProcessed) => {
|
||||
parentPort!.postMessage({ type: 'progress', filesProcessed });
|
||||
parentPort!.postMessage(
|
||||
encodeMessage(MessageTag.Progress, { type: 'progress', filesProcessed }),
|
||||
);
|
||||
});
|
||||
parentPort!.postMessage({ type: 'result', data: result });
|
||||
parentPort!.postMessage(encodeMessage(MessageTag.Result, { type: 'result', data: result }));
|
||||
return;
|
||||
}
|
||||
|
||||
// Sub-batch mode: { type: 'sub-batch', files: [...] }
|
||||
if (msg.type === 'sub-batch') {
|
||||
const result = processBatch(msg.files, (filesProcessed) => {
|
||||
parentPort!.postMessage({
|
||||
type: 'progress',
|
||||
filesProcessed: cumulativeProcessed + filesProcessed,
|
||||
});
|
||||
parentPort!.postMessage(
|
||||
encodeMessage(MessageTag.Progress, {
|
||||
type: 'progress',
|
||||
filesProcessed: cumulativeProcessed + filesProcessed,
|
||||
}),
|
||||
);
|
||||
});
|
||||
cumulativeProcessed += result.fileCount;
|
||||
mergeResult(accumulated, result);
|
||||
// Signal ready for next sub-batch
|
||||
parentPort!.postMessage({ type: 'sub-batch-done' });
|
||||
parentPort!.postMessage(encodeMessage(MessageTag.SubBatchDone, { type: 'sub-batch-done' }));
|
||||
return;
|
||||
}
|
||||
|
||||
// Flush: send accumulated results
|
||||
if (msg.type === 'flush') {
|
||||
parentPort!.postMessage({ type: 'result', data: accumulated });
|
||||
parentPort!.postMessage(
|
||||
encodeMessage(MessageTag.Result, { type: 'result', data: accumulated }),
|
||||
);
|
||||
// Reset for potential reuse
|
||||
accumulated = {
|
||||
nodes: [],
|
||||
|
|
@ -2509,6 +2549,6 @@ parentPort!.on('message', (msg: WorkerIncomingMessage) => {
|
|||
}
|
||||
} catch (err) {
|
||||
const message = err instanceof Error ? err.message : String(err);
|
||||
parentPort!.postMessage({ type: 'error', error: message });
|
||||
parentPort!.postMessage(encodeMessage(MessageTag.Error, { type: 'error', error: message }));
|
||||
}
|
||||
});
|
||||
|
|
|
|||
|
|
@ -113,10 +113,21 @@ export function encodeMessage(tag: MessageTagValue, payload: unknown): Buffer {
|
|||
}
|
||||
|
||||
/**
|
||||
* Decode a single message from a Buffer. The buffer must start with a
|
||||
* complete protocol frame; trailing bytes beyond the declared length
|
||||
* are ignored (callers receiving a concatenated stream should slice at
|
||||
* `PROTOCOL_HEADER_BYTES + length` before decoding the next frame).
|
||||
* Decode a single message from a Buffer (or any Uint8Array containing a
|
||||
* frame). The buffer must start with a complete protocol frame; trailing
|
||||
* bytes beyond the declared length are ignored (callers receiving a
|
||||
* concatenated stream should slice at `PROTOCOL_HEADER_BYTES + length`
|
||||
* before decoding the next frame).
|
||||
*
|
||||
* Accepts `Uint8Array` rather than only `Buffer` because Node's
|
||||
* worker_threads `postMessage` uses structured clone, which strips the
|
||||
* `Buffer` prototype: a `Buffer` sent over the wire arrives on the
|
||||
* receiving thread as a plain `Uint8Array`. Buffer extends Uint8Array,
|
||||
* so when the input is already a Buffer the readUInt / subarray fast
|
||||
* paths still apply; when the input is a bare Uint8Array, we adopt its
|
||||
* underlying memory via `Buffer.from(view.buffer, view.byteOffset,
|
||||
* view.byteLength)` (a zero-copy view, not a clone) so the rest of the
|
||||
* decode runs through the same code path.
|
||||
*
|
||||
* Throws {@link ProtocolDecodeError} for any of:
|
||||
* - buffer shorter than the 5-byte header
|
||||
|
|
@ -124,10 +135,13 @@ export function encodeMessage(tag: MessageTagValue, payload: unknown): Buffer {
|
|||
* - declared payload length exceeds available bytes
|
||||
* - payload bytes are not valid UTF-8 JSON
|
||||
*/
|
||||
export function decodeMessage(buf: Buffer): {
|
||||
export function decodeMessage(input: Uint8Array): {
|
||||
tag: MessageTagValue;
|
||||
payload: unknown;
|
||||
} {
|
||||
const buf: Buffer = Buffer.isBuffer(input)
|
||||
? input
|
||||
: Buffer.from(input.buffer, input.byteOffset, input.byteLength);
|
||||
if (buf.length < PROTOCOL_HEADER_BYTES) {
|
||||
throw new ProtocolDecodeError(
|
||||
`frame too small for header: got ${buf.length} bytes, need ${PROTOCOL_HEADER_BYTES}`,
|
||||
|
|
|
|||
|
|
@ -5,6 +5,28 @@ import { fileURLToPath } from 'node:url';
|
|||
|
||||
import { logger } from '../../logger.js';
|
||||
import { createQuarantine } from './quarantine.js';
|
||||
import { encodeMessage, decodeMessage, MessageTag } from './protocol.js';
|
||||
|
||||
/**
|
||||
* Decode an incoming worker message. Production `parse-worker.ts` (post-U17)
|
||||
* always sends Buffer-encoded frames via `encodeMessage`. Test scaffolding
|
||||
* (FakeWorkers in worker-pool-resilience.test.ts and siblings) still emits
|
||||
* raw POJO messages — the pool tolerates both so the existing in-process
|
||||
* mock workers don't need rewriting. If a future commit wants to strictly
|
||||
* require encoded frames (rejecting POJO), this helper is the single point
|
||||
* to harden.
|
||||
*/
|
||||
function decodeIncomingWorkerMessage(raw: unknown): WorkerOutgoingMessage {
|
||||
// `instanceof Uint8Array` (not `Buffer.isBuffer`): Node's worker_threads
|
||||
// postMessage structured-clones the payload, which strips the Buffer
|
||||
// prototype, so a frame sent as `Buffer` arrives here as a plain
|
||||
// Uint8Array. Buffer extends Uint8Array so this check covers both,
|
||||
// and `decodeMessage` adopts a Uint8Array view zero-copy.
|
||||
if (raw instanceof Uint8Array) {
|
||||
return decodeMessage(raw).payload as WorkerOutgoingMessage;
|
||||
}
|
||||
return raw as WorkerOutgoingMessage;
|
||||
}
|
||||
export interface WorkerPool {
|
||||
/**
|
||||
* Dispatch items across workers. Items are split into bounded jobs, each job
|
||||
|
|
@ -319,7 +341,20 @@ function waitForWorkerReady(worker: Worker): Promise<void> {
|
|||
worker.removeListener('exit', onExit);
|
||||
worker.removeListener('messageerror', onMessageError);
|
||||
};
|
||||
const onMessage = (msg: unknown) => {
|
||||
const onMessage = (raw: unknown) => {
|
||||
// U17: production parse-worker.ts emits the ready handshake as a
|
||||
// Buffer-encoded frame; FakeWorkers in the test suite still emit
|
||||
// POJO. Tolerate both. ProtocolDecodeError on a malformed frame is
|
||||
// swallowed locally — the wait keeps listening, and the eventual
|
||||
// timeout / exit / error handlers catch a genuinely-broken worker.
|
||||
let msg: unknown = raw;
|
||||
if (raw instanceof Uint8Array) {
|
||||
try {
|
||||
msg = decodeMessage(raw).payload;
|
||||
} catch {
|
||||
return;
|
||||
}
|
||||
}
|
||||
if (typeof msg === 'object' && msg !== null && (msg as { type?: unknown }).type === 'ready') {
|
||||
cleanup();
|
||||
resolve();
|
||||
|
|
@ -1092,9 +1127,29 @@ export const createWorkerPool = (
|
|||
// stale generation. The guard catches future-refactor mistakes.
|
||||
const slotGen = slotGenerations[workerIndex];
|
||||
|
||||
const handler = (msg: WorkerOutgoingMessage) => {
|
||||
const handler = (raw: unknown) => {
|
||||
if (slotGenerations[workerIndex] !== slotGen) return;
|
||||
if (settled || stopped) return;
|
||||
// U17: production parse-worker.ts emits Buffer-encoded messages;
|
||||
// FakeWorkers in the test suite still emit POJOs. Tolerate both
|
||||
// via `decodeIncomingWorkerMessage`. A malformed frame from a
|
||||
// real worker is treated as a worker-side bug — the protocol
|
||||
// error escapes and is caught by the `messageerror` handler
|
||||
// below, routing through the existing recovery layer.
|
||||
let msg: WorkerOutgoingMessage;
|
||||
try {
|
||||
msg = decodeIncomingWorkerMessage(raw);
|
||||
} catch (err) {
|
||||
settled = true;
|
||||
cleanup();
|
||||
void recoverAndResume(
|
||||
`Worker ${workerIndex} protocol decode error: ${
|
||||
err instanceof Error ? err.message : String(err)
|
||||
}`,
|
||||
resolveExcludePaths(),
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (msg.type === 'starting-file') {
|
||||
inFlightPath = msg.path;
|
||||
resetIdleTimer();
|
||||
|
|
@ -1111,7 +1166,7 @@ export const createWorkerPool = (
|
|||
} else if (msg.type === 'sub-batch-done') {
|
||||
waitingForFlush = true;
|
||||
resetIdleTimer();
|
||||
worker.postMessage({ type: 'flush' });
|
||||
worker.postMessage(encodeMessage(MessageTag.DispatchJob, { type: 'flush' }));
|
||||
} else if (msg.type === 'error') {
|
||||
settled = true;
|
||||
cleanup();
|
||||
|
|
@ -1214,7 +1269,9 @@ export const createWorkerPool = (
|
|||
cleanup();
|
||||
return;
|
||||
}
|
||||
worker.postMessage({ type: 'sub-batch', files: job.items });
|
||||
worker.postMessage(
|
||||
encodeMessage(MessageTag.DispatchJob, { type: 'sub-batch', files: job.items }),
|
||||
);
|
||||
};
|
||||
|
||||
for (const slotIndex of activeSlots) runWorker(slotIndex);
|
||||
|
|
|
|||
|
|
@ -30,13 +30,43 @@ const DIST_WORKER = path.resolve(
|
|||
);
|
||||
const hasDistWorker = fs.existsSync(DIST_WORKER);
|
||||
|
||||
// Prepend the M4 ready handshake to every ad-hoc test worker source so the
|
||||
// pool's `waitForWorkerReady` resolves immediately for replacement spawns.
|
||||
// Production `parse-worker.ts` emits the same handshake at top-of-script
|
||||
// before installing its message handler. Without it, every test that triggers
|
||||
// a replacement (worker crash + recover) would hit the 5s WORKER_READY_TIMEOUT_MS
|
||||
// and fail with "Replacement worker startup failed and no slots remain".
|
||||
const READY_PREAMBLE = `require('node:worker_threads').parentPort.postMessage({ type: 'ready' });\n`;
|
||||
// Prepend two things to every ad-hoc test worker source:
|
||||
//
|
||||
// 1. The M4 ready handshake so the pool's `waitForWorkerReady` resolves
|
||||
// immediately for replacement spawns. Production `parse-worker.ts`
|
||||
// emits the same handshake at top-of-script before installing its
|
||||
// message handler. Without it, every test that triggers a
|
||||
// replacement (worker crash + recover) would hit the 5s
|
||||
// WORKER_READY_TIMEOUT_MS and fail with "Replacement worker startup
|
||||
// failed and no slots remain".
|
||||
//
|
||||
// 2. (U17) A transparent decode wrapper around `parentPort.on('message', ...)`
|
||||
// so the 9 ad-hoc test worker scripts don't need rewriting now that
|
||||
// the pool sends Buffer-encoded `sub-batch` / `flush` dispatches.
|
||||
// Test workers continue to receive POJO objects in their existing
|
||||
// `msg.type === 'sub-batch'` checks. The inline mini-decoder mirrors
|
||||
// protocol.ts's wire layout (1-byte tag + 4-byte LE uint32 length +
|
||||
// UTF-8 JSON body) — duplicated here because ad-hoc test workers
|
||||
// `require` Node builtins only; they can't import the compiled
|
||||
// dist/protocol.js without knowing its absolute path. The pool is
|
||||
// tolerant of POJO incoming, so outgoing `parentPort.postMessage`
|
||||
// calls in the test scripts stay as POJO (unchanged).
|
||||
const READY_PREAMBLE = `
|
||||
const { parentPort: __pp } = require('node:worker_threads');
|
||||
const __decodeFrame = (raw) => {
|
||||
// structured clone strips the Buffer prototype; raw arrives as Uint8Array.
|
||||
if (!(raw instanceof Uint8Array)) return raw;
|
||||
const buf = Buffer.from(raw.buffer, raw.byteOffset, raw.byteLength);
|
||||
const length = buf.readUInt32LE(1);
|
||||
return JSON.parse(buf.subarray(5, 5 + length).toString('utf8'));
|
||||
};
|
||||
const __origOn = __pp.on.bind(__pp);
|
||||
__pp.on = (event, handler) => {
|
||||
if (event !== 'message') return __origOn(event, handler);
|
||||
return __origOn(event, (raw) => handler(__decodeFrame(raw)));
|
||||
};
|
||||
__pp.postMessage({ type: 'ready' });
|
||||
`;
|
||||
|
||||
function writeReadyWorker(workerPath: string, source: string): void {
|
||||
fs.writeFileSync(workerPath, READY_PREAMBLE + source);
|
||||
|
|
@ -892,27 +922,45 @@ describe('worker pool integration', () => {
|
|||
});
|
||||
|
||||
it('drops a slot after maxRespawnsPerSlot and continues on the survivor', async () => {
|
||||
// 2-worker pool, budget=1. Worker A crashes twice on its assigned
|
||||
// chunk; slot A is dropped. Slot B handles the requeued remainder
|
||||
// successfully. Validates the per-slot drop + wakeIdleSlots flow
|
||||
// under real worker timing.
|
||||
// 2-worker pool, budget=1. The worker that's assigned the chunk
|
||||
// containing `a.ts` crashes on a.ts (quarantines a), respawns, gets
|
||||
// the requeued remainder containing `b.ts`, crashes on b.ts
|
||||
// (quarantines b). That slot's respawn budget is now exhausted, so
|
||||
// the pool drops it. The OTHER slot — assigned the chunk with
|
||||
// [c,d] — never sees the poison files and completes its work
|
||||
// normally. Validates the per-slot drop + wakeIdleSlots flow under
|
||||
// real worker timing.
|
||||
//
|
||||
// The path-based crash trigger replaces an earlier shared-counter-
|
||||
// file design that was a write-write race between the two workers:
|
||||
// pre-U17 timing happened to land on counter=2 by the end of
|
||||
// round 1 (so round-2 workers saw counter==2 and didn't crash),
|
||||
// but the post-U17 protocol-decoding latency shifted the window
|
||||
// so round 2's first worker read counter=1 and crashed too,
|
||||
// producing 3 quarantines instead of 2. Switching to a path-based
|
||||
// trigger removes the inter-worker race entirely — the outcome
|
||||
// depends only on which chunk contains the poison files, which is
|
||||
// deterministic given the dispatch ordering of [a,b,c,d] with
|
||||
// subBatchSize=2.
|
||||
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-resilience-slot-drop-'));
|
||||
const counterPath = path.join(tempDir, 'crash-count.txt');
|
||||
const workerPath = path.join(tempDir, 'worker.js');
|
||||
writeReadyWorker(
|
||||
workerPath,
|
||||
`
|
||||
const fs = require('node:fs');
|
||||
const { parentPort } = require('node:worker_threads');
|
||||
const counterPath = ${JSON.stringify(counterPath)};
|
||||
let current = [];
|
||||
parentPort.on('message', (msg) => {
|
||||
if (msg && msg.type === 'sub-batch') {
|
||||
current = msg.files.map((f) => f.path);
|
||||
const prior = fs.existsSync(counterPath) ? Number(fs.readFileSync(counterPath, 'utf-8')) : 0;
|
||||
if (prior < 2) {
|
||||
fs.writeFileSync(counterPath, String(prior + 1));
|
||||
parentPort.postMessage({ type: 'starting-file', path: current[0] });
|
||||
// Crash deterministically on the poison files. The pool
|
||||
// filters quarantined paths from subsequent re-dispatches,
|
||||
// so the first crash quarantines a.ts and the requeue then
|
||||
// contains b.ts; the second crash quarantines b.ts and the
|
||||
// slot's respawn budget is exhausted. Worker handling the
|
||||
// [c,d] chunk never enters this branch.
|
||||
const poison = current.find((p) => p === 'a.ts' || p === 'b.ts');
|
||||
if (poison) {
|
||||
parentPort.postMessage({ type: 'starting-file', path: poison });
|
||||
process.exit(134);
|
||||
}
|
||||
parentPort.postMessage({ type: 'progress', filesProcessed: current.length });
|
||||
|
|
@ -941,18 +989,12 @@ describe('worker pool integration', () => {
|
|||
{ path: 'd.ts', content: '' },
|
||||
]);
|
||||
|
||||
// Two crashes both attributed to their starting-file (a.ts, then
|
||||
// whatever the requeued chunk's items[0] was — b.ts or c.ts
|
||||
// depending on dispatch order). At minimum a.ts is quarantined.
|
||||
const quarantine = pool.getQuarantinedPaths?.() ?? [];
|
||||
expect(quarantine.length).toBe(2);
|
||||
expect(quarantine).toContain('a.ts');
|
||||
// All non-quarantined files eventually parsed.
|
||||
// Deterministically: a.ts crashes round 1, b.ts crashes round 2.
|
||||
const quarantine = (pool.getQuarantinedPaths?.() ?? []).sort();
|
||||
expect(quarantine).toEqual(['a.ts', 'b.ts']);
|
||||
// All non-quarantined files eventually parsed by the survivor slot.
|
||||
const allPaths = results.flatMap((r) => r.paths).sort();
|
||||
const expectedAllowed = ['a.ts', 'b.ts', 'c.ts', 'd.ts'].filter(
|
||||
(p) => !quarantine.includes(p),
|
||||
);
|
||||
expect(allPaths.sort()).toEqual(expectedAllowed.sort());
|
||||
expect(allPaths).toEqual(['c.ts', 'd.ts']);
|
||||
} finally {
|
||||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||||
}
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@ import {
|
|||
resolveWorkerPoolOptions,
|
||||
resolveAutoPoolSize,
|
||||
} from '../../src/core/ingestion/workers/worker-pool.js';
|
||||
import { decodeMessage } from '../../src/core/ingestion/workers/protocol.js';
|
||||
|
||||
/**
|
||||
* Minimal `node:worker_threads` Worker double for unit-testing the pool's
|
||||
|
|
@ -43,7 +44,16 @@ class FakeWorker extends EventEmitter {
|
|||
});
|
||||
}
|
||||
|
||||
postMessage(msg: unknown): void {
|
||||
postMessage(rawMsg: unknown): void {
|
||||
// U17: production pool now sends Buffer-encoded dispatch frames.
|
||||
// Decode them here so this in-process mock can keep its existing
|
||||
// POJO-shaped action-scripting API — the action queue still sees
|
||||
// `{type, files}` shapes regardless of whether the pool encoded
|
||||
// the message on the way in. Store the DECODED payload in
|
||||
// `seenMessages` so test-side introspection assertions (which
|
||||
// expect `msg.type` / `msg.files`) keep working after the wire
|
||||
// format flipped to Buffer.
|
||||
const msg = Buffer.isBuffer(rawMsg) ? decodeMessage(rawMsg).payload : rawMsg;
|
||||
this.seenMessages.push(msg);
|
||||
if (typeof msg !== 'object' || msg === null) return;
|
||||
const m = msg as { type?: string; files?: { path: string }[] };
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ import { pathToFileURL } from 'node:url';
|
|||
import fs from 'node:fs';
|
||||
import os from 'node:os';
|
||||
import { createWorkerPool } from '../../src/core/ingestion/workers/worker-pool.js';
|
||||
import { decodeMessage } from '../../src/core/ingestion/workers/protocol.js';
|
||||
|
||||
type FakeAction =
|
||||
| { kind: 'crash-after-starting'; startingPath: string; code: number }
|
||||
|
|
@ -42,7 +43,9 @@ class FakeWorker extends EventEmitter {
|
|||
this.emit('message', { type: 'ready' });
|
||||
});
|
||||
}
|
||||
postMessage(msg: unknown): void {
|
||||
postMessage(rawMsg: unknown): void {
|
||||
// U17: decode Buffer-encoded dispatches; pool is now strict-encoded.
|
||||
const msg = rawMsg instanceof Uint8Array ? decodeMessage(rawMsg).payload : rawMsg;
|
||||
if (typeof msg !== 'object' || msg === null) return;
|
||||
const m = msg as { type?: string };
|
||||
if (m.type !== 'sub-batch') return;
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ import { pathToFileURL } from 'node:url';
|
|||
import fs from 'node:fs';
|
||||
import os from 'node:os';
|
||||
import { createWorkerPool } from '../../src/core/ingestion/workers/worker-pool.js';
|
||||
import { decodeMessage } from '../../src/core/ingestion/workers/protocol.js';
|
||||
|
||||
/**
|
||||
* Minimal FakeWorker for this test: emit `starting-file` for the script's
|
||||
|
|
@ -48,7 +49,9 @@ class FakeWorker extends EventEmitter {
|
|||
this.emit('message', { type: 'ready' });
|
||||
});
|
||||
}
|
||||
postMessage(msg: unknown): void {
|
||||
postMessage(rawMsg: unknown): void {
|
||||
// U17: decode Buffer-encoded dispatches; pool is now strict-encoded.
|
||||
const msg = rawMsg instanceof Uint8Array ? decodeMessage(rawMsg).payload : rawMsg;
|
||||
if (typeof msg !== 'object' || msg === null) return;
|
||||
const m = msg as { type?: string };
|
||||
if (m.type !== 'sub-batch') return;
|
||||
|
|
|
|||
|
|
@ -133,3 +133,44 @@ describe('worker IPC protocol — decode error paths (U16)', () => {
|
|||
}
|
||||
});
|
||||
});
|
||||
|
||||
describe('worker IPC protocol — Uint8Array decode path (U17)', () => {
|
||||
// Pins the U17 production fix: Node's worker_threads `postMessage`
|
||||
// structured-clones the payload, which strips the Buffer prototype, so
|
||||
// a frame sent as Buffer arrives on the receiver as a bare Uint8Array.
|
||||
// `decodeMessage` must accept a Uint8Array view zero-copy. Without this
|
||||
// path, every Buffer-encoded message sent through worker_threads would
|
||||
// silently fail to decode and the worker would never reply to the
|
||||
// dispatch — observed pre-fix as a 30s test timeout on every real
|
||||
// dispatch through dist/parse-worker.js.
|
||||
it('decodes a Uint8Array view (no Buffer prototype) identically to the Buffer original', () => {
|
||||
const original = encodeMessage(MessageTag.Result, { fileCount: 5, paths: ['a.ts', 'b.ts'] });
|
||||
// Strip the Buffer prototype while keeping the same backing memory —
|
||||
// mirrors what structured clone does to a Buffer on the receive side.
|
||||
const stripped = new Uint8Array(original.buffer, original.byteOffset, original.byteLength);
|
||||
expect(Buffer.isBuffer(stripped)).toBe(false);
|
||||
expect(stripped).toBeInstanceOf(Uint8Array);
|
||||
|
||||
const decoded = decodeMessage(stripped);
|
||||
expect(decoded.tag).toBe(MessageTag.Result);
|
||||
expect(decoded.payload).toEqual({ fileCount: 5, paths: ['a.ts', 'b.ts'] });
|
||||
});
|
||||
|
||||
it('decodes a Uint8Array that views a slice of a larger ArrayBuffer (non-zero byteOffset)', () => {
|
||||
// Pins the zero-copy adoption path. If `decodeMessage` were to call
|
||||
// `Buffer.from(uint8)` (which copies and zeroes the offset) instead
|
||||
// of `Buffer.from(buf.buffer, buf.byteOffset, buf.byteLength)`, the
|
||||
// header parse would still succeed on the copy, but the wider
|
||||
// ArrayBuffer-with-offset case below also catches `Buffer.from(uint8)`
|
||||
// semantics regressions on Node versions that materialize structured
|
||||
// clones into shared ArrayBuffers with offsets.
|
||||
const original = encodeMessage(MessageTag.Progress, { filesProcessed: 11 });
|
||||
const padded = new Uint8Array(original.byteLength + 8);
|
||||
padded.set(original, 4); // place the frame at offset 4 in a wider buffer
|
||||
const view = new Uint8Array(padded.buffer, 4, original.byteLength);
|
||||
|
||||
const decoded = decodeMessage(view);
|
||||
expect(decoded.tag).toBe(MessageTag.Progress);
|
||||
expect(decoded.payload).toEqual({ filesProcessed: 11 });
|
||||
});
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue