Skip to content
Draft
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
47 changes: 33 additions & 14 deletions src/lib/agent/agent-interface.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ import {
AgentErrorType,
REMARK_INSTRUCTION,
RESUME_INSTRUCTION,
runErrorType,
} from './signals';
import { classifyAuthFailure } from '@lib/errors';
import { isGrantRevoked } from '@lib/auth-session-state';
Expand Down Expand Up @@ -865,6 +866,18 @@ export async function runAgent(
return {};
};

// How this run failed, so the `agent aborted` event below can name it.
// Every terminal path sets it through `failWith`; a run that reaches
// `completeWithSuccess` leaves it unset and fires no abort event.
let failureMode: AgentErrorType | undefined;
const failWith = (
error: AgentErrorType,
message?: string,
): { error: AgentErrorType; message?: string } => {
failureMode = error;
return message === undefined ? { error } : { error, message };
};

// Abort controller — lets us force-kill the SDK query when we detect an
// [ABORT] signal in the agent's output. Also stashes the reason so the
// runner can surface it via outroData after we unwind.
Expand Down Expand Up @@ -1325,27 +1338,27 @@ export async function runAgent(
if (yaraViolationReason) {
logToFile('Agent error: YARA_VIOLATION');
spinner.stop('Security check stopped the setup');
return { error: AgentErrorType.YARA_VIOLATION };
return failWith(AgentErrorType.YARA_VIOLATION);
}

// If the middleware caught an [ABORT] and aborted the SDK query, surface
// it as a structured error before checking other signals.
if (abortReason) {
spinner.stop('Wizard aborted');
return { error: AgentErrorType.ABORT, message: abortReason };
return failWith(AgentErrorType.ABORT, abortReason);
}

// Check for error markers in the agent's output
if (signals.has('MCP_MISSING')) {
logToFile('Agent error: MCP_MISSING');
spinner.stop('Agent could not access PostHog MCP');
return { error: AgentErrorType.MCP_MISSING };
return failWith(AgentErrorType.MCP_MISSING);
}

if (signals.has('RESOURCE_MISSING')) {
logToFile('Agent error: RESOURCE_MISSING');
spinner.stop('Agent could not access setup resource');
return { error: AgentErrorType.RESOURCE_MISSING };
return failWith(AgentErrorType.RESOURCE_MISSING);
}

// A clean success result already arrived. The Claude SDK can emit a second
Expand All @@ -1365,13 +1378,13 @@ export async function runAgent(
if (signals.hasApiErrorStatus(429)) {
logToFile('Agent error: RATE_LIMIT');
spinner.stop('Rate limit exceeded');
return { error: AgentErrorType.RATE_LIMIT, message: apiErrorMessage };
return failWith(AgentErrorType.RATE_LIMIT, apiErrorMessage);
}

if (signals.hasApiError()) {
logToFile('Agent error: API_ERROR');
spinner.stop('API error occurred');
return { error: AgentErrorType.API_ERROR, message: apiErrorMessage };
return failWith(AgentErrorType.API_ERROR, apiErrorMessage);
}

return completeWithSuccess();
Expand All @@ -1385,14 +1398,14 @@ export async function runAgent(
if (yaraViolationReason) {
logToFile('Agent error: YARA_VIOLATION');
spinner.stop('Security check stopped the setup');
return { error: AgentErrorType.YARA_VIOLATION };
return failWith(AgentErrorType.YARA_VIOLATION);
}

// If the middleware caught an [ABORT] and triggered abortController.abort(),
// the SDK will throw an AbortError — surface it as a clean abort result.
if (abortReason) {
spinner.stop('Wizard aborted');
return { error: AgentErrorType.ABORT, message: abortReason };
return failWith(AgentErrorType.ABORT, abortReason);
}

// If we already received a successful result, the error is from SDK cleanup
Expand All @@ -1410,27 +1423,33 @@ export async function runAgent(
if (signals.hasApiErrorStatus(429)) {
logToFile('Agent error (caught): RATE_LIMIT');
spinner.stop('Rate limit exceeded');
return { error: AgentErrorType.RATE_LIMIT, message: apiErrorMessage };
return failWith(AgentErrorType.RATE_LIMIT, apiErrorMessage);
}

if (signals.hasApiError()) {
logToFile('Agent error (caught): API_ERROR');
spinner.stop('API error occurred');
return { error: AgentErrorType.API_ERROR, message: apiErrorMessage };
return failWith(AgentErrorType.API_ERROR, apiErrorMessage);
}

// No API error found, re-throw the original exception
// No API error found, re-throw the original exception. The runner codes
// it, but the abort event still needs a mode, so classify it here on the
// same rule the pi harness uses.
failureMode = runErrorType(errorMessage);
spinner.stop(errorMessage);
getUI().log.error(`Error: ${(error as Error).message}`);
logToFile('Agent run failed:', error);
debug('Full error:', error);
throw error;
} finally {
// Always capture run duration, even on abort/error, so we can alert on
// long runs where the user gave up before completion.
if (!receivedSuccessResult) {
// Only a run that actually failed aborts. A run can finish without an SDK
// success result and still complete — keying on that result reported those
// runs as aborts as well, so the event counted far more failures than
// happened and named none of them.
if (failureMode !== undefined) {
const durationMs = Date.now() - startTime;
analytics.wizardCapture('agent aborted', {
failure_mode: failureMode,
duration_ms: durationMs,
duration_seconds: Math.round(durationMs / 1000),
model: agentConfig.model,
Expand Down
94 changes: 92 additions & 2 deletions src/lib/agent/runner/harness/pi/__tests__/completion.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { completionFailure, runErrorType } from '../completion';
import { AgentErrorType } from '@lib/agent/signals';
import { completionFailure, nudgeWhileUnfinished } from '../completion';
import { AgentErrorType, runErrorType } from '@lib/agent/signals';

describe('completionFailure', () => {
it('fails a no-op run (zero tool calls) as NO_PROGRESS', () => {
Expand Down Expand Up @@ -57,3 +57,93 @@ describe('runErrorType', () => {
expect(runErrorType('')).toBe(AgentErrorType.API_ERROR);
});
});

describe('nudgeWhileUnfinished', () => {
const run = (args: {
max: number;
unfinished: () => boolean;
progress: () => number;
send: (nudge: number) => Promise<void>;
}) => nudgeWhileUnfinished({ ...args, backoffMs: 0 });

it('stops on the first nudge that produces no tool call and no output', async () => {
let sent = 0;
const result = await run({
max: 20,
unfinished: () => true,
// Work never advances: every nudge came back empty.
progress: () => 0,
send: () => {
sent += 1;
return Promise.resolve();
},
});
expect(result).toEqual({ nudges: 1, dead: true });
expect(sent).toBe(1);
});

it('keeps nudging while each nudge does work, then stops when finished', async () => {
let work = 0;
let open = 3;
const result = await run({
max: 20,
unfinished: () => open > 0,
progress: () => work,
send: () => {
work += 1;
open -= 1;
return Promise.resolve();
},
});
expect(result).toEqual({ nudges: 3, dead: false });
});

it('never sends more nudges than the cap', async () => {
let work = 0;
const result = await run({
max: 4,
unfinished: () => true,
progress: () => work,
send: () => {
work += 1;
return Promise.resolve();
},
});
expect(result).toEqual({ nudges: 4, dead: false });
});

it('sends nothing when the work is already finished', async () => {
const send = vi.fn();
const result = await run({
max: 20,
unfinished: () => false,
progress: () => 0,
send: () => {
send();
return Promise.resolve();
},
});
expect(result).toEqual({ nudges: 0, dead: false });
expect(send).not.toHaveBeenCalled();
});

it('waits between live nudges instead of spinning', async () => {
let work = 0;
let open = 3;
const started = Date.now();
const result = await nudgeWhileUnfinished({
max: 20,
backoffMs: 20,
unfinished: () => open > 0,
progress: () => work,
send: () => {
work += 1;
open -= 1;
return Promise.resolve();
},
});
expect(result.nudges).toBe(3);
// Two waits between three nudges.
expect(Date.now() - started).toBeGreaterThanOrEqual(35);
});
});
52 changes: 41 additions & 11 deletions src/lib/agent/runner/harness/pi/completion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,18 +13,48 @@ export function completionFailure(args: {
}

/**
* Which error type a thrown pi run reports.
* Pause between completion-guard nudges. A dead turn comes back in
* milliseconds, so without it the guard spins through its whole cap in
* seconds and the run aborts before the model had any chance to act.
*/
export const NUDGE_BACKOFF_MS = 1_000;

export interface NudgeRun {
/** Nudges sent. */
nudges: number;
/** True when a nudge came back with no tool call and no assistant output. */
dead: boolean;
}

/**
* Re-prompt a pi session while its work is unfinished.
*
* Both pi entry points classified the caught message inline with the same two
* needles, and both did it *after* firing `agent aborted` — so the event that
* announces the abort could not name it. Pulling the classification out lets
* each entry point decide the type first and hand it to the event, and keeps
* the one rule in one place.
* pi's `prompt()` resolves on the first turn without a tool call, which an
* agent mid-plan does emit, so the guard nudges it to carry on. A nudge that
* returns without a tool call and without assistant output is dead: the model
* has nothing more to give, and every further nudge returns just as empty, so
* the cap drains and the run aborts with no work done. Stop on the first dead
* nudge, and space live ones out.
*/
export function runErrorType(message: string): AgentErrorType {
const lower = message.toLowerCase();
if (lower.includes('rate limit') || lower.includes('429')) {
return AgentErrorType.RATE_LIMIT;
export async function nudgeWhileUnfinished(args: {
max: number;
/** True while the guard should keep nudging. */
unfinished: () => boolean;
/** Work the run has done: tool calls plus assistant output. */
progress: () => number;
send: (nudge: number) => Promise<void>;
backoffMs?: number;
}): Promise<NudgeRun> {
const backoffMs = args.backoffMs ?? NUDGE_BACKOFF_MS;
let nudges = 0;
while (nudges < args.max && args.unfinished()) {
if (nudges > 0) {
await new Promise((resolve) => setTimeout(resolve, backoffMs));
}
const before = args.progress();
nudges += 1;
await args.send(nudges);
if (args.progress() === before) return { nudges, dead: true };
}
return AgentErrorType.API_ERROR;
return { nudges, dead: false };
}
48 changes: 36 additions & 12 deletions src/lib/agent/runner/harness/pi/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,11 @@ import {
} from '@lib/constants';
import { analytics } from '@utils/analytics';
import { AgentErrorType } from '@lib/agent/agent-interface';
import { AgentSignals, REMARK_INSTRUCTION } from '@lib/agent/signals';
import {
AgentSignals,
REMARK_INSTRUCTION,
runErrorType,
} from '@lib/agent/signals';
import { AgentOutputSignals } from '@lib/agent/output-signals';
import { assembleCommandments } from '../../switchboard/commandments';
import { gatewayAuth, type GatewayAuth } from '@lib/gateway-session';
Expand All @@ -42,7 +46,11 @@ import type {
} from '../types';
import type { BootstrapResult } from '@lib/agent/runner/shared/types';
import type { TaskStore } from './tasks';
import { completionFailure, runErrorType } from './completion';
import {
completionFailure,
nudgeWhileUnfinished,
type NudgeRun,
} from './completion';

/** Injects the MCP server `instructions` pi-mcp-adapter drops (project env, skill steer, tool domains) into the system prompt, falling back to a bootstrap-derived project block when the warm-connect captured none. */
function piMcpContext(
Expand Down Expand Up @@ -219,6 +227,9 @@ export const piBackend: AgentHarness = {
// Tool calls across the whole run. Zero means the agent only ever produced
// text and never acted — a no-op that leaves the project untouched.
let toolCalls = 0;
// Assistant turns that carried text. A turn with neither text nor a tool
// call did nothing, and the completion guard reads that as a dead nudge.
let assistantOutputs = 0;
const runDurations = () => {
const durationMs = Date.now() - startTime;
return {
Expand All @@ -236,6 +247,8 @@ export const piBackend: AgentHarness = {
// Not `reason`: the linear sequence emits this same event with a `reason`
// holding the agent's free-text [ABORT] string, and one property cannot be
// both a closed enum and unbounded prose without making either unreadable.
// What the completion guard did, read by the abort event below.
let nudgeRun: NudgeRun = { nudges: 0, dead: false };
const captureAborted = (failureMode: AgentErrorType) =>
analytics.wizardCapture('agent aborted', {
failure_mode: failureMode,
Expand Down Expand Up @@ -507,6 +520,7 @@ export const piBackend: AgentHarness = {
turns.noteAssistantTurn(event.message);
const assistant = extractText(event.message).trim();
if (assistant) {
assistantOutputs += 1;
logToFile(`[pi] assistant: ${assistant.slice(0, 1000)}`);
applyOutroMarkers(assistant);
// Surface [STATUS] lines into the live spinner + status history,
Expand Down Expand Up @@ -566,17 +580,23 @@ export const piBackend: AgentHarness = {
// Completion guard: pi's prompt() resolves the moment the model returns
// a turn with no tool call (e.g. a lone [STATUS] line), even mid-plan.
// While tasks remain open and we're under the cap, nudge it to continue.
let continueNudges = 0;
while (
continueNudges < MAX_CONTINUE_NUDGES &&
!security.state.criticalViolation &&
hasOpenTasks(wizardTaskTools.store)
) {
continueNudges += 1;
nudgeRun = await nudgeWhileUnfinished({
max: MAX_CONTINUE_NUDGES,
unfinished: () =>
!security.state.criticalViolation &&
hasOpenTasks(wizardTaskTools.store),
progress: () => toolCalls + assistantOutputs,
send: (nudge) => {
logToFile(
`[pi] completion guard: tasks still open, nudge ${nudge}/${MAX_CONTINUE_NUDGES}`,
);
return turns.prompt(CONTINUE_INSTRUCTION);
},
});
if (nudgeRun.dead) {
logToFile(
`[pi] completion guard: tasks still open, nudge ${continueNudges}/${MAX_CONTINUE_NUDGES}`,
`[pi] completion guard: nudge ${nudgeRun.nudges} produced nothing; stopping`,
);
await turns.prompt(CONTINUE_INSTRUCTION);
}

// Best-effort remark ask — a failed turn never fails a successful run.
Expand Down Expand Up @@ -621,7 +641,11 @@ export const piBackend: AgentHarness = {
if (failure === AgentErrorType.INCOMPLETE_TASKS) {
spinner.stop('Agent stopped before finishing');
logToFile('[pi] incomplete: tasks left open');
analytics.wizardCapture('agent incomplete tasks', { open_tasks: true });
analytics.wizardCapture('agent incomplete tasks', {
open_tasks: true,
nudges: nudgeRun.nudges,
dead_nudge: nudgeRun.dead,
});
captureAborted(failure);
return { error: failure };
}
Expand Down
Loading
Loading