diff --git a/.changeset/eager-newts-purr.md b/.changeset/eager-newts-purr.md new file mode 100644 index 000000000..0021a5a35 --- /dev/null +++ b/.changeset/eager-newts-purr.md @@ -0,0 +1,5 @@ +--- +type: Fixed +pr: 1009 +--- +**`gsd-tools` no longer throws `EAGAIN` or truncates output under heavy load** — the CLI's stdout/stderr writes now retry the transient `EAGAIN`/`EINTR` errnos and handle short writes when the output stream is a full non-blocking pipe (e.g. the parallel test runner), instead of throwing or silently dropping bytes. diff --git a/src/io.cts b/src/io.cts index 97c7a43a0..33e1890d9 100644 --- a/src/io.cts +++ b/src/io.cts @@ -69,6 +69,63 @@ function reapStaleTempFiles(prefix = 'gsd-', { maxAgeMs = 5 * 60 * 1000, dirsOnl // ─── Output helpers ─────────────────────────────────────────────────────────── +/** + * Transient write errnos. When stdout/stderr is a NON-BLOCKING pipe — as it is + * under the parallel `node --test` runner on Linux CI — a full pipe buffer makes + * `fs.writeSync` throw EAGAIN, and a signal can interrupt it with EINTR. Both + * clear on retry once the reader drains. This is the same transient class the + * STATE.md lock path already retries (ACQUIRE_LOCK_RETRY_ERRNOS, #3776); #1008. + */ +const WRITE_RETRY_ERRNOS = new Set(['EAGAIN', 'EINTR']); + +// Bounded so a pathological never-draining fd cannot spin forever. Each retry +// yields the thread for ~1ms via Atomics.wait (the project's sync-sleep idiom — +// see clock.cts realClock.sleep), so the cap is ~1s of total back-pressure wait. +const WRITE_MAX_RETRIES = 1000; +const WRITE_RETRY_BACKOFF_MS = 1; + +// Sleep buffer is lazily allocated on the FIRST back-pressure retry (rare — only +// when a non-blocking pipe is full) and then reused. Keeping it out of module +// load costs nothing on the overwhelmingly common no-retry path and avoids +// perturbing SharedArrayBuffer-allocation accounting in other modules (perf-316). +let _writeSleepBuf: Int32Array | null = null; +function backoffOnce(): void { + if (_writeSleepBuf === null) _writeSleepBuf = new Int32Array(new SharedArrayBuffer(4)); + Atomics.wait(_writeSleepBuf, 0, 0, WRITE_RETRY_BACKOFF_MS); +} + +/** + * Write the entire payload to `fd`, tolerating non-blocking-pipe back-pressure. + * + * `fs.writeSync` does NOT block on a non-blocking pipe: a full buffer throws + * EAGAIN, and a partially-drained buffer returns a SHORT count (fewer bytes than + * requested). The previous bare `fs.writeSync(fd, string)` call assumed it always + * blocked until the kernel accepted every byte — false under load, which both + * threw spurious errors and risked silently truncating output (#1008). + * + * This loops on short counts (advancing the offset) and retries EAGAIN/EINTR with + * a brief Atomics.wait backoff that yields the thread so the reader can drain. + * Non-transient errors (e.g. EPIPE) propagate unchanged. + */ +function writeAllSync(fd: number, data: string): void { + const buf = Buffer.from(data, 'utf8'); + let offset = 0; + let retries = 0; + while (offset < buf.length) { + try { + offset += fs.writeSync(fd, buf, offset, buf.length - offset); + } catch (err) { + const code = (err as NodeJS.ErrnoException).code ?? ''; + if (WRITE_RETRY_ERRNOS.has(code) && retries < WRITE_MAX_RETRIES) { + retries += 1; + backoffOnce(); + continue; + } + throw err; + } + } +} + function output(result: unknown, raw: boolean, rawValue?: unknown): void { let data: string; if (raw && rawValue !== undefined) { @@ -89,10 +146,10 @@ function output(result: unknown, raw: boolean, rawValue?: unknown): void { } } // process.stdout.write() is async when stdout is a pipe — process.exit() - // can tear down the process before the reader consumes the buffer. - // fs.writeSync(1, ...) blocks until the kernel accepts the bytes, and - // skipping process.exit() lets the event loop drain naturally. - fs.writeSync(1, data); + // can tear down the process before the reader consumes the buffer. writeAllSync + // pushes every byte synchronously (looping short counts, retrying EAGAIN/EINTR), + // and skipping process.exit() lets the event loop drain naturally. + writeAllSync(1, data); } /** @@ -155,9 +212,9 @@ function getJsonErrorMode(): boolean { return _jsonErrorMode; } function error(message: string, reason: ErrorReasonValue = ERROR_REASON.UNKNOWN): never { if (_jsonErrorMode) { const payload = JSON.stringify({ ok: false, reason, message }) + '\n'; - fs.writeSync(2, payload); + writeAllSync(2, payload); } else { - fs.writeSync(2, 'Error: ' + message + '\n'); + writeAllSync(2, 'Error: ' + message + '\n'); } process.exit(1); } diff --git a/tests/io.test.cjs b/tests/io.test.cjs index 3825be1c1..8663ac825 100644 --- a/tests/io.test.cjs +++ b/tests/io.test.cjs @@ -339,3 +339,113 @@ describe('core.cjs re-export shims', () => { assert.strictEqual(core.GSD_TEMP_DIR, io.GSD_TEMP_DIR); }); }); + +// ─── bug #1008: output()/error() tolerate a full / slow non-blocking pipe ───── +// +// The pre-fix bare `fs.writeSync(fd, data)` assumed it blocks until the kernel +// accepts every byte — false when fd is a non-blocking pipe (the parallel +// node:test runner on Linux): a full pipe throws EAGAIN and a partially-drained +// pipe returns a SHORT count. These behavioral tests inject fs.writeSync via +// mock.method (the approved fault-injection seam) and assert the observable +// contract (no throw, full payload, real errors still surface). They are red +// against the pre-fix io.cjs (throw / truncate). + +// Normalize either writeSync call form to the chunk it emits: +// buffer form: writeSync(fd, buffer, offset, length) ← the fixed writeAllSync loop +// string form: writeSync(fd, string) ← the pre-fix bare call +function bug1008ChunkOf(data, offset, length) { + if (Buffer.isBuffer(data)) { + const start = offset ?? 0; + const end = length === undefined ? data.length : start + length; + return data.subarray(start, end).toString('utf8'); + } + return String(data); +} + +function bug1008WriteError(code, errno) { + const e = new Error(`${code}: write`); + e.code = code; + e.errno = errno; + e.syscall = 'write'; + return e; +} + +describe('bug #1008: io.output() tolerates a full / slow non-blocking pipe', () => { + test('retries on EAGAIN and emits the full payload without throwing', (t) => { + const written = []; + let calls = 0; + t.mock.method(fs, 'writeSync', (fd, data, offset, length) => { + calls += 1; + if (calls === 1) throw bug1008WriteError('EAGAIN', -11); // pipe momentarily full + const chunk = bug1008ChunkOf(data, offset, length); + written.push(chunk); + return Buffer.byteLength(chunk, 'utf8'); + }); + + const payload = { ok: true, n: 42 }; + assert.doesNotThrow(() => io.output(payload, false)); + assert.ok(calls >= 2, `expected a retry after EAGAIN, got ${calls} call(s)`); + assert.equal(written.join(''), JSON.stringify(payload, null, 2), 'full payload must reach the fd'); + }); + + test('retries on EINTR (signal-interrupted write) too', (t) => { + const written = []; + let calls = 0; + t.mock.method(fs, 'writeSync', (fd, data, offset, length) => { + calls += 1; + if (calls === 1) throw bug1008WriteError('EINTR', -4); + const chunk = bug1008ChunkOf(data, offset, length); + written.push(chunk); + return Buffer.byteLength(chunk, 'utf8'); + }); + + assert.doesNotThrow(() => io.output('plain', true, 'PLAIN-RAW')); + assert.equal(written.join(''), 'PLAIN-RAW'); + }); + + test('handles short (partial) writes without truncating', (t) => { + const written = []; + const CAP = 3; // each writeSync accepts at most 3 bytes, like a draining pipe + t.mock.method(fs, 'writeSync', (fd, data, offset, length) => { + const chunk = bug1008ChunkOf(data, offset, length); + const part = chunk.slice(0, CAP); + written.push(part); + return Buffer.byteLength(part, 'utf8'); + }); + + const payload = { message: 'a reasonably long ascii payload to force many short writes' }; + io.output(payload, false); + assert.equal(written.join(''), JSON.stringify(payload, null, 2), 'no bytes may be dropped on short writes'); + }); + + test('does NOT swallow a genuine, non-transient write error (EPIPE)', (t) => { + t.mock.method(fs, 'writeSync', () => { throw bug1008WriteError('EPIPE', -32); }); + assert.throws( + () => io.output({ ok: true }, false), + (err) => err.code === 'EPIPE', + 'real (non-transient) errors must still surface', + ); + }); +}); + +describe('bug #1008: io.error() tolerates a full non-blocking stderr pipe', () => { + test('retries on EAGAIN, emits the full message, and still exits', (t) => { + const written = []; + let calls = 0; + let exitCode = null; + t.mock.method(process, 'exit', (code) => { exitCode = code; }); // neutralize the hard exit + t.mock.method(fs, 'writeSync', (fd, data, offset, length) => { + calls += 1; + if (calls === 1) throw bug1008WriteError('EAGAIN', -11); + assert.equal(fd, 2, 'error() must write to stderr'); + const chunk = bug1008ChunkOf(data, offset, length); + written.push(chunk); + return Buffer.byteLength(chunk, 'utf8'); + }); + + assert.doesNotThrow(() => io.error('boom', io.ERROR_REASON.UNKNOWN)); + assert.ok(calls >= 2, 'error() should retry after EAGAIN'); + assert.equal(written.join(''), 'Error: boom\n'); + assert.equal(exitCode, 1, 'error() must still exit(1) after a retried write'); + }); +});