perf(workers): zero-copy file content transfer via transferList (U19)

Pool dispatch now hoists `{path, content: string}[]` file contents OUT
of the U17 JSON envelope into separately-allocated `Uint8Array`s whose
ArrayBuffers are passed to `worker.postMessage`'s `transferList` for
zero-copy ownership transfer. The envelope itself carries only
lightweight metadata (`{path, byteLength}` per file) and is structure-
cloned the same as before.

What this saves vs U17 baseline:

- **JSON.stringify of file contents on main thread** drops to zero —
  the envelope is now O(paths + sizes), not O(total bytes). For a 200-
  file sub-batch of 10 KB TS files, that's ~2 MB of escape processing
  per dispatch that disappears. JSON.stringify's per-character branch
  on quotes/backslashes/control chars is roughly 2x slower than
  UTF-8 transcode in TextEncoder, so the replacement is a CPU win
  even though it adds a single TextEncoder.encode per file.
- **Structured-clone memcpy of file contents** drops to zero — the
  contents' backing ArrayBuffers are ownership-transferred, not copied
  into the worker's heap. The envelope's struct-clone cost is now
  proportional to metadata size only.
- **JSON.parse on worker thread** likewise no longer scales with
  content size. Worker decodes each `Uint8Array` to string via
  `TextDecoder` lazily at the parse boundary — runs on the worker
  thread, parallel with continued main-thread work, vs U17's
  sequential JSON.parse blocking the worker before processBatch can
  start.

Pipelining: TextEncoder.encode (main) and TextDecoder.decode (worker)
can both run while the OTHER side is doing useful work. Under U17,
struct-clone was a synchronous main-thread blocker.

The ArrayBuffer ownership contract is load-bearing:

- File-content `Uint8Array`s are allocated via `TextEncoder.encode`,
  NOT `Buffer.from(str, 'utf8')`. TextEncoder produces a dedicated
  ArrayBuffer per call; `Buffer.from(str)` carves from Node's shared
  `Buffer.poolSize` slab for small strings, so transferring one
  pool-backed Buffer's ArrayBuffer would detach every other Buffer
  that shares that slab — silent data corruption.
- The envelope itself is NOT transferred. It MAY be pool-backed by
  `encodeMessage`, and at ~30-80 bytes/file the struct-clone cost is
  negligible. Not transferring avoids the same detach-collateral risk
  the contents path is careful to dodge.

Detection is strict: every input element must have both `path: string`
and `content: string`. A single non-conforming element disqualifies
the whole batch from the transfer path and falls back to the legacy
single-Uint8Array `encodeMessage` envelope. Safer than partial
transfer (which would split a sub-batch into mixed-shape messages
the worker can't reassemble).

`parse-worker.ts` `decodeIncomingMessage` recognizes the hybrid
`{envelope, contents}` shape, decodes the envelope, zips metadata
positionally with the contents array, decodes UTF-8 → string per file,
and hands the reassembled `ParseWorkerInput[]` to the existing
`processBatch`. Identical downstream behavior to U17 — the IPC
optimization is invisible above this line.

Test scaffolding (3 FakeWorkers + 1 integration-test preamble) gain a
`decodeDispatchedMessage` helper that tolerates BOTH shapes (legacy
single-frame Uint8Array AND the new hybrid envelope+contents) so the
in-process unit mocks keep their existing action-scripting API and the
9 ad-hoc integration test workers keep their `msg.type === 'sub-batch'`
handlers unchanged.

`buildDispatchMessage` is now exported from worker-pool.ts so its
contract can be tested in isolation. A new
`test/unit/worker-pool-transferlist.test.ts` pins:
  - hybrid shape produced for parse-worker inputs
  - transferList carries one ArrayBuffer per file in input order
  - envelope decodes to metadata only (no `content` field)
  - content bytes round-trip byte-for-byte through UTF-8 (ASCII,
    multi-byte, surrogate-pair emoji)
  - each content's ArrayBuffer is independently allocated (no pool
    sharing) — the load-bearing transfer-safety invariant
  - non-parse shapes, empty arrays, and mixed-conformance arrays all
    fall back to the legacy single-frame path

All 271 test files (6166 unit + integration tests) pass.
This commit is contained in:
Gergo Magyar 2026-05-20 12:00:19 +01:00
parent 7744e75b77
commit 27d750b86d
7 changed files with 391 additions and 9 deletions

View file

@ -2485,6 +2485,36 @@ function decodeIncomingMessage(raw: unknown): WorkerIncomingMessage {
if (raw instanceof Uint8Array) {
return decodeMessage(raw).payload as WorkerIncomingMessage;
}
// U19: hybrid envelope+contents shape. Pool dispatch hoists file
// contents OUT of the JSON envelope and into a `contents: Uint8Array[]`
// companion whose ArrayBuffers were transferred zero-copy. The
// envelope (also a Uint8Array, structured-cloned not transferred)
// carries `{type:'sub-batch', files:[{path, byteLength}]}` — metadata
// only. Reassemble by zipping the metadata files with the contents
// array positionally, decoding UTF-8 → string at this point (one
// toString per file, executed on the worker thread in parallel with
// ongoing main-thread work).
if (
raw !== null &&
typeof raw === 'object' &&
(raw as { envelope?: unknown }).envelope instanceof Uint8Array &&
Array.isArray((raw as { contents?: unknown }).contents)
) {
const envelope = (raw as { envelope: Uint8Array }).envelope;
const contents = (raw as { contents: Uint8Array[] }).contents;
const decoded = decodeMessage(envelope).payload as {
type: string;
files: Array<{ path: string; byteLength: number }>;
};
if (decoded.type === 'sub-batch' && Array.isArray(decoded.files)) {
const decoder = new TextDecoder('utf-8');
const files: ParseWorkerInput[] = decoded.files.map((meta, i) => ({
path: meta.path,
content: decoder.decode(contents[i]),
}));
return { type: 'sub-batch', files };
}
}
return raw as WorkerIncomingMessage;
}

View file

@ -27,6 +27,72 @@ function decodeIncomingWorkerMessage(raw: unknown): WorkerOutgoingMessage {
}
return raw as WorkerOutgoingMessage;
}
/**
* U19: zero-copy dispatch builder.
*
* For the parse-worker shape `{path: string, content: string}[]`, hoists
* file contents OUT of the JSON envelope into separate Uint8Arrays and
* returns a `transferList` of their ArrayBuffers so `worker.postMessage`
* transfers ownership zero-copy. The envelope itself carries only
* lightweight metadata (path + byteLength per file), so the JSON
* round-trip cost no longer scales with total file-content size.
*
* For non-parse shapes (test scaffolding sending arbitrary items), falls
* back to the legacy `encodeMessage` path with the full payload inside
* the JSON envelope — no transfer, no shape assumptions.
*
* Ownership notes:
* - Each content `Uint8Array` is produced via `TextEncoder.encode`,
* which allocates its own ArrayBuffer. This is intentional: Node's
* `Buffer.from(str, 'utf8')` and `Buffer.alloc(size)` may carve out
* of the shared `Buffer.poolSize` pool, and transferring a pool-
* backed ArrayBuffer would detach every other Buffer that happens
* to share the same pool slab. TextEncoder bypasses the pool.
* - The envelope (`encodeMessage` output) is NOT transferred. It MAY
* be pool-backed; structured-cloning it is cheap (envelope is
* ~30-80 bytes per file, dominated by path strings, no content),
* and avoiding transfer here means we can't accidentally detach an
* unrelated Buffer that shares the pool.
*/
export function buildDispatchMessage<T>(items: readonly T[]): {
message: Uint8Array | { envelope: Uint8Array; contents: Uint8Array[] };
transferList?: ArrayBuffer[];
} {
const isParseWorkerShape =
items.length > 0 &&
items.every(
(it) =>
it != null &&
typeof it === 'object' &&
typeof (it as { path?: unknown }).path === 'string' &&
typeof (it as { content?: unknown }).content === 'string',
);
if (!isParseWorkerShape) {
return {
message: encodeMessage(MessageTag.DispatchJob, { type: 'sub-batch', files: items }),
};
}
const encoder = new TextEncoder();
const filesMeta: Array<{ path: string; byteLength: number }> = [];
const contents: Uint8Array[] = [];
for (const item of items as unknown as Array<{ path: string; content: string }>) {
const u8 = encoder.encode(item.content);
filesMeta.push({ path: item.path, byteLength: u8.byteLength });
contents.push(u8);
}
const envelope = encodeMessage(MessageTag.DispatchJob, {
type: 'sub-batch',
files: filesMeta,
});
const transferList: ArrayBuffer[] = contents.map((c) => c.buffer as ArrayBuffer);
return {
message: { envelope, contents },
transferList,
};
}
export interface WorkerPool {
/**
* Dispatch items across workers. Items are split into bounded jobs, each job
@ -1269,9 +1335,12 @@ export const createWorkerPool = (
cleanup();
return;
}
worker.postMessage(
encodeMessage(MessageTag.DispatchJob, { type: 'sub-batch', files: job.items }),
);
const { message, transferList } = buildDispatchMessage(job.items);
if (transferList) {
worker.postMessage(message, transferList);
} else {
worker.postMessage(message);
}
};
for (const slotIndex of activeSlots) runWorker(slotIndex);

View file

@ -40,7 +40,7 @@ const hasDistWorker = fs.existsSync(DIST_WORKER);
// 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', ...)`
// 2. (U17/U19) 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
@ -51,15 +51,45 @@ const hasDistWorker = fs.existsSync(DIST_WORKER);
// 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).
//
// U19 adds a second incoming shape: the hybrid
// `{envelope: Uint8Array, contents: Uint8Array[]}` produced by
// `buildDispatchMessage` when file contents are transferred
// zero-copy. The wrapper detects this shape, decodes the envelope
// header to get the file metadata, and zips the metadata with the
// transferred Uint8Arrays — decoding each content back to a UTF-8
// string so test workers see the legacy POJO layout
// `{type:'sub-batch', files:[{path, content: string}]}`.
const READY_PREAMBLE = `
const { parentPort: __pp } = require('node:worker_threads');
const __decodeFrame = (raw) => {
const __decodeProtocolBuf = (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 __decoder = new TextDecoder('utf-8');
const __decodeFrame = (raw) => {
if (raw instanceof Uint8Array) return __decodeProtocolBuf(raw);
if (
raw && typeof raw === 'object' &&
raw.envelope instanceof Uint8Array &&
Array.isArray(raw.contents)
) {
const env = __decodeProtocolBuf(raw.envelope);
if (env && env.type === 'sub-batch' && Array.isArray(env.files)) {
return {
type: 'sub-batch',
files: env.files.map((m, i) => ({
path: m.path,
content: __decoder.decode(raw.contents[i]),
})),
};
}
return env;
}
return raw;
};
const __origOn = __pp.on.bind(__pp);
__pp.on = (event, handler) => {
if (event !== 'message') return __origOn(event, handler);

View file

@ -12,6 +12,45 @@ import {
} from '../../src/core/ingestion/workers/worker-pool.js';
import { decodeMessage } from '../../src/core/ingestion/workers/protocol.js';
/**
* Decode whatever shape the pool sent — supports both:
* - U17 single-frame: a Uint8Array protocol frame (payload inside JSON)
* - U19 hybrid: `{envelope: Uint8Array, contents: Uint8Array[]}` where
* file contents were hoisted out of JSON into a transferList. Test
* action logic only inspects `msg.files[*].path`, so contents are
* decoded back to strings for shape parity with the legacy POJO.
*/
function decodeDispatchedMessage(rawMsg: unknown): unknown {
if (rawMsg instanceof Uint8Array) {
return decodeMessage(rawMsg).payload;
}
if (
rawMsg !== null &&
typeof rawMsg === 'object' &&
(rawMsg as { envelope?: unknown }).envelope instanceof Uint8Array &&
Array.isArray((rawMsg as { contents?: unknown }).contents)
) {
const env = (rawMsg as { envelope: Uint8Array }).envelope;
const contents = (rawMsg as { contents: Uint8Array[] }).contents;
const decoded = decodeMessage(env).payload as {
type: string;
files: Array<{ path: string; byteLength: number }>;
};
if (decoded.type === 'sub-batch' && Array.isArray(decoded.files)) {
const decoder = new TextDecoder('utf-8');
return {
type: 'sub-batch',
files: decoded.files.map((m, i) => ({
path: m.path,
content: decoder.decode(contents[i]),
})),
};
}
return decoded;
}
return rawMsg;
}
/**
* Minimal `node:worker_threads` Worker double for unit-testing the pool's
* resilience layers (auto-respawn, circuit breaker, quarantine, retry
@ -53,7 +92,7 @@ class FakeWorker extends EventEmitter {
// `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;
const msg = decodeDispatchedMessage(rawMsg);
this.seenMessages.push(msg);
if (typeof msg !== 'object' || msg === null) return;
const m = msg as { type?: string; files?: { path: string }[] };

View file

@ -29,6 +29,35 @@ import os from 'node:os';
import { createWorkerPool } from '../../src/core/ingestion/workers/worker-pool.js';
import { decodeMessage } from '../../src/core/ingestion/workers/protocol.js';
function decodeDispatchedMessage(rawMsg: unknown): unknown {
if (rawMsg instanceof Uint8Array) return decodeMessage(rawMsg).payload;
if (
rawMsg !== null &&
typeof rawMsg === 'object' &&
(rawMsg as { envelope?: unknown }).envelope instanceof Uint8Array &&
Array.isArray((rawMsg as { contents?: unknown }).contents)
) {
const env = (rawMsg as { envelope: Uint8Array }).envelope;
const contents = (rawMsg as { contents: Uint8Array[] }).contents;
const decoded = decodeMessage(env).payload as {
type: string;
files: Array<{ path: string; byteLength: number }>;
};
if (decoded.type === 'sub-batch') {
const decoder = new TextDecoder('utf-8');
return {
type: 'sub-batch',
files: decoded.files.map((m, i) => ({
path: m.path,
content: decoder.decode(contents[i]),
})),
};
}
return decoded;
}
return rawMsg;
}
type FakeAction =
| { kind: 'crash-after-starting'; startingPath: string; code: number }
| { kind: 'parse-ok'; files: { path: string }[] };
@ -45,7 +74,7 @@ class FakeWorker extends EventEmitter {
}
postMessage(rawMsg: unknown): void {
// U17: decode Buffer-encoded dispatches; pool is now strict-encoded.
const msg = rawMsg instanceof Uint8Array ? decodeMessage(rawMsg).payload : rawMsg;
const msg = decodeDispatchedMessage(rawMsg);
if (typeof msg !== 'object' || msg === null) return;
const m = msg as { type?: string };
if (m.type !== 'sub-batch') return;

View file

@ -0,0 +1,156 @@
/**
* U19 — Zero-copy transferList dispatch builder.
*
* `worker-pool.ts`'s `buildDispatchMessage` is the U19 boundary between the
* pool's generic `dispatch<T>(items)` and the parse-worker-specific
* postMessage payload shape. For items shaped as `{path, content: string}[]`
* (the parse-worker contract), file contents are hoisted OUT of the JSON
* envelope into separately-allocated `Uint8Array`s whose ArrayBuffers go
* into the `transferList` for zero-copy ownership transfer. For any other
* shape, the builder falls back to the legacy `encodeMessage` path with
* the full payload inside a single JSON-encoded protocol frame.
*
* These tests pin the contract:
* - parse-worker shape produces hybrid envelope + transferList
* - non-parse shape stays on the legacy single-Uint8Array path
* - content bytes round-trip byte-for-byte through encode+decode
* - each content buffer is independently allocated (not shared via
* Node's `Buffer.poolSize` slab) so transferring one cannot detach
* another
* - empty items array returns the legacy path (no shape inference on
* zero elements)
*/
import { describe, it, expect } from 'vitest';
import {
buildDispatchMessage,
// Internal type — re-export not needed; just exercise observable behaviour.
} from '../../src/core/ingestion/workers/worker-pool.js';
import { decodeMessage } from '../../src/core/ingestion/workers/protocol.js';
describe('worker pool — buildDispatchMessage (U19)', () => {
it('parse-worker shape returns hybrid envelope + transferList of one buffer per file', () => {
const items = [
{ path: 'a.ts', content: 'export const A = 1;' },
{ path: 'b.ts', content: 'export const B = 2;' },
];
const { message, transferList } = buildDispatchMessage(items);
// Hybrid shape: not a bare Uint8Array.
expect(message).not.toBeInstanceOf(Uint8Array);
expect(message).toEqual(
expect.objectContaining({
envelope: expect.any(Uint8Array),
contents: expect.any(Array),
}),
);
const hybrid = message as { envelope: Uint8Array; contents: Uint8Array[] };
expect(hybrid.contents).toHaveLength(2);
expect(hybrid.contents[0]).toBeInstanceOf(Uint8Array);
expect(hybrid.contents[1]).toBeInstanceOf(Uint8Array);
// transferList carries one ArrayBuffer per file, in the same order
// as `contents`. Identity check is the strict contract — transferring
// a different ArrayBuffer reference would no-op the ownership swap.
expect(transferList).toHaveLength(2);
expect(transferList?.[0]).toBe(hybrid.contents[0].buffer);
expect(transferList?.[1]).toBe(hybrid.contents[1].buffer);
});
it('envelope decodes to {type:"sub-batch", files:[{path, byteLength}]} metadata with NO content field', () => {
const items = [{ path: 'src/foo.ts', content: 'hello world' }];
const { message } = buildDispatchMessage(items);
const { envelope } = message as { envelope: Uint8Array };
const decoded = decodeMessage(envelope).payload as {
type: string;
files: Array<{ path: string; byteLength: number; content?: unknown }>;
};
expect(decoded.type).toBe('sub-batch');
expect(decoded.files).toHaveLength(1);
expect(decoded.files[0].path).toBe('src/foo.ts');
expect(decoded.files[0].byteLength).toBe(11); // 'hello world' = 11 UTF-8 bytes
// Content must NOT appear in the JSON envelope — that's the whole
// point of the hybrid shape. If a future refactor accidentally
// duplicates content into both envelope and transferList, this
// assertion catches it.
expect('content' in decoded.files[0]).toBe(false);
});
it('content bytes round-trip byte-for-byte through TextDecoder', () => {
// Mix ASCII, multi-byte UTF-8 (café = c-a-f-é where é is 2 bytes),
// and an emoji (4 UTF-8 bytes) to cover the encoder boundaries.
const items = [
{ path: 'a.ts', content: 'plain ASCII' },
{ path: 'b.ts', content: 'café au lait' },
{ path: 'c.ts', content: 'rocket: 🚀 emoji' },
];
const { message } = buildDispatchMessage(items);
const { contents } = message as { contents: Uint8Array[] };
const decoder = new TextDecoder('utf-8');
expect(decoder.decode(contents[0])).toBe('plain ASCII');
expect(decoder.decode(contents[1])).toBe('café au lait');
expect(decoder.decode(contents[2])).toBe('rocket: 🚀 emoji');
});
it('each content buffer owns a dedicated ArrayBuffer (no shared Buffer pool slab)', () => {
// Pin the transfer-safety contract: TextEncoder allocates each
// Uint8Array on its own ArrayBuffer, so transferring one cannot
// detach the backing of another. If a future refactor swaps to
// `Buffer.from(str, 'utf8')` (which carves from `Buffer.poolSize`
// slabs for small strings), small files would share an
// ArrayBuffer and transferList would detach unrelated content.
// This test allocates many small files — the slab would normally
// batch them — and verifies each ArrayBuffer is distinct.
const items = Array.from({ length: 8 }, (_, i) => ({
path: `f${i}.ts`,
content: `tiny ${i}`,
}));
const { message } = buildDispatchMessage(items);
const { contents } = message as { contents: Uint8Array[] };
const buffers = new Set(contents.map((c) => c.buffer));
expect(buffers.size).toBe(8);
// Each content's view covers the entire ArrayBuffer (no offset).
for (const c of contents) {
expect(c.byteOffset).toBe(0);
expect(c.byteLength).toBe(c.buffer.byteLength);
}
});
it('non-parse-worker shape falls back to the legacy single-frame path with no transferList', () => {
// Items lacking a string `content` field don't match the
// parse-worker shape detector and must round-trip through the
// existing encodeMessage path unchanged.
const items = [{ id: 1, payload: 'arbitrary' }];
const { message, transferList } = buildDispatchMessage(items);
expect(message).toBeInstanceOf(Uint8Array);
expect(transferList).toBeUndefined();
// The whole input array is embedded inside the JSON envelope on
// this path, so a decode recovers it verbatim under `files`.
const decoded = decodeMessage(message as Uint8Array).payload as {
type: string;
files: typeof items;
};
expect(decoded.type).toBe('sub-batch');
expect(decoded.files).toEqual(items);
});
it('empty items array falls back to the legacy path (no shape inference on zero elements)', () => {
// The shape detector requires at least one element so empty
// dispatches don't get false-positively routed through the
// transferList builder.
const { message, transferList } = buildDispatchMessage([]);
expect(message).toBeInstanceOf(Uint8Array);
expect(transferList).toBeUndefined();
});
it('mixed-shape items (some missing content) fall back to the legacy path', () => {
// Strict shape detection: every element must have a string content.
// A single non-conforming element disqualifies the transfer path —
// safer than partially transferring some and embedding others.
const items = [{ path: 'a.ts', content: 'ok' }, { path: 'b.ts' /* no content */ }];
const { message, transferList } = buildDispatchMessage(items);
expect(message).toBeInstanceOf(Uint8Array);
expect(transferList).toBeUndefined();
});
});

View file

@ -28,6 +28,35 @@ import os from 'node:os';
import { createWorkerPool } from '../../src/core/ingestion/workers/worker-pool.js';
import { decodeMessage } from '../../src/core/ingestion/workers/protocol.js';
function decodeDispatchedMessage(rawMsg: unknown): unknown {
if (rawMsg instanceof Uint8Array) return decodeMessage(rawMsg).payload;
if (
rawMsg !== null &&
typeof rawMsg === 'object' &&
(rawMsg as { envelope?: unknown }).envelope instanceof Uint8Array &&
Array.isArray((rawMsg as { contents?: unknown }).contents)
) {
const env = (rawMsg as { envelope: Uint8Array }).envelope;
const contents = (rawMsg as { contents: Uint8Array[] }).contents;
const decoded = decodeMessage(env).payload as {
type: string;
files: Array<{ path: string; byteLength: number }>;
};
if (decoded.type === 'sub-batch') {
const decoder = new TextDecoder('utf-8');
return {
type: 'sub-batch',
files: decoded.files.map((m, i) => ({
path: m.path,
content: decoder.decode(contents[i]),
})),
};
}
return decoded;
}
return rawMsg;
}
/**
* Minimal FakeWorker for this test: emit `starting-file` for the script's
* configured path, then either exit (death → quarantine the in-flight
@ -51,7 +80,7 @@ class FakeWorker extends EventEmitter {
}
postMessage(rawMsg: unknown): void {
// U17: decode Buffer-encoded dispatches; pool is now strict-encoded.
const msg = rawMsg instanceof Uint8Array ? decodeMessage(rawMsg).payload : rawMsg;
const msg = decodeDispatchedMessage(rawMsg);
if (typeof msg !== 'object' || msg === null) return;
const m = msg as { type?: string };
if (m.type !== 'sub-batch') return;