mirror of
https://github.com/garrytan/gstack.git
synced 2026-10-02 17:40:02 +02:00
* feat: add surface-aware exploratory QA and ship documentation gates * test: preserve delegated QA setup authority after main integration * fix(qa): clarify exploration order and preserve report artifacts * test(qa): follow the shared setup reference directly * refactor(ship): make verification and recovery routes explicit * test(ship): align evidence and review guards with explicit routes * fix(workflows): clarify ship recovery and functional QA evidence * fix(workflows): clarify approval recovery and full QA coverage * refactor(workflows): order review transactions and clarify ship state * fix(ship): clarify final verification and fail closed at publication * fix(evals): attribute native atomic documentation writes * fix(ship): clarify recovery and documentation lifecycle guidance * fix(test): preserve observed native placeholder styling in CI * fix(codex): report watchdog timeouts without a process-exit race * Checkpoint functional QA implementation and workflow validation repairs * Fix documentation and shared-review fixture contracts * docs: clarify judge reuse and evaluation supervision * test: align review evidence and selected case contracts * test: verify append-only documentation checkpoints and recovery * fix: qualify QA workflows and CI validation repairs * fix: launch shared-libs fixture scripts on Windows * fix: qualify QA deadlines, fixture isolation, and shard cleanup * fix: preserve qualified QA and cancellation repairs * fix: enforce functional fixture authority and share strict event decoding * fix: retain free-test evidence and explain recovery * fix: reject malformed native evidence after decoder consolidation * test: use reliable capture for telemetry privacy filters * test: refresh measured quick coverage and document validation costs * Fix native fixture receipts and preserve VM validation evidence * Align negative judge controls with upstream clarity policy * Fix report-only QA preparation and public evidence handling * Clarify QA-only preparation and current-report preservation * Stream Ship quality judgments with an explicit 64k response contract * Validate compact judge reasoning locally with supported wire schema * Align functional QA fixture instructions with evidence acceptance * Bind native browser diagnostics to execution evidence and align review verdicts * Preserve native diagnostic line boundaries * Serialize functional QA evidence from native captures * Keep large QA evidence fixture payload out of Windows argv
717 lines
31 KiB
TypeScript
717 lines
31 KiB
TypeScript
/**
|
|
* Claude CLI subprocess runner for skill E2E testing.
|
|
*
|
|
* Spawns `claude -p` as a completely independent process (not via Agent SDK),
|
|
* so it works inside Claude Code sessions. Pipes prompt via stdin, streams
|
|
* NDJSON output for real-time progress, scans for browse errors.
|
|
*/
|
|
|
|
import * as fs from 'fs';
|
|
import * as path from 'path';
|
|
import * as os from 'os';
|
|
import { spawn } from 'child_process';
|
|
import { Readable } from 'node:stream';
|
|
import { createHash, type Hash } from 'node:crypto';
|
|
import { getProjectEvalDir } from './eval-store';
|
|
import { hermeticChildEnv, isHermeticEnabled } from './hermetic-env';
|
|
import { killProcessGroup } from '../../scripts/test-strict-output';
|
|
import { resolveEvalModel } from '../../lib/eval-model';
|
|
|
|
const GSTACK_DEV_DIR = path.join(os.homedir(), '.gstack-dev');
|
|
const HEARTBEAT_PATH = path.join(GSTACK_DEV_DIR, 'e2e-live.json'); // heartbeat stays global
|
|
const PROJECT_DIR = path.dirname(getProjectEvalDir()); // ~/.gstack/projects/$SLUG/
|
|
|
|
/** Sanitize test name for use as filename: strip leading slashes, replace / with - */
|
|
export function sanitizeTestName(name: string): string {
|
|
return name.replace(/^\/+/, '').replace(/\//g, '-');
|
|
}
|
|
|
|
/** Atomic write: write to .tmp then rename. Non-fatal on error. */
|
|
function atomicWriteSync(filePath: string, data: string): void {
|
|
const tmp = filePath + '.tmp';
|
|
fs.writeFileSync(tmp, data);
|
|
fs.renameSync(tmp, filePath);
|
|
}
|
|
|
|
export interface CostEstimate {
|
|
inputChars: number;
|
|
outputChars: number;
|
|
estimatedTokens: number;
|
|
estimatedCost: number; // USD
|
|
turnsUsed: number;
|
|
}
|
|
|
|
export interface SkillTestResult {
|
|
toolCalls: Array<{ tool: string; input: any; output: string }>;
|
|
browseErrors: string[];
|
|
exitReason: string;
|
|
duration: number;
|
|
output: string;
|
|
costEstimate: CostEstimate;
|
|
transcript: any[];
|
|
/** Which model was used for this test (added for Sonnet/Opus split diagnostics) */
|
|
model: string;
|
|
/** Time from spawn to first NDJSON line, in ms (added for rate-limit diagnostics) */
|
|
firstResponseMs: number;
|
|
/** Peak latency between consecutive tool calls, in ms */
|
|
maxInterTurnMs: number;
|
|
}
|
|
|
|
/** Local default startup grace: 90s covers observed API queue latency
|
|
* (60-90s receipts) without letting a dead API burn a 600s budget. */
|
|
export const STARTUP_GRACE_MS = 90_000;
|
|
/** CI floor (TODOS-filed): shared runners queue harder; killing startup
|
|
* before 300s in CI converts ordinary queueing into false failures.
|
|
* Pinned by test/session-runner-startup-grace.test.ts. */
|
|
export const STARTUP_GRACE_CI_FLOOR_MS = 300_000;
|
|
/** Existing pipe-drain allowance; never adds model work time. */
|
|
export { SESSION_DRAIN_GRACE_MS } from './session-drain-policy';
|
|
import { SESSION_DRAIN_GRACE_MS } from './session-drain-policy';
|
|
|
|
const BROWSE_ERROR_PATTERNS = [
|
|
/Unknown command: \w+/,
|
|
/Unknown snapshot flag: .+/,
|
|
/ERROR: browse binary not found/,
|
|
/Server failed to start/,
|
|
/no such file or directory.*\bbrowse(?:\.exe)?(?=$|[\s'":),])/i,
|
|
];
|
|
|
|
// --- Testable NDJSON parser ---
|
|
|
|
export interface ParsedNDJSON {
|
|
transcript: any[];
|
|
resultLine: any | null;
|
|
turnCount: number;
|
|
toolCallCount: number;
|
|
toolCalls: Array<{ tool: string; input: any; output: string }>;
|
|
}
|
|
|
|
function toolResultText(content: unknown): string | null {
|
|
if (typeof content === 'string') return content;
|
|
if (!Array.isArray(content)) return null;
|
|
return content.flatMap(block => block?.type === 'text' && typeof block.text === 'string'
|
|
? [block.text] : []).join('\n');
|
|
}
|
|
|
|
/**
|
|
* Parse an array of NDJSON lines into structured transcript data.
|
|
* Pure function — no I/O, no side effects. Used by both the streaming
|
|
* reader and unit tests.
|
|
*/
|
|
export function parseNDJSON(lines: string[]): ParsedNDJSON {
|
|
const transcript: any[] = [];
|
|
let resultLine: any = null;
|
|
let turnCount = 0;
|
|
let toolCallCount = 0;
|
|
const toolCalls: ParsedNDJSON['toolCalls'] = [];
|
|
const callsByParent = new Map<string | null, Map<string, ParsedNDJSON['toolCalls'][number]>>();
|
|
|
|
for (const line of lines) {
|
|
if (!line.trim()) continue;
|
|
try {
|
|
const event = JSON.parse(line);
|
|
transcript.push(event);
|
|
|
|
// Track turns and tool calls from assistant events
|
|
if (event.type === 'assistant') {
|
|
turnCount++;
|
|
const content = event.message?.content || [];
|
|
for (const item of content) {
|
|
if (item.type === 'tool_use') {
|
|
toolCallCount++;
|
|
const call = {
|
|
tool: item.name || 'unknown',
|
|
input: item.input || {},
|
|
output: '',
|
|
};
|
|
toolCalls.push(call);
|
|
if (typeof item.id === 'string') {
|
|
// Forwarded subagent events may reuse a parent's tool-use ID.
|
|
const parent = event.parent_tool_use_id ?? null;
|
|
let calls = callsByParent.get(parent);
|
|
if (!calls) callsByParent.set(parent, calls = new Map());
|
|
calls.set(item.id, call);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
if (event.type === 'user' && Array.isArray(event.message?.content)) {
|
|
const results = event.message.content.filter((item: any) => item?.type === 'tool_result');
|
|
const calls = callsByParent.get(event.parent_tool_use_id ?? null);
|
|
for (const result of results) {
|
|
const call = calls?.get(result.tool_use_id);
|
|
if (!call) continue;
|
|
// A sole Agent/Task result also carries the clean verdict separately
|
|
// from the message's agentId/usage wrapper. Keep only public text.
|
|
const verdict = results.length === 1 && ['Agent', 'Task'].includes(call.tool)
|
|
? toolResultText(event.tool_use_result?.content) : null;
|
|
call.output = verdict ?? toolResultText(result.content) ?? '';
|
|
}
|
|
}
|
|
|
|
if (event.type === 'result') resultLine = event;
|
|
} catch { /* skip malformed lines */ }
|
|
}
|
|
|
|
return { transcript, resultLine, turnCount, toolCallCount, toolCalls };
|
|
}
|
|
|
|
function truncate(s: string, max: number): string {
|
|
return s.length > max ? s.slice(0, max) + '…' : s;
|
|
}
|
|
|
|
/** Diagnostic-only projection. Partial input never becomes a complete tool call. */
|
|
function publicStreamProjection(startTime: number): (line: string) => string {
|
|
let messageId: string | undefined;
|
|
const blocks = new Map<number, { type: string; tool?: string; bytes: number; hash: Hash }>();
|
|
return (line) => {
|
|
let row: any;
|
|
try { row = JSON.parse(line); } catch {
|
|
// A truncated line may contain private reasoning or unfinished tool input.
|
|
return JSON.stringify({ type: 'public_stream_diagnostic', kind: 'unparseable_line',
|
|
elapsedMs: Date.now() - startTime, bytes: Buffer.byteLength(line) });
|
|
}
|
|
if (row === null || typeof row !== 'object' || Array.isArray(row)) {
|
|
return JSON.stringify({ type: 'public_stream_diagnostic', kind: 'non_object_line',
|
|
elapsedMs: Date.now() - startTime, bytes: Buffer.byteLength(line) });
|
|
}
|
|
if (row.type === 'stream_event') {
|
|
const event = row.event ?? {};
|
|
if (event.type === 'message_start') {
|
|
messageId = event.message?.id;
|
|
blocks.clear();
|
|
}
|
|
const index = event.index;
|
|
if (event.type === 'content_block_start' && Number.isInteger(index)) {
|
|
blocks.set(index, { type: event.content_block?.type, tool: event.content_block?.name,
|
|
bytes: 0, hash: createHash('sha256') });
|
|
}
|
|
const block = blocks.get(index);
|
|
if (event.type === 'content_block_delta' && block?.type === 'tool_use'
|
|
&& event.delta?.type === 'input_json_delta' && typeof event.delta.partial_json === 'string') {
|
|
const chunk = Buffer.from(event.delta.partial_json);
|
|
block.bytes += chunk.length;
|
|
block.hash.update(chunk);
|
|
}
|
|
const diagnostic = { type: 'public_stream_diagnostic', kind: event.type,
|
|
session_id: row.session_id, messageId, index, elapsedMs: Date.now() - startTime,
|
|
blockType: block?.type, toolName: block?.tool, deltaType: event.delta?.type,
|
|
...(block?.type === 'tool_use' ? { inputBytes: block.bytes,
|
|
inputSha256: block.hash.copy().digest('hex') } : {}),
|
|
...(event.type === 'message_delta' ? { stopReason: event.delta?.stop_reason } : {}),
|
|
};
|
|
if (event.type === 'content_block_stop') blocks.delete(index);
|
|
return JSON.stringify(diagnostic);
|
|
}
|
|
if (Array.isArray(row.message?.content)) {
|
|
row.message.content = row.message.content.map((block: any) =>
|
|
block.type === 'thinking' || block.type === 'redacted_thinking'
|
|
? { type: block.type, omitted: true } : block);
|
|
}
|
|
return JSON.stringify(row);
|
|
};
|
|
}
|
|
|
|
// --- Main runner ---
|
|
|
|
export async function runSkillTest(options: {
|
|
prompt: string;
|
|
workingDirectory: string;
|
|
maxTurns?: number;
|
|
/** Optional harness contract appended to the native system prompt. */
|
|
appendSystemPrompt?: string;
|
|
/** Opt-in section completion reserve; describes the existing entry deadline. */
|
|
completionReserveMs?: number;
|
|
/** Approval allowlist; does not restrict which tools the model can see. */
|
|
allowedTools?: string[];
|
|
/** Optional built-in tool availability. Omit to preserve the CLI defaults. */
|
|
tools?: string[];
|
|
/** Opt-in public block timing/input-size diagnostics; never completion evidence. */
|
|
publicStreamDiagnostics?: boolean;
|
|
timeout?: number;
|
|
testName?: string;
|
|
runId?: string;
|
|
/** Model to use. Defaults to the frontier eval model (overridable via EVALS_MODEL env). */
|
|
model?: string;
|
|
/** Extra env vars merged into the spawned claude -p process. Useful for
|
|
* per-test GSTACK_HOME overrides so the test doesn't have to spell out
|
|
* env setup in the prompt itself. */
|
|
env?: Record<string, string>;
|
|
/** Startup-phase deadline: if NO NDJSON byte arrives within this window,
|
|
* the run is killed EARLY with exitReason 'timeout_startup' instead of
|
|
* burning the whole work budget waiting on an API that is not answering
|
|
* (the recurring '0 turns / $0.00' class — four budget-bump receipts).
|
|
* Defaults to min(STARTUP_GRACE_MS, timeout); the CI floor is higher
|
|
* because CI queueing is real. Startup and model work share timeout;
|
|
* only pipe cleanup may use the separate SESSION_DRAIN_GRACE_MS. */
|
|
startupGraceMs?: number;
|
|
/** Cancel the owned process group when an enclosing attempt expires. */
|
|
signal?: AbortSignal;
|
|
nativeLifecycle?: {
|
|
onSpawn(pid: number): void;
|
|
onSettled(input: { deadline: number; exited: boolean }): Promise<void>;
|
|
};
|
|
}): Promise<SkillTestResult> {
|
|
const startTime = Date.now();
|
|
options.signal?.throwIfAborted();
|
|
const {
|
|
prompt,
|
|
workingDirectory,
|
|
maxTurns = 15,
|
|
allowedTools = ['Bash', 'Read', 'Write'],
|
|
timeout = 120_000,
|
|
testName,
|
|
runId,
|
|
env: extraEnv,
|
|
signal,
|
|
} = options;
|
|
// The CI floor is a FLOOR, not a default: an explicit startupGraceMs below
|
|
// 300s in CI would re-open the queueing-becomes-false-red hole the floor
|
|
// exists for (review finding — the name promised a clamp the code lacked).
|
|
// Local runs honor the caller verbatim; timeout still caps everything.
|
|
const requestedGrace = options.startupGraceMs ?? (process.env.CI ? STARTUP_GRACE_CI_FLOOR_MS : STARTUP_GRACE_MS);
|
|
const startupGraceMs = Math.min(
|
|
process.env.CI ? Math.max(requestedGrace, STARTUP_GRACE_CI_FLOOR_MS) : requestedGrace,
|
|
timeout,
|
|
);
|
|
const model = options.model ?? process.env.EVALS_MODEL ?? resolveEvalModel('capture');
|
|
|
|
const deadline = startTime + timeout;
|
|
const startedAt = new Date().toISOString();
|
|
let systemPrompt = options.appendSystemPrompt;
|
|
if (options.completionReserveMs !== undefined) {
|
|
const reserve = options.completionReserveMs;
|
|
if (!Number.isFinite(timeout) || timeout <= 0 || !Number.isFinite(reserve) || reserve <= 0 || reserve >= timeout) {
|
|
throw new Error('Section completion reserve must be positive and smaller than the existing work timeout');
|
|
}
|
|
if (!allowedTools.includes('Bash') || (options.tools !== undefined && !options.tools.includes('Bash'))) {
|
|
throw new Error('Section completion clock requires Bash in the declared tools and approval allowlist');
|
|
}
|
|
const notice = `Section completion clock (fixture contract):
|
|
Runner entry UTC: ${new Date(startTime).toISOString()}
|
|
Hard deadline UTC: ${new Date(deadline).toISOString()}
|
|
Completion reserve starts UTC: ${new Date(deadline - reserve).toISOString()}
|
|
Setup, CLI startup and API queueing consume this same window; it never resets.
|
|
Before source Reads and after each saved checkpoint, use Bash to run exactly \`date -u +%Y-%m-%dT%H:%M:%SZ\`. Compare that observed UTC time with the times above. When remaining time is at most ${reserve / 1000} seconds, prioritize the remaining required completion outputs and verification. No required content or gate may be skipped. If the clock read fails, report timing unavailable; do not invent remaining time or restart the deadline.`;
|
|
systemPrompt = systemPrompt ? `${systemPrompt}\n\n${notice}` : notice;
|
|
}
|
|
|
|
// Set up per-run log directory if runId is provided
|
|
let runDir: string | null = null;
|
|
const safeName = testName ? sanitizeTestName(testName) : null;
|
|
if (runId) {
|
|
try {
|
|
runDir = path.join(PROJECT_DIR, 'e2e-runs', runId);
|
|
fs.mkdirSync(runDir, { recursive: true });
|
|
} catch { /* non-fatal */ }
|
|
}
|
|
|
|
// Spawn claude -p with streaming NDJSON output. Prompt piped via stdin to
|
|
// avoid shell escaping issues. --verbose is required for stream-json mode.
|
|
const args = [
|
|
'-p',
|
|
'--model', model,
|
|
'--output-format', 'stream-json',
|
|
'--verbose',
|
|
'--dangerously-skip-permissions',
|
|
'--max-turns', String(maxTurns),
|
|
'--allowed-tools', ...allowedTools,
|
|
];
|
|
if (systemPrompt) args.push('--append-system-prompt', systemPrompt);
|
|
// --allowed-tools controls approval, including when permissions are skipped;
|
|
// only --tools removes unrelated built-ins such as Agent, Bash, and Skill.
|
|
// Keep this opt-in: existing workflow evals intentionally use CLI defaults.
|
|
if (options.tools !== undefined) args.push('--tools', options.tools.join(','));
|
|
if (options.publicStreamDiagnostics) args.push('--include-partial-messages');
|
|
// Hermetic children get zero MCP servers (no --mcp-config is passed).
|
|
// Gated on the same call-time check as the env scrub so EVALS_HERMETIC=0
|
|
// restores operator MCP along with the operator env.
|
|
if (isHermeticEnabled()) args.push('--strict-mcp-config');
|
|
|
|
// Spawn claude directly with array-form args (no shell interpolation).
|
|
// node:child_process spawn (not Bun.spawn): `detached` puts the child in
|
|
// its OWN process group, so the timeout handler can killpg the whole tree.
|
|
// Bun.spawn has no detached option, and its bare proc.kill() signalled only
|
|
// claude itself — tool subprocesses claude spawned survived as orphans
|
|
// burning shared API rate for the rest of the shard's lifetime.
|
|
// Prompt is piped via stdin to avoid temp files and shell escaping.
|
|
const childEnv = hermeticChildEnv({ GSTACK_HEADLESS: '1', ...extraEnv });
|
|
signal?.throwIfAborted();
|
|
// Entry, logging, and hermetic setup consume the same startup/work budget.
|
|
if (Date.now() >= Math.min(deadline, startTime + startupGraceMs)) {
|
|
return {
|
|
toolCalls: [], browseErrors: [], exitReason: 'timeout_startup', duration: Date.now() - startTime,
|
|
output: '', transcript: [], model, firstResponseMs: 0, maxInterTurnMs: 0,
|
|
costEstimate: { inputChars: prompt.length, outputChars: 0, estimatedTokens: 0, estimatedCost: 0, turnsUsed: 0 },
|
|
};
|
|
}
|
|
const proc = spawn('claude', args, {
|
|
cwd: workingDirectory,
|
|
// Hermetic by default (see test/helpers/hermetic-env.ts): operator
|
|
// session context (CONDUCTOR_*, CLAUDECODE, ~/.claude config, ~/.gstack)
|
|
// never reaches the child; EVALS_HERMETIC=0 restores the legacy env.
|
|
// Default GSTACK_HEADLESS=1 so eval/E2E runs classify as headless (BLOCK on an
|
|
// AskUserQuestion failure rather than emit a prose question no human reads). A
|
|
// suite exercising the INTERACTIVE prose-fallback path opts out by passing
|
|
// `env: { GSTACK_HEADLESS: '' }` — extraEnv wins because it spreads last.
|
|
env: childEnv,
|
|
stdio: ['pipe', 'pipe', 'pipe'],
|
|
detached: process.platform !== 'win32',
|
|
});
|
|
const stdoutWeb = Readable.toWeb(proc.stdout!) as ReadableStream<Uint8Array>;
|
|
const reader = stdoutWeb.getReader();
|
|
// Exit is independent of descendant-held pipes. Both stream drains and a
|
|
// failed kill that never emits exit have a bounded completion path.
|
|
let exitCode: number | undefined;
|
|
let stdoutDone = false;
|
|
let stderrDone = false;
|
|
let stderrEnded = false;
|
|
let drainExpired = false;
|
|
let streamError: Error | undefined;
|
|
let processError: Error | undefined;
|
|
let drainTimer: ReturnType<typeof setTimeout> | undefined;
|
|
let drainDeadline = Infinity;
|
|
let releaseDrain!: () => void;
|
|
const forcedDrain = new Promise<void>(resolve => { releaseDrain = resolve; });
|
|
let releaseExit!: () => void;
|
|
const procExited = new Promise<void>(resolve => { releaseExit = resolve; });
|
|
let releaseStderr!: () => void;
|
|
const stderrClosed = new Promise<void>(resolve => { releaseStderr = resolve; });
|
|
|
|
// Two-phase timeout. Phase 1 (startup): no NDJSON byte yet — a shorter
|
|
// deadline kills a non-answering API run EARLY and names it, instead of
|
|
// the old single timer burning the full work budget to produce an opaque
|
|
// '0 turns / $0.00' failure. Phase 2 (work): armed by the read loop when
|
|
// the FIRST byte arrives, for the REMAINING budget — model work always
|
|
// stays inside timeout; pipe cleanup has a separate bounded allowance.
|
|
let stderr = '';
|
|
let exitReason = 'unknown';
|
|
let timedOut = false;
|
|
let timedOutInStartup = false;
|
|
let phaseTimer: ReturnType<typeof setTimeout>;
|
|
|
|
const closePipes = () => {
|
|
reader.cancel().catch(() => { /* already closed */ });
|
|
proc.stdin!.destroy();
|
|
proc.stdout!.destroy();
|
|
proc.stderr!.destroy();
|
|
};
|
|
const expireDrain = () => {
|
|
drainExpired ||= !stdoutDone || !stderrDone;
|
|
killProcessGroup(proc, 'SIGKILL');
|
|
closePipes();
|
|
releaseDrain();
|
|
};
|
|
const armDrain = () => {
|
|
// Exit after cancellation must not restart the five-second allowance.
|
|
drainDeadline = Math.min(drainDeadline, Date.now() + SESSION_DRAIN_GRACE_MS);
|
|
clearTimeout(drainTimer);
|
|
drainTimer = setTimeout(expireDrain, Math.max(0, drainDeadline - Date.now()));
|
|
};
|
|
|
|
const killRun = (startupPhase: boolean): void => {
|
|
// Labeling and unblocking are SEPARATE concerns: a timer firing after
|
|
// the child already exited must not relabel a real exit (auth error,
|
|
// crash) as a timeout — but it must STILL group-kill and cancel the
|
|
// reader, or an orphan holding the pipes re-creates the exact
|
|
// blocked-drain hang this runner fixed (an early `return` here was the
|
|
// bug the adversarial pass caught in the first version of this guard).
|
|
if (exitCode === undefined) {
|
|
if (!timedOut) timedOutInStartup = startupPhase;
|
|
timedOut = true;
|
|
} else drainExpired ||= !stdoutDone || !stderrDone;
|
|
// Group SIGKILL (mirrors runShardChild): claude AND every tool
|
|
// subprocess it spawned die together — a bare proc.kill() left orphans
|
|
// that inherited our stdout/stderr pipes and kept the API burning
|
|
// (observed: a 600s timeout stretching past 1400s while an orphan held
|
|
// the pipes open).
|
|
killProcessGroup(proc, 'SIGKILL');
|
|
// Belt and braces with the group kill: even if an orphan survives (EPERM
|
|
// fallback path), cancel() unblocks the read loop below.
|
|
closePipes();
|
|
armDrain();
|
|
};
|
|
const onAbort = () => killRun(false);
|
|
const onExit = (code: number | null, exitSignal: NodeJS.Signals | null) => {
|
|
exitCode = code ?? (exitSignal ? 128 + (os.constants.signals[exitSignal] ?? 0) : 1);
|
|
clearTimeout(phaseTimer);
|
|
armDrain();
|
|
releaseExit();
|
|
};
|
|
const onError = (error: Error) => {
|
|
processError = error;
|
|
stderr += error.message;
|
|
exitCode = 1;
|
|
closePipes();
|
|
releaseExit();
|
|
releaseDrain();
|
|
};
|
|
const onStderr = (chunk: string) => { stderr += chunk; };
|
|
const onStderrDone = () => { stderrDone = true; releaseStderr(); };
|
|
const onStderrEnd = () => { stderrEnded = true; onStderrDone(); };
|
|
const onStderrClose = () => {
|
|
// Close releases the drain, but only end proves stderr reached EOF.
|
|
if (!stderrEnded) streamError ??= new Error('stderr closed before EOF');
|
|
onStderrDone();
|
|
};
|
|
const onStreamError = (error: Error) => { streamError = error; };
|
|
proc.on('exit', onExit);
|
|
proc.on('error', onError);
|
|
proc.stderr!.setEncoding('utf8');
|
|
proc.stderr!.on('data', onStderr);
|
|
proc.stderr!.on('end', onStderrEnd).on('close', onStderrClose).on('error', onStreamError);
|
|
proc.stdout!.on('error', onStreamError);
|
|
signal?.addEventListener('abort', onAbort, { once: true });
|
|
phaseTimer = setTimeout(() => killRun(true), Math.max(0, startTime + startupGraceMs - Date.now()));
|
|
proc.stdin!.on('error', () => { /* exit handling reports early child failure */ });
|
|
if (signal?.aborted || Date.now() >= deadline) onAbort();
|
|
else if (!options.nativeLifecycle) proc.stdin!.end(prompt);
|
|
/** Called once by the read loop on the first NDJSON byte. */
|
|
const armWorkPhase = (elapsedMs: number): void => {
|
|
clearTimeout(phaseTimer);
|
|
if (exitCode === undefined && !timedOut) {
|
|
phaseTimer = setTimeout(() => killRun(false), Math.max(0, timeout - elapsedMs));
|
|
}
|
|
};
|
|
|
|
// Stream NDJSON from stdout for real-time progress
|
|
const collectedLines: string[] = [];
|
|
let liveTurnCount = 0;
|
|
let liveToolCount = 0;
|
|
let firstResponseMs = 0;
|
|
let workPhaseArmed = false;
|
|
let lastToolTime = 0;
|
|
let maxInterTurnMs = 0;
|
|
const decoder = new TextDecoder();
|
|
let buf = '';
|
|
const projectLine = options.publicStreamDiagnostics ? publicStreamProjection(startTime) : (line: string) => line;
|
|
let lifecycleFailure: unknown;
|
|
|
|
try {
|
|
options.nativeLifecycle?.onSpawn(proc.pid!);
|
|
if (options.nativeLifecycle && !signal?.aborted && Date.now() < deadline) proc.stdin!.end(prompt);
|
|
try {
|
|
while (true) {
|
|
const { done, value } = await reader.read();
|
|
if (done) break;
|
|
buf += decoder.decode(value, { stream: true });
|
|
const lines = buf.split('\n');
|
|
buf = lines.pop() || '';
|
|
for (const rawLine of lines) {
|
|
if (!rawLine.trim()) continue;
|
|
const line = projectLine(rawLine);
|
|
collectedLines.push(line);
|
|
|
|
// Track time to first NDJSON line (measures latency from spawn to first Claude response)
|
|
if (!workPhaseArmed) {
|
|
// Flag, not `firstResponseMs === 0`: a first line landing in the
|
|
// same millisecond as spawn would read as "not yet seen" and leave
|
|
// the startup timer live for the whole run (claude adversarial).
|
|
workPhaseArmed = true;
|
|
firstResponseMs = Date.now() - startTime;
|
|
// First byte: startup phase over — arm the work phase for the
|
|
// REMAINING budget (total wall stays <= timeout).
|
|
armWorkPhase(firstResponseMs);
|
|
}
|
|
|
|
// Real-time progress to stderr + persistent logs
|
|
try {
|
|
const event = JSON.parse(line);
|
|
if (event.type === 'assistant') {
|
|
liveTurnCount++;
|
|
const content = event.message?.content || [];
|
|
for (const item of content) {
|
|
if (item.type === 'tool_use') {
|
|
liveToolCount++;
|
|
const now = Date.now();
|
|
const elapsed = Math.round((now - startTime) / 1000);
|
|
// Track inter-turn latency (tool call to tool call)
|
|
if (lastToolTime > 0) {
|
|
const interTurn = now - lastToolTime;
|
|
if (interTurn > maxInterTurnMs) maxInterTurnMs = interTurn;
|
|
}
|
|
lastToolTime = now;
|
|
const progressLine = ` [${elapsed}s] turn ${liveTurnCount} tool #${liveToolCount}: ${item.name}(${truncate(JSON.stringify(item.input || {}), 80)})\n`;
|
|
process.stderr.write(progressLine);
|
|
|
|
// Persist progress.log
|
|
if (runDir) {
|
|
try { fs.appendFileSync(path.join(runDir, 'progress.log'), progressLine); } catch { /* non-fatal */ }
|
|
}
|
|
|
|
// Write heartbeat (atomic)
|
|
if (runId && testName) {
|
|
try {
|
|
const toolDesc = `${item.name}(${truncate(JSON.stringify(item.input || {}), 60)})`;
|
|
atomicWriteSync(HEARTBEAT_PATH, JSON.stringify({
|
|
runId,
|
|
pid: proc.pid,
|
|
startedAt,
|
|
currentTest: testName,
|
|
status: 'running',
|
|
turn: liveTurnCount,
|
|
toolCount: liveToolCount,
|
|
lastTool: toolDesc,
|
|
lastToolAt: new Date().toISOString(),
|
|
elapsedSec: elapsed,
|
|
}, null, 2) + '\n');
|
|
} catch { /* non-fatal */ }
|
|
}
|
|
}
|
|
}
|
|
}
|
|
} catch { /* skip — parseNDJSON will handle it later */ }
|
|
|
|
// Append raw NDJSON line to per-test transcript file
|
|
if (runDir && safeName) {
|
|
try { fs.appendFileSync(path.join(runDir, `${safeName}.ndjson`), line + '\n'); } catch { /* non-fatal */ }
|
|
}
|
|
}
|
|
}
|
|
} catch (error) { streamError = error as Error; }
|
|
stdoutDone = true;
|
|
|
|
// Flush remaining buffer
|
|
if (buf.trim()) {
|
|
const line = projectLine(buf);
|
|
collectedLines.push(line);
|
|
if (options.publicStreamDiagnostics && runDir && safeName) {
|
|
try { fs.appendFileSync(path.join(runDir, `${safeName}.ndjson`), line + '\n'); } catch { /* non-fatal */ }
|
|
}
|
|
}
|
|
|
|
await Promise.race([Promise.all([procExited, stderrClosed]), forcedDrain]);
|
|
} catch (error) {
|
|
lifecycleFailure = error;
|
|
throw error;
|
|
} finally {
|
|
if (options.nativeLifecycle) {
|
|
killProcessGroup(proc, 'SIGKILL');
|
|
closePipes();
|
|
armDrain();
|
|
try {
|
|
await Promise.race([procExited, forcedDrain]);
|
|
let hookTimer: ReturnType<typeof setTimeout> | undefined;
|
|
try {
|
|
await Promise.race([
|
|
options.nativeLifecycle.onSettled({ deadline: drainDeadline, exited: exitCode !== undefined && !processError }),
|
|
new Promise<never>((_, reject) => {
|
|
hookTimer = setTimeout(() => reject(new Error('native lifecycle settlement deadline exceeded')), Math.max(0, drainDeadline - Date.now()));
|
|
}),
|
|
]);
|
|
} finally { clearTimeout(hookTimer); }
|
|
} catch (error) {
|
|
if (lifecycleFailure) throw new AggregateError([lifecycleFailure, error], 'native lifecycle failed');
|
|
throw error;
|
|
} finally {
|
|
clearTimeout(phaseTimer);
|
|
clearTimeout(drainTimer);
|
|
signal?.removeEventListener('abort', onAbort);
|
|
proc.removeListener('exit', onExit);
|
|
proc.stderr!.removeListener('data', onStderr);
|
|
}
|
|
}
|
|
clearTimeout(phaseTimer);
|
|
clearTimeout(drainTimer);
|
|
signal?.removeEventListener('abort', onAbort);
|
|
killProcessGroup(proc, 'SIGKILL');
|
|
closePipes();
|
|
proc.removeListener('exit', onExit);
|
|
proc.stderr!.removeListener('data', onStderr);
|
|
}
|
|
|
|
if (timedOut) {
|
|
// 'timeout_startup' = the API never sent a byte inside the grace — an
|
|
// availability problem, not a test failure worth reading transcripts
|
|
// for. Distinct so triage (and WS10's inconclusive classification) can
|
|
// key off it without receipts archaeology.
|
|
exitReason = timedOutInStartup ? 'timeout_startup' : 'timeout';
|
|
} else if (exitCode === 0) {
|
|
exitReason = drainExpired ? 'error_output_drain' : streamError ? 'error_output_stream' : 'success';
|
|
if (drainExpired) stderr += `\nOutput drain exceeded ${SESSION_DRAIN_GRACE_MS}ms after exit 0`;
|
|
else if (streamError) stderr += `\nOutput stream failed: ${streamError.message}`;
|
|
} else {
|
|
exitReason = `exit_code_${exitCode ?? 1}`;
|
|
}
|
|
|
|
const duration = Date.now() - startTime;
|
|
|
|
// Parse all collected NDJSON lines
|
|
const parsed = parseNDJSON(collectedLines);
|
|
const { transcript, resultLine, toolCalls } = parsed;
|
|
const browseErrors: string[] = [];
|
|
|
|
const allText = toolCalls.filter(call => call.tool === 'Bash').map(call => call.output).join('\n') + '\n' + stderr;
|
|
for (const pattern of BROWSE_ERROR_PATTERNS) {
|
|
const match = allText.match(pattern);
|
|
if (match) {
|
|
browseErrors.push(match[0].slice(0, 200));
|
|
}
|
|
}
|
|
|
|
// Native Claude uses exit 1 for its structured max-turns result. Preserve
|
|
// that existing semantic outcome only after a complete, uncancelled drain;
|
|
// success-shaped payloads must never override a process/stream failure.
|
|
const maxTurnsExit = exitCode === 1 && resultLine?.subtype === 'error_max_turns'
|
|
&& resultLine.is_error === true && !timedOut && !signal?.aborted
|
|
&& stdoutDone && stderrDone && !drainExpired && !streamError && !processError;
|
|
if (maxTurnsExit) {
|
|
exitReason = 'error_max_turns';
|
|
} else if (resultLine && exitReason === 'success') {
|
|
if (resultLine.subtype === 'success' && resultLine.is_error) {
|
|
// claude -p can return subtype=success with is_error=true (e.g. API connection failure)
|
|
exitReason = 'error_api';
|
|
} else if (resultLine.subtype === 'success' && exitCode === 0 && !timedOut) {
|
|
exitReason = 'success';
|
|
} else if (resultLine.subtype && resultLine.subtype !== 'success') {
|
|
// Preserve known subtypes like error_max_turns even if is_error is set
|
|
exitReason = resultLine.subtype;
|
|
}
|
|
}
|
|
|
|
// Save failure transcript to persistent run directory (or fallback to workingDirectory)
|
|
if (browseErrors.length > 0 || exitReason !== 'success') {
|
|
try {
|
|
const failureDir = runDir || path.join(workingDirectory, '.gstack', 'test-transcripts');
|
|
fs.mkdirSync(failureDir, { recursive: true });
|
|
const failureName = safeName
|
|
? `${safeName}-failure.json`
|
|
: `e2e-${new Date().toISOString().replace(/[:.]/g, '-')}.json`;
|
|
fs.writeFileSync(
|
|
path.join(failureDir, failureName),
|
|
JSON.stringify({
|
|
prompt: prompt.slice(0, 500),
|
|
testName: testName || 'unknown',
|
|
exitReason,
|
|
browseErrors,
|
|
duration,
|
|
turnAtTimeout: timedOut ? liveTurnCount : undefined,
|
|
lastToolCall: liveToolCount > 0 ? `tool #${liveToolCount}` : undefined,
|
|
stderr: stderr.slice(0, 2000),
|
|
result: resultLine ? { type: resultLine.type, subtype: resultLine.subtype, result: resultLine.result?.slice?.(0, 500) } : null,
|
|
}, null, 2),
|
|
);
|
|
} catch { /* non-fatal */ }
|
|
}
|
|
|
|
// Cost from result line (exact) or estimate from chars
|
|
const turnsUsed = resultLine?.num_turns || 0;
|
|
const estimatedCost = resultLine?.total_cost_usd || 0;
|
|
const inputChars = prompt.length;
|
|
const outputChars = (resultLine?.result || '').length;
|
|
const estimatedTokens = (resultLine?.usage?.input_tokens || 0)
|
|
+ (resultLine?.usage?.output_tokens || 0)
|
|
+ (resultLine?.usage?.cache_read_input_tokens || 0);
|
|
|
|
const costEstimate: CostEstimate = {
|
|
inputChars,
|
|
outputChars,
|
|
estimatedTokens,
|
|
estimatedCost: Math.round((estimatedCost) * 100) / 100,
|
|
turnsUsed,
|
|
};
|
|
|
|
return { toolCalls, browseErrors, exitReason, duration, output: resultLine?.result || '', costEstimate, transcript, model, firstResponseMs, maxInterTurnMs };
|
|
}
|