feat(workers): introduce protocol.ts wire-format module (U16, IPC scaffold)

Defines the binary frame for worker-thread IPC as an isolated, fully-tested
module. Production wiring is deferred to U17 — shipping the wire-format
contract first de-risks the migration by establishing a single source of
truth for the byte layout. Resolves the scaffold half of PR #1693 review
R12.

Wire layout (per message, single buffer):

  +---------+-----------+---------------------+
  | tag     | length    | payload bytes …     |
  | 1 byte  | 4 bytes   |                     |
  +---------+-----------+---------------------+

  tag    : MessageTag enum value (0x01 DispatchJob ... 0x08 Ready)
  length : little-endian uint32 byte count for the payload region
  payload: UTF-8 JSON-encoded value, possibly "null"

Why JSON for the body (rather than 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
the same encodeMessage / decodeMessage surface without changing the
frame.

API:
  - MessageTag (const object): stable byte tags 0x01..0x08
  - PROTOCOL_HEADER_BYTES = 5
  - ProtocolDecodeError extends Error: distinct class so U17's
    pool-side handler can route protocol violations through the
    existing messageerror recovery layer (U3 H1) distinctly from
    other failure classes
  - encodeMessage(tag, payload): Buffer
  - decodeMessage(buf): { tag, payload }
  - Uses Buffer#subarray instead of the deprecated Buffer#slice

Tests (18, all exact-equality per DoD §2.7):
  - byte layout (tag at offset 0, length LE uint32 at offset 1)
  - empty/null payload encodes to 5-byte header + 4-byte "null" body
  - round-trip for every MessageTag with representative payloads
  - non-ASCII path string (UTF-8 byte-length boundary)
  - 9 MB payload (well past the existing 8 MB sub-batch budget)
  - decode errors surface as ProtocolDecodeError, not generic Error:
      * buffer < header size
      * tag outside valid range
      * declared length exceeds buffer
      * payload bytes are not valid JSON
  - error class name is preserved through prototype chain so callers
    can `err instanceof ProtocolDecodeError` reliably
This commit is contained in:
Gergo Magyar 2026-05-20 10:01:14 +01:00
parent 47060c2408
commit ac90a133d0
2 changed files with 293 additions and 0 deletions

View file

@ -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 };
}

View file

@ -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');
}
});
});