From 7744e75b779440e15abea2d6ad27658a4b77148f Mon Sep 17 00:00:00 2001 From: Gergo Magyar Date: Wed, 20 May 2026 11:30:27 +0100 Subject: [PATCH] feat(workers): wire protocol.ts encoded IPC into parse-worker + pool (U17) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../core/ingestion/workers/parse-worker.ts | 70 +++++++++--- .../src/core/ingestion/workers/protocol.ts | 24 ++++- .../src/core/ingestion/workers/worker-pool.ts | 65 +++++++++++- gitnexus/test/integration/worker-pool.test.ts | 100 +++++++++++++----- .../test/unit/worker-pool-resilience.test.ts | 12 ++- .../unit/worker-pool-slot-generation.test.ts | 5 +- .../worker-pool-windows-quarantine.test.ts | 5 +- gitnexus/test/unit/workers/protocol.test.ts | 41 +++++++ 8 files changed, 266 insertions(+), 56 deletions(-) diff --git a/gitnexus/src/core/ingestion/workers/parse-worker.ts b/gitnexus/src/core/ingestion/workers/parse-worker.ts index a096cf597..cbdca83fc 100644 --- a/gitnexus/src/core/ingestion/workers/parse-worker.ts +++ b/gitnexus/src/core/ingestion/workers/parse-worker.ts @@ -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