diff --git a/gitnexus/src/core/ingestion/workers/parse-worker.ts b/gitnexus/src/core/ingestion/workers/parse-worker.ts index cbdca83fc..cd1307693 100644 --- a/gitnexus/src/core/ingestion/workers/parse-worker.ts +++ b/gitnexus/src/core/ingestion/workers/parse-worker.ts @@ -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; } diff --git a/gitnexus/src/core/ingestion/workers/worker-pool.ts b/gitnexus/src/core/ingestion/workers/worker-pool.ts index 790b0dc0b..38a7e7651 100644 --- a/gitnexus/src/core/ingestion/workers/worker-pool.ts +++ b/gitnexus/src/core/ingestion/workers/worker-pool.ts @@ -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(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); diff --git a/gitnexus/test/integration/worker-pool.test.ts b/gitnexus/test/integration/worker-pool.test.ts index 5896e4ca8..04b4372fd 100644 --- a/gitnexus/test/integration/worker-pool.test.ts +++ b/gitnexus/test/integration/worker-pool.test.ts @@ -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); diff --git a/gitnexus/test/unit/worker-pool-resilience.test.ts b/gitnexus/test/unit/worker-pool-resilience.test.ts index 6857aaf76..e2dea546e 100644 --- a/gitnexus/test/unit/worker-pool-resilience.test.ts +++ b/gitnexus/test/unit/worker-pool-resilience.test.ts @@ -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 }[] }; diff --git a/gitnexus/test/unit/worker-pool-slot-generation.test.ts b/gitnexus/test/unit/worker-pool-slot-generation.test.ts index 4919dbe90..6b71635df 100644 --- a/gitnexus/test/unit/worker-pool-slot-generation.test.ts +++ b/gitnexus/test/unit/worker-pool-slot-generation.test.ts @@ -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; diff --git a/gitnexus/test/unit/worker-pool-transferlist.test.ts b/gitnexus/test/unit/worker-pool-transferlist.test.ts new file mode 100644 index 000000000..b48eb619a --- /dev/null +++ b/gitnexus/test/unit/worker-pool-transferlist.test.ts @@ -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(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(); + }); +}); diff --git a/gitnexus/test/unit/worker-pool-windows-quarantine.test.ts b/gitnexus/test/unit/worker-pool-windows-quarantine.test.ts index ebd2ca9c5..6ba903224 100644 --- a/gitnexus/test/unit/worker-pool-windows-quarantine.test.ts +++ b/gitnexus/test/unit/worker-pool-windows-quarantine.test.ts @@ -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;