Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion packages/sdk/src/named-gate-lowering.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
22 changes: 18 additions & 4 deletions packages/sdk/tests/authored-parallel-agents.test.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand All @@ -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');
Expand All @@ -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 {
Expand Down Expand Up @@ -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');
});
Expand Down
33 changes: 33 additions & 0 deletions packages/sdk/tests/named-gate-diagnostics.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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');
});
});
});
Loading