diff --git a/packages/sdk/src/named-gate-lowering.ts b/packages/sdk/src/named-gate-lowering.ts index 8a4ac503..56200f6b 100644 --- a/packages/sdk/src/named-gate-lowering.ts +++ b/packages/sdk/src/named-gate-lowering.ts @@ -125,10 +125,17 @@ process.exit(result.status===0?0:1);`; // The other gate that spawns a child. Its `wc` stderr was piped and then // discarded, so a broken or absent `wc` was reported as a word count out // of bounds. A bounded suffix of what it said travels with the cause. + // + // EPIPE is not "could not run": it means wc ran and exited before + // reading all of `input` (it crashed, was killed, or never reads + // stdin). spawnSync still reports its status, signal and output, so + // those decide the verdict. Whether the write raced wc's exit is + // timing, so it must not change the verdict; reporting it hid what wc + // did (#611). body = `const result=cp.spawnSync('wc',['-w'],{input:text,encoding:'utf8',env:{...process.env,LC_ALL:'C'}}); const noise=(result.stderr||'').trim().slice(-200); const said=noise===''?'':': '+noise; -if(result.error)fail('could not run wc -w: '+(result.error.code||result.error.message)+said); +if(result.error&&result.error.code!=='EPIPE')fail('could not run wc -w: '+(result.error.code||result.error.message)+said); if(result.signal)fail('wc -w was terminated by '+result.signal+said); if(result.status!==0)fail('wc -w exited '+result.status+said); const count=result.stdout.trim(); diff --git a/packages/sdk/tests/authored-parallel-agents.test.ts b/packages/sdk/tests/authored-parallel-agents.test.ts index 12ef531e..c3503368 100644 --- a/packages/sdk/tests/authored-parallel-agents.test.ts +++ b/packages/sdk/tests/authored-parallel-agents.test.ts @@ -1,4 +1,4 @@ -import { mkdirSync, readFileSync, symlinkSync, writeFileSync } from 'node:fs'; +import { existsSync, mkdirSync, readFileSync, symlinkSync, writeFileSync } from 'node:fs'; import { join, resolve } from 'node:path'; import { afterEach, describe, expect, it } from 'vitest'; import { flow } from '@relayflows/surface'; @@ -19,12 +19,14 @@ async function slowAgents(capacity: number, failFirstSession = false) { const fixture = chainFixture(); closes.push(() => fixture.close()); const spans = join(fixture.root, 'spans.jsonl'); + const started = join(fixture.root, 'started'); writeFileSync(fixture.wrapper, `#!/usr/bin/env node import { receiveWrapperRequest } from ${JSON.stringify(resolve('../../testdata/preflight/wrapper-session.mjs'))}; import { appendFileSync, existsSync, writeFileSync } from 'node:fs'; if (process.argv[2] === 'auth') process.exit(0); const request = await receiveWrapperRequest(); if (request) { + writeFileSync(${JSON.stringify(started)}, ''); const start = Date.now(); await new Promise(done => setTimeout(done, 400)); appendFileSync(${JSON.stringify(spans)}, JSON.stringify({ start, end: Date.now() }) + '\\n'); @@ -50,7 +52,14 @@ if (request) { closes.push(() => llm.close()); const readSpans = () => readFileSync(spans, 'utf8').trim().split('\n') .map(line => JSON.parse(line) as { start: number; end: number }); - return { fixture, client, agent, readSpans }; + /** Resolves once a session has received its request: an agent holds a slot and is running. */ + const firstStarted = async () => { + for (const deadline = Date.now() + 5_000; !existsSync(started);) { + if (Date.now() > deadline) throw new Error('no agent session started within 5s'); + await new Promise(done => setTimeout(done, 10)); + } + }; + return { fixture, client, agent, readSpans, firstStarted }; } function peakOverlap(spans: Array<{ start: number; end: number }>): number { @@ -135,11 +144,16 @@ describe('authored steps under local workers with capacity', () => { }); it('never starts queued agents once the body has failed', async () => { - const { fixture, client, agent, readSpans } = await slowAgents(1); + const { fixture, client, agent, readSpans, firstStarted } = await slowAgents(1); const failing = flow('fails-while-queued', async f => { await Promise.all([ ...['a', 'b', 'c'].map(lens => f.agent(`review-${lens}`, { task: `Review for ${lens}` })), - new Promise((_, reject) => setTimeout(() => reject(new Error('body failed')), 100)), + // Fail once one agent is really running, not after a fixed delay. Each + // call runs its preflight before it asks for a slot; under load that + // outlasted 100ms, so the body failed while no agent held the slot, + // all three were (correctly) refused, and spans.jsonl was never + // written (#611). + firstStarted().then(() => { throw new Error('body failed'); }), ]); f.done('success'); }); diff --git a/packages/sdk/tests/named-gate-diagnostics.test.ts b/packages/sdk/tests/named-gate-diagnostics.test.ts index b35c60eb..d5f661d6 100644 --- a/packages/sdk/tests/named-gate-diagnostics.test.ts +++ b/packages/sdk/tests/named-gate-diagnostics.test.ts @@ -279,4 +279,37 @@ describe('word_count_bounds reports why its own child failed', () => { expect(capture.status).toBe(1); expect(capture.stderr).toContain('SIGKILL'); }); + + // #611: a wc that exits without reading stdin made spawnSync's input write + // fail with EPIPE, reported as "could not run wc -w: EPIPE" in place of what + // wc did. With 'one two' that depended on who won the race. spawnSync's + // stdin is a socketpair: on macOS its buffer is small, so this input makes + // wc exit first every time; on Linux the buffer holds more than one env + // string can carry, so the race is not forced there. These hold either way. + describe('when wc exits before reading all of its input', () => { + // Below Linux's 128 KiB limit on one env string. + const longText = 'word '.repeat(20_000); + const run = (script: string) => + runGate(command(), { output: deterministicEnvelope(longText) }, { path: stubWordCount(script) }); + + it('still reports output that is not a count', () => { + const capture = run("printf 'not a number'"); + expect(capture.status).toBe(1); + expect(capture.stderr).toContain('not a number'); + expect(capture.stderr).not.toContain('EPIPE'); + }); + + it('still reports its exit status and what it said', () => { + const capture = run("printf 'wc: read error' >&2; exit 2"); + expect(capture.status).toBe(1); + expect(capture.stderr).toContain('exited 2'); + expect(capture.stderr).toContain('wc: read error'); + }); + + it('still reports the signal that ended it', () => { + const capture = run('kill -9 $$'); + expect(capture.status).toBe(1); + expect(capture.stderr).toContain('SIGKILL'); + }); + }); });