From 264c4580cac0c731a204ff4aae48e164bd7e5422 Mon Sep 17 00:00:00 2001 From: Miya Date: Sat, 12 Sep 2026 15:13:30 +0200 Subject: [PATCH] feat(sdk): deterministic failure diagnostic surfaces exit code + stderr excerpt (#276) --- kernel/relayflowd/src/exec_det.rs | 28 ++- packages/sdk/src/cli/deterministic-failure.ts | 68 +++++++ packages/sdk/src/cli/run.ts | 30 ++- packages/sdk/src/failure-kinds.ts | 9 + .../deterministic-failure-diagnostic.test.ts | 174 ++++++++++++++++++ 5 files changed, 301 insertions(+), 8 deletions(-) create mode 100644 packages/sdk/src/cli/deterministic-failure.ts create mode 100644 packages/sdk/tests/deterministic-failure-diagnostic.test.ts diff --git a/kernel/relayflowd/src/exec_det.rs b/kernel/relayflowd/src/exec_det.rs index 8445c793e..e907945d7 100644 --- a/kernel/relayflowd/src/exec_det.rs +++ b/kernel/relayflowd/src/exec_det.rs @@ -118,6 +118,10 @@ pub(crate) fn execute_placed_with_input( "stdout_tail": tail(&stdout), "stderr_tail": tail(&stderr), }); + // Failed completions deliberately null their reusable output. Preserve + // command evidence in the existing diagnostic field before that happens. + let trajectory_tail = (output["exit_code"] != 0) + .then(|| json!({ "exit_code": output["exit_code"], "stderr_tail": output["stderr_tail"] })); AttemptResult { human_intervention: false, output, @@ -125,7 +129,7 @@ pub(crate) fn execute_placed_with_input( completed_by: "kernel".to_owned(), end_pins: None, effects: vec![], - trajectory_tail: None, + trajectory_tail, failure_reason: timed_out.then_some(CompletionReason::Timeout), failure_detail: timed_out .then(|| format!("step exceeded its {} ms timeout", timeout.as_millis())), @@ -189,6 +193,28 @@ mod tests { assert_eq!(result.output["exit_code"], 0); assert_eq!(result.output["stdout_tail"], "hello"); assert_eq!(result.failure_reason, None); + assert_eq!(result.trajectory_tail, None); + } + + #[test] + fn failed_command_evidence_survives_completion() { + use relayflowd_core::machine::{Action, completion_actions}; + + let step: StepSpec = serde_json::from_value(json!({ + "id": "fail-command", "type": "deterministic", + "command": "printf 'shakedown intentional failure' >&2; exit 7" + })) + .unwrap(); + let actions = completion_actions("run", &step, 1, 0, execute(&step), 0); + let Action::Append(completed) = &actions[0] else { + panic!("expected completion") + }; + assert_eq!(completed.payload["output"], serde_json::Value::Null); + assert_eq!(completed.payload["trajectory_tail"]["exit_code"], 7); + assert_eq!( + completed.payload["trajectory_tail"]["stderr_tail"], + "shakedown intentional failure" + ); } #[test] diff --git a/packages/sdk/src/cli/deterministic-failure.ts b/packages/sdk/src/cli/deterministic-failure.ts new file mode 100644 index 000000000..53173d386 --- /dev/null +++ b/packages/sdk/src/cli/deterministic-failure.ts @@ -0,0 +1,68 @@ +import type { StepFailedDetails } from '../failure-kinds.js'; +import type { JournalClient } from '../journal-client.js'; + +const STDERR_TAIL_BYTES = 1_024; + +/** Read completion evidence through the existing protocol, including later pages. */ +export async function deterministicFailureDetails( + client: JournalClient, + runId: string, +): Promise { + const snapshot = await client.runGet(runId); + if (!Object.values(snapshot.steps).some(step => step.type === 'deterministic')) return undefined; + let fromSeq = 1; + const failures = new Map(); + while (true) { + const { entries } = await client.journalRead(runId, fromSeq, 100); + if (entries.length === 0) break; + for (const raw of entries) { + const entry = record(raw); + const seq = entry?.['seq']; + if (typeof seq !== 'number' || !Number.isSafeInteger(seq) || seq < fromSeq) { + throw new Error('invalid journal sequence in command failure inspection'); + } + fromSeq = seq + 1; + const stepId = entry?.['step_id']; + if (entry?.['entry_type'] !== 'step.completed' || typeof stepId !== 'string' + || snapshot.steps[stepId]?.type !== 'deterministic') continue; + // A later completion supersedes an earlier failed attempt. + failures.delete(stepId); + const payload = record(entry['payload']); + // Failed completions null reusable output; command evidence survives in + // trajectory_tail. Accept output as well for existing completion records. + const output = record(payload?.['trajectory_tail']) ?? record(payload?.['output']); + const exitCode = output?.['exit_code']; + if (payload?.['disposition'] !== 'step_done' || payload['completionReason'] === 'success' + || typeof exitCode !== 'number' || !Number.isSafeInteger(exitCode) || exitCode === 0) continue; + failures.set(stepId, { + stepId, + exitCode, + stderrTail: stderrTail(output?.['stderr_tail']), + hint: `flows replay ${shellQuote(runId)} --at ${shellQuote(stepId)}`, + }); + } + } + return [...failures.values()].at(-1); +} + +function record(value: unknown): Record | undefined { + return value !== null && typeof value === 'object' && !Array.isArray(value) + ? value as Record : undefined; +} + +function stderrTail(value: unknown): string { + if (typeof value !== 'string') return ''; + const bytes = Buffer.from(value, 'utf8'); + let start = Math.max(0, bytes.length - STDERR_TAIL_BYTES); + // Drop a partial leading code point, avoiding replacement-byte expansion. + while (start < bytes.length && (bytes[start]! & 0xc0) === 0x80) start += 1; + // Preserve tabs/newlines; replace binary controls (including ESC and CR), + // C1 controls and Unicode formatting controls without growing the excerpt. + return bytes.subarray(start).toString('utf8') + .replace(/[\p{Cc}\p{Cf}]/gu, character => character === '\n' || character === '\t' ? character : '?'); +} + +function shellQuote(value: string): string { + return /^[A-Za-z0-9_-]+$/.test(value) && !value.startsWith('-') + ? value : `'${value.replace(/'/g, "'\\''")}'`; +} diff --git a/packages/sdk/src/cli/run.ts b/packages/sdk/src/cli/run.ts index baabbf83e..f340491f7 100644 --- a/packages/sdk/src/cli/run.ts +++ b/packages/sdk/src/cli/run.ts @@ -10,7 +10,8 @@ import { socketPathFor } from '../daemon-connection.js'; import { ensureDaemon, type EnsureDaemonOptions } from '../daemon-lifecycle.js'; import { isAuthoredFlowPath } from '../direct-input.js'; import { daemonRefusal } from './daemon-refusal.js'; -import type { RunFailureKind, RunWarningKind } from '../failure-kinds.js'; +import type { RunFailureKind, RunWarningKind, StepFailedDetails } from '../failure-kinds.js'; +import { deterministicFailureDetails } from './deterministic-failure.js'; import { JournalClient, JournalProtocolError } from '../journal-client.js'; import { attachLocalAgent } from '../local-agent.js'; import type { @@ -32,7 +33,7 @@ export interface ParkedStep { type: Extract; } -export interface RunDiagnostic { +export interface RunDiagnostic extends StepFailedDetails { severity: 'refusal' | 'failure' | 'parked' | 'warning'; kind: RunFailureKind | RunWarningKind | RunCompletionReason; message: string; @@ -345,15 +346,30 @@ export async function classifyOutcome( if (report.ok) return { exitCode: 0, report }; if (current.status === 'failed' && current.completion_reason !== null) { + const diagnostic: RunDiagnostic = { + severity: 'failure', + kind: current.completion_reason, + message: `Run "${current.run_id}" failed with completionReason: ${current.completion_reason}.`, + }; + if (current.completion_reason === 'step_failed') { + try { + const details = await deterministicFailureDetails(client, current.run_id); + if (details !== undefined) { + Object.assign(diagnostic, details); + diagnostic.message += ` Step ${JSON.stringify(details.stepId)} exit=${details.exitCode}.` + + (details.stderrTail ? `\nStderr (last 1,024 bytes):\n${details.stderrTail}` : '') + + `\nInspect: ${details.hint}`; + } + } catch (error) { + // Inspection must not erase the already known run failure. + diagnostic.message += ` Could not inspect command failure: ${errorMessage(error)}`; + } + } return { exitCode: 1, report: { ...report, - diagnostics: [...report.diagnostics, { - severity: 'failure', - kind: current.completion_reason, - message: `Run "${current.run_id}" failed with completionReason: ${current.completion_reason}.`, - }], + diagnostics: [...report.diagnostics, diagnostic], }, }; } diff --git a/packages/sdk/src/failure-kinds.ts b/packages/sdk/src/failure-kinds.ts index e6c73fd92..b301e5f89 100644 --- a/packages/sdk/src/failure-kinds.ts +++ b/packages/sdk/src/failure-kinds.ts @@ -118,6 +118,15 @@ export type PreflightWarningKind = (typeof PREFLIGHT_WARNING_KINDS)[number]; export type RunFailureKind = (typeof RUN_FAILURE_KINDS)[number]; export type RunWarningKind = (typeof RUN_WARNING_KINDS)[number]; +/** Optional evidence on the existing step_failed diagnostic, not a new kind. */ +export interface StepFailedDetails { + stepId?: string; + exitCode?: number; + /** Terminal-safe UTF-8 excerpt, at most 1,024 bytes. */ + stderrTail?: string; + hint?: string; +} + const CHECK_FAILURE_KIND_SET: ReadonlySet = new Set(CHECK_FAILURE_KINDS); const RUN_FAILURE_KIND_SET: ReadonlySet = new Set(RUN_FAILURE_KINDS); diff --git a/packages/sdk/tests/deterministic-failure-diagnostic.test.ts b/packages/sdk/tests/deterministic-failure-diagnostic.test.ts new file mode 100644 index 000000000..fbcc060e1 --- /dev/null +++ b/packages/sdk/tests/deterministic-failure-diagnostic.test.ts @@ -0,0 +1,174 @@ +import { spawnSync } from 'node:child_process'; +import { mkdtempSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { runCli } from '../src/cli.js'; +import { classifyOutcome, emptyReport, type RunDiagnostic, type RunReport } from '../src/cli/run.js'; +import { socketPathFor } from '../src/daemon-connection.js'; +import * as daemonLifecycle from '../src/daemon-lifecycle.js'; +import { JournalClient } from '../src/journal-client.js'; +import { PROTOCOL_VERSION, type RunOutcome } from '../src/protocol.js'; + +const failure: RunOutcome = { + run_id: 'run-failed', status: 'failed', completion_reason: 'step_failed', completed_steps: 0, +}; +const directories: string[] = []; + +afterEach(() => { + vi.restoreAllMocks(); + for (const dir of directories.splice(0)) rmSync(dir, { recursive: true, force: true }); +}); + +function completion(seq: number, stderr = '', exitCode = 7, stepId = 'fail-command', disposition = 'step_done') { + return { + seq, entry_type: 'step.completed', step_id: stepId, + payload: { + completionReason: exitCode === 0 ? 'success' : 'verification_failed', disposition, + output: exitCode === 0 ? { exit_code: exitCode, stderr_tail: stderr } : null, + ...(exitCode !== 0 ? { trajectory_tail: { exit_code: exitCode, stderr_tail: stderr } } : {}), + }, + }; +} + +function stub(pages: unknown[][], type: 'deterministic' | 'agent' | 'llm' = 'deterministic') { + const runGet = vi.fn(async () => ({ steps: { 'fail-command': { type, state: 'done' } } })); + const journalRead = vi.fn(async () => ({ entries: pages.shift() ?? [] })); + return { client: { runGet, journalRead } as unknown as JournalClient, runGet, journalRead }; +} + +async function classify(client: JournalClient, outcome = failure) { + return classifyOutcome(client, 'run', outcome, emptyReport('run'), '/tmp/diagnostic.sock', {}); +} + +async function diagnosticFor(stderr: string) { + const { client } = stub([[completion(1, stderr)]]); + return (await classify(client)).report.diagnostics.at(-1) as RunDiagnostic; +} + +describe('deterministic failure diagnostic', () => { + it.each([false, true])('surfaces command exit, stderr and replay hint through the CLI (json=%s)', async json => { + const dir = mkdtempSync(join(tmpdir(), 'flows-failure-')); + directories.push(dir); + writeFileSync(join(dir, 'flows.json'), '{}'); + const command = 'printf "shakedown intentional failure" >&2; exit 7'; + const path = join(dir, 'runtime-error.yaml'); + writeFileSync(path, `version: "0.1.0"\nname: runtime-error\nsteps:\n - id: fail-command\n type: deterministic\n command: '${command}'\n`); + // Execute the repro and supply the kernel's completion shape at the + // client boundary. This tests CLI formatting independently of transport. + const result = spawnSync('/bin/sh', ['-c', command], { encoding: 'utf8' }); + expect(result.status).toBe(7); + vi.spyOn(daemonLifecycle, 'ensureDaemon').mockResolvedValue({ + kind: 'attached', socketPath: socketPathFor(dir), connection: null, + }); + vi.spyOn(JournalClient.prototype, 'connect').mockResolvedValue(undefined); + vi.spyOn(JournalClient.prototype, 'hello').mockResolvedValue({ protocol: PROTOCOL_VERSION, server: 'relayflowd-test' }); + vi.spyOn(JournalClient.prototype, 'runStart').mockResolvedValue(failure); + vi.spyOn(JournalClient.prototype, 'runGet').mockResolvedValue({ + run_id: failure.run_id, status: 'failed', budget: { tokens_in: 0, tokens_out: 0, dollars: '0' }, + steps: { 'fail-command': { type: 'deterministic', state: 'done' } }, + }); + vi.spyOn(JournalClient.prototype, 'journalRead').mockImplementation(async (_runId, fromSeq) => ({ + entries: fromSeq === 1 ? [completion(1, result.stderr, result.status!)] : [], + })); + const stdout: string[] = []; + const stderr: string[] = []; + const code = await runCli(['run', '--no-spawn', '--data-dir', dir, ...(json ? ['--json'] : []), path], { + stdout: line => stdout.push(line), stderr: line => stderr.push(line), + }); + expect(code).toBe(1); + if (json) { + const report = JSON.parse(stdout.join('\n')) as RunReport; + expect(report.diagnostics).toContainEqual(expect.objectContaining({ + kind: 'step_failed', stepId: 'fail-command', exitCode: 7, + stderrTail: 'shakedown intentional failure', + hint: 'flows replay run-failed --at fail-command', + })); + } else { + const output = stderr.join('\n'); + expect(output).toContain('FAILED [step_failed]'); + expect(output).toContain('exit=7'); + expect(output).toContain('shakedown intentional failure'); + expect(output).toContain('flows replay run-failed --at fail-command'); + } + }); + + it('keeps the last 1,024 bytes, not the prefix', async () => { + const diagnostic = await diagnosticFor('discarded-prefix' + 'x'.repeat(2_000) + 'last error'); + expect(diagnostic.stderrTail).toBe('x'.repeat(1_014) + 'last error'); + expect(Buffer.byteLength(diagnostic.stderrTail!)).toBe(1_024); + expect(diagnostic.message).not.toContain('discarded-prefix'); + }); + + it('truncates at UTF-8 boundaries without splitting multibyte characters', async () => { + const diagnostic = await diagnosticFor('🙂'.repeat(400) + 'é'); + expect(diagnostic.stderrTail).toBe('🙂'.repeat(255) + 'é'); + expect(Buffer.byteLength(diagnostic.stderrTail!)).toBeLessThanOrEqual(1_024); + expect(diagnostic.stderrTail).not.toContain('\ufffd'); + }); + + it('replaces binary and terminal controls and remains byte bounded', async () => { + const binary = Buffer.from([0, 0xff, 0x1b, 13, 8]).toString('utf8'); + const diagnostic = await diagnosticFor(binary.repeat(1_000) + '\u009b31m\u202eEND\n\t'); + expect(diagnostic.stderrTail).toMatch(/END\n\t$/); + expect(diagnostic.stderrTail).not.toMatch(/[\x00-\x08\x0b-\x1f\x7f-\x9f\p{Cf}]/u); + expect(Buffer.byteLength(diagnostic.stderrTail!)).toBeLessThanOrEqual(1_024); + }); + + it('includes an empty stderr field for a silent nonzero exit', async () => { + expect(await diagnosticFor('')).toMatchObject({ exitCode: 7, stderrTail: '' }); + }); + + it('reads later pages and ignores failed attempts superseded by success', async () => { + const { client, journalRead } = stub([ + [completion(1, 'old error', 8, 'fail-command', 'retry')], + [completion(2, '', 0)], + [completion(3, 'final error', 7)], + ]); + const report = (await classify(client)).report; + expect(report.diagnostics.at(-1)).toMatchObject({ exitCode: 7, stderrTail: 'final error' }); + expect(journalRead.mock.calls).toEqual([ + ['run-failed', 1, 100], ['run-failed', 2, 100], ['run-failed', 3, 100], ['run-failed', 4, 100], + ]); + }); + + it('does not report recovered attempts or successful command output', async () => { + const { client } = stub([[completion(1, 'old error'), completion(2, 'success stderr', 0)]]); + expect((await classify(client)).report.diagnostics.at(-1)).not.toHaveProperty('exitCode'); + }); + + it.each(['agent', 'llm'] as const)('leaves %s failure diagnostics unchanged', async type => { + const { client, journalRead } = stub([[completion(1, 'worker stderr')]], type); + expect((await classify(client)).report.diagnostics.at(-1)).toEqual({ + severity: 'failure', kind: 'step_failed', + message: 'Run "run-failed" failed with completionReason: step_failed.', + }); + expect(journalRead).not.toHaveBeenCalled(); + }); + + it('does not inspect the journal or change diagnostics on success', async () => { + const { client, runGet, journalRead } = stub([]); + const execution = await classify(client, { ...failure, status: 'completed', completion_reason: 'success', completed_steps: 1 }); + expect(execution.exitCode).toBe(0); + expect(execution.report.diagnostics).toEqual([]); + expect(runGet).not.toHaveBeenCalled(); + expect(journalRead).not.toHaveBeenCalled(); + }); + + it('preserves step_failed and explains failed journal inspection', async () => { + const { client, journalRead } = stub([]); + journalRead.mockRejectedValueOnce(new Error('journal unavailable')); + const execution = await classify(client); + expect(execution.exitCode).toBe(1); + expect(execution.report.diagnostics.at(-1)).toMatchObject({ + kind: 'step_failed', message: expect.stringContaining('Could not inspect command failure: journal unavailable'), + }); + }); + + it('stops on non-advancing journal pages', async () => { + const { client } = stub([[completion(1)], [completion(1)]]); + expect((await classify(client)).report.diagnostics.at(-1)).toMatchObject({ + kind: 'step_failed', message: expect.stringContaining('invalid journal sequence'), + }); + }); +});