diff --git a/gitnexus/src/core/ingestion/workers/protocol.ts b/gitnexus/src/core/ingestion/workers/protocol.ts new file mode 100644 index 000000000..953e055f2 --- /dev/null +++ b/gitnexus/src/core/ingestion/workers/protocol.ts @@ -0,0 +1,158 @@ +/** + * Worker-thread IPC wire format (U16, scaffold for U17 migration). + * + * This module defines the binary frame for messages exchanged between + * the worker pool and `parse-worker.ts`. It is intentionally NOT wired + * into `worker-pool.ts` / `parse-worker.ts` in this commit — the + * production migration is U17. Shipping the wire-format contract first, + * as an isolated and fully-tested module, de-risks U17 by establishing + * a single source of truth for the byte layout. + * + * Wire layout — one message per buffer, 5-byte header + variable body: + * + * +---------+-----------+---------------------+ + * | tag | length | payload bytes … | + * | 1 byte | 4 bytes | | + * +---------+-----------+---------------------+ + * + * - `tag` is one of the {@link MessageTag} values (0x01..0x08). + * - `length` is a little-endian uint32 byte count for the payload + * region (excludes the header). + * - `payload` is a UTF-8 JSON-encoded value, possibly `"null"`. + * + * Why JSON for the body (vs per-shape binary encoders): the doc-review + * adversarial reviewer (A2) flagged that a true per-shape binary encoder + * for the result message — which carries nested heterogeneous + * extracted-call / import / heritage / route arrays — would be 500-1500 + * LOC and a substantial maintenance burden. The honest perf win the + * IPC repack targets is moving file CONTENTS via `ArrayBuffer` + * `transferList` (zero-copy ownership transfer for the largest single + * piece of state in any message). That win is captured by U17 layering + * `transferList` over the bulk file-content payload while keeping this + * module's framing for the surrounding metadata. If U18 benchmark data + * shows the JSON body is itself a bottleneck after U17 lands, a + * follow-up unit can swap to per-shape binary encoding behind this + * same `encodeMessage` / `decodeMessage` surface without changing the + * frame. + */ + +/** + * 1-byte type tag identifying the message shape on the wire. Values are + * stable: never re-number an existing tag. New variants append at the + * next unused byte and {@link isValidTag} below grows accordingly. + */ +export const MessageTag = { + /** main -> worker: dispatch a sub-batch of files to parse. */ + DispatchJob: 0x01, + /** worker -> main: parsed result for the sub-batch. */ + Result: 0x02, + /** worker -> main: incremental progress count. */ + Progress: 0x03, + /** worker -> main: authoritative in-flight file path for the pool's + * attribution layer. */ + StartingFile: 0x04, + /** worker -> main: sub-batch fully processed; pool may send flush. */ + SubBatchDone: 0x05, + /** worker -> main: non-fatal warning message. */ + Warning: 0x06, + /** worker -> main: fatal error (the worker is about to bail). */ + Error: 0x07, + /** worker -> main: top-of-script init complete; pool may dispatch. */ + Ready: 0x08, +} as const; + +export type MessageTagValue = (typeof MessageTag)[keyof typeof MessageTag]; + +/** Header size: 1-byte tag + 4-byte little-endian length. */ +export const PROTOCOL_HEADER_BYTES = 5; + +/** + * Thrown by {@link decodeMessage} when an incoming buffer cannot be + * parsed as a valid protocol frame. Distinguishable from other errors so + * U17's pool-side handler can route protocol violations through the + * existing `messageerror` recovery layer (U3 H1) instead of treating + * them as silent data loss. + */ +export class ProtocolDecodeError extends Error { + constructor(message: string) { + super(message); + this.name = 'ProtocolDecodeError'; + } +} + +const MIN_TAG = 0x01; +const MAX_TAG = 0x08; + +function isValidTag(byte: number): byte is MessageTagValue { + return byte >= MIN_TAG && byte <= MAX_TAG; +} + +/** + * Encode a single message into a Buffer following the wire layout. + * `payload` is JSON-stringified; pass `undefined` or `null` for messages + * with no body. The returned Buffer can be sent verbatim via + * `worker.postMessage(buf)` once U17 swaps the dispatch layer over. + * + * Throws `RangeError` if the payload encodes to more than 4 GiB (the + * uint32 length-field ceiling). In practice the structured-clone budget + * caps payloads well below this, but the check makes the boundary + * explicit instead of silently truncating. + */ +export function encodeMessage(tag: MessageTagValue, payload: unknown): Buffer { + const json = Buffer.from(JSON.stringify(payload ?? null), 'utf8'); + if (json.length > 0xffffffff) { + throw new RangeError( + `protocol payload exceeds uint32 length cap (${json.length} > ${0xffffffff})`, + ); + } + const buf = Buffer.alloc(PROTOCOL_HEADER_BYTES + json.length); + buf.writeUInt8(tag, 0); + buf.writeUInt32LE(json.length, 1); + json.copy(buf, PROTOCOL_HEADER_BYTES); + return buf; +} + +/** + * 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). + * + * Throws {@link ProtocolDecodeError} for any of: + * - buffer shorter than the 5-byte header + * - tag byte outside the valid range + * - declared payload length exceeds available bytes + * - payload bytes are not valid UTF-8 JSON + */ +export function decodeMessage(buf: Buffer): { + tag: MessageTagValue; + payload: unknown; +} { + if (buf.length < PROTOCOL_HEADER_BYTES) { + throw new ProtocolDecodeError( + `frame too small for header: got ${buf.length} bytes, need ${PROTOCOL_HEADER_BYTES}`, + ); + } + const tag = buf.readUInt8(0); + if (!isValidTag(tag)) { + throw new ProtocolDecodeError( + `unknown message tag: 0x${tag.toString(16).padStart(2, '0')} (valid range: 0x01..0x08)`, + ); + } + const length = buf.readUInt32LE(1); + const payloadEnd = PROTOCOL_HEADER_BYTES + length; + if (buf.length < payloadEnd) { + throw new ProtocolDecodeError( + `truncated payload: header declared ${length} bytes, buffer has ${buf.length - PROTOCOL_HEADER_BYTES}`, + ); + } + const json = buf.subarray(PROTOCOL_HEADER_BYTES, payloadEnd).toString('utf8'); + let payload: unknown; + try { + payload = JSON.parse(json); + } catch (err) { + const reason = err instanceof Error ? err.message : String(err); + throw new ProtocolDecodeError(`payload is not valid JSON: ${reason}`); + } + return { tag, payload }; +} diff --git a/gitnexus/test/unit/workers/protocol.test.ts b/gitnexus/test/unit/workers/protocol.test.ts new file mode 100644 index 000000000..b080ba5f5 --- /dev/null +++ b/gitnexus/test/unit/workers/protocol.test.ts @@ -0,0 +1,135 @@ +/** + * U16 — Wire-format round-trip + error-path coverage for + * `core/ingestion/workers/protocol.ts`. + * + * Module-only at this stage (U17 wires it into worker-pool.ts / + * parse-worker.ts). The tests pin: + * - exact byte layout: tag (1) + length (4 LE) + UTF-8 JSON body + * - round-trip for every message tag + * - decode errors surface as ProtocolDecodeError (not generic Error) + * so U17's pool-side handler can route them through messageerror + * recovery distinctly from other failure classes + * - boundary cases: empty body, large body + */ +import { describe, it, expect } from 'vitest'; +import { + MessageTag, + PROTOCOL_HEADER_BYTES, + ProtocolDecodeError, + encodeMessage, + decodeMessage, + type MessageTagValue, +} from '../../../src/core/ingestion/workers/protocol.js'; + +describe('worker IPC protocol — encodeMessage byte layout (U16)', () => { + it('writes the tag as a single byte at offset 0', () => { + const buf = encodeMessage(MessageTag.Ready, null); + expect(buf.readUInt8(0)).toBe(MessageTag.Ready); + }); + + it('writes the payload length as little-endian uint32 at offset 1', () => { + const buf = encodeMessage(MessageTag.Progress, { filesProcessed: 42 }); + // The JSON encoding of {"filesProcessed":42} is a fixed-width + // ASCII string; the encoder uses Buffer.byteLength on the UTF-8 + // bytes, so reading it back via the same path is the load-bearing + // assertion. Header length must equal payloadEnd - header. + const declared = buf.readUInt32LE(1); + expect(declared).toBe(buf.length - PROTOCOL_HEADER_BYTES); + }); + + it('encodes an empty/null payload as 5-byte header + 4-byte "null" body', () => { + // JSON.stringify(null) === "null" (4 bytes). The header carries 4 + // as the length so decoders can advance the cursor consistently. + const buf = encodeMessage(MessageTag.SubBatchDone, null); + expect(buf.length).toBe(PROTOCOL_HEADER_BYTES + 4); + expect(buf.readUInt32LE(1)).toBe(4); + expect(buf.subarray(PROTOCOL_HEADER_BYTES).toString('utf8')).toBe('null'); + }); +}); + +describe('worker IPC protocol — round-trip per tag (U16)', () => { + const cases: ReadonlyArray<{ tag: MessageTagValue; name: string; payload: unknown }> = [ + { + tag: MessageTag.DispatchJob, + name: 'DispatchJob', + payload: { files: [{ path: 'a.ts', content: 'x' }] }, + }, + { tag: MessageTag.Result, name: 'Result', payload: { fileCount: 3, nodes: [{ id: 'n1' }] } }, + { tag: MessageTag.Progress, name: 'Progress', payload: { filesProcessed: 7 } }, + { tag: MessageTag.StartingFile, name: 'StartingFile', payload: { path: 'src/foo.ts' } }, + { tag: MessageTag.SubBatchDone, name: 'SubBatchDone', payload: null }, + { tag: MessageTag.Warning, name: 'Warning', payload: { message: 'unparseable file' } }, + { tag: MessageTag.Error, name: 'Error', payload: { error: 'native crash' } }, + { tag: MessageTag.Ready, name: 'Ready', payload: null }, + ]; + + for (const { tag, name, payload } of cases) { + it(`round-trips ${name} payload unchanged`, () => { + const buf = encodeMessage(tag, payload); + const decoded = decodeMessage(buf); + expect(decoded.tag).toBe(tag); + // toEqual covers structural equality across nested objects/arrays — + // payloads that mutated on the way through would surface as a + // diff in the assertion, not a silent corruption. + expect(decoded.payload).toEqual(payload ?? null); + }); + } + + it('round-trips a non-ASCII path string (UTF-8 boundary)', () => { + const buf = encodeMessage(MessageTag.StartingFile, { path: 'src/café.ts' }); + const decoded = decodeMessage(buf); + expect((decoded.payload as { path: string }).path).toBe('src/café.ts'); + }); + + it('round-trips a payload near the structured-clone sub-batch budget (8 MB)', () => { + // The pool's existing sub-batch byte budget is 8 MB. Verify the + // protocol does not impose a tighter limit by encoding/decoding a + // ~9 MB JSON payload successfully. + const big = 'x'.repeat(9 * 1024 * 1024); + const buf = encodeMessage(MessageTag.Warning, { message: big }); + const decoded = decodeMessage(buf); + expect((decoded.payload as { message: string }).message.length).toBe(big.length); + }); +}); + +describe('worker IPC protocol — decode error paths (U16)', () => { + it('throws ProtocolDecodeError when the buffer is smaller than the 5-byte header', () => { + expect(() => decodeMessage(Buffer.from([0x01, 0x00]))).toThrow(ProtocolDecodeError); + }); + + it('throws ProtocolDecodeError on a tag byte outside the valid range', () => { + const buf = Buffer.alloc(PROTOCOL_HEADER_BYTES + 4); + buf.writeUInt8(0xff, 0); // not a defined tag + buf.writeUInt32LE(4, 1); + Buffer.from('null', 'utf8').copy(buf, PROTOCOL_HEADER_BYTES); + expect(() => decodeMessage(buf)).toThrow(ProtocolDecodeError); + }); + + it('throws ProtocolDecodeError when the declared payload length exceeds the buffer', () => { + // Header claims a 1000-byte body but the buffer only has 4 actual + // payload bytes — a truncated frame. + const buf = Buffer.alloc(PROTOCOL_HEADER_BYTES + 4); + buf.writeUInt8(MessageTag.Ready, 0); + buf.writeUInt32LE(1000, 1); + expect(() => decodeMessage(buf)).toThrow(ProtocolDecodeError); + }); + + it('throws ProtocolDecodeError when payload bytes are not valid JSON', () => { + const garbage = Buffer.from('{not-json}', 'utf8'); + const buf = Buffer.alloc(PROTOCOL_HEADER_BYTES + garbage.length); + buf.writeUInt8(MessageTag.Warning, 0); + buf.writeUInt32LE(garbage.length, 1); + garbage.copy(buf, PROTOCOL_HEADER_BYTES); + expect(() => decodeMessage(buf)).toThrow(ProtocolDecodeError); + }); + + it('preserves the error class name so callers can route protocol violations distinctly', () => { + try { + decodeMessage(Buffer.alloc(2)); + throw new Error('should have thrown'); + } catch (err) { + expect(err).toBeInstanceOf(ProtocolDecodeError); + expect((err as Error).name).toBe('ProtocolDecodeError'); + } + }); +});