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
28 changes: 27 additions & 1 deletion kernel/relayflowd/src/exec_det.rs
Original file line number Diff line number Diff line change
Expand Up @@ -118,14 +118,18 @@ 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,
budget: Budget::default(),
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())),
Expand Down Expand Up @@ -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]
Expand Down
68 changes: 68 additions & 0 deletions packages/sdk/src/cli/deterministic-failure.ts
Original file line number Diff line number Diff line change
@@ -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<StepFailedDetails | undefined> {
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<string, StepFailedDetails>();
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<string, unknown> | undefined {
return value !== null && typeof value === 'object' && !Array.isArray(value)
? value as Record<string, unknown> : 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, "'\\''")}'`;
}
30 changes: 23 additions & 7 deletions packages/sdk/src/cli/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -32,7 +33,7 @@ export interface ParkedStep {
type: Extract<StepType, 'llm' | 'agent'>;
}

export interface RunDiagnostic {
export interface RunDiagnostic extends StepFailedDetails {
severity: 'refusal' | 'failure' | 'parked' | 'warning';
kind: RunFailureKind | RunWarningKind | RunCompletionReason;
message: string;
Expand Down Expand Up @@ -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],
},
};
}
Expand Down
9 changes: 9 additions & 0 deletions packages/sdk/src/failure-kinds.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string> = new Set(CHECK_FAILURE_KINDS);
const RUN_FAILURE_KIND_SET: ReadonlySet<string> = new Set(RUN_FAILURE_KINDS);

Expand Down
174 changes: 174 additions & 0 deletions packages/sdk/tests/deterministic-failure-diagnostic.test.ts
Original file line number Diff line number Diff line change
@@ -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'),
});
});
});
Loading