Skip to content

Commit 3f44694

Browse files
committed
Harden provider failure handling
1 parent 04ffbac commit 3f44694

11 files changed

Lines changed: 294 additions & 193 deletions

src/exec/runner.ts

Lines changed: 24 additions & 60 deletions
Original file line numberDiff line numberDiff line change
@@ -13,12 +13,7 @@ import { noopAuditStore, permissiveAuthorize } from "@intx/agent/testing";
1313
import { getLogger } from "@intx/log";
1414
import { createOptimizedContextStore } from "../session/optimized-context-store.js";
1515
import { type } from "arktype";
16-
import {
17-
buildCodexSource,
18-
buildOpenAISource,
19-
buildXaiSource,
20-
type Config,
21-
} from "../config/index.js";
16+
import { type Config } from "../config/index.js";
2217
import {
2318
loadLocalSettings,
2419
resolveLocalSettingsPath,
@@ -67,7 +62,10 @@ import { createAgentToolset, type AgentToolset, type OperatorResult } from "../a
6762
import { createAgentWithLiveToolDispatch } from "../agent/live-tool-dispatch.js";
6863
import { liveTelemetry } from "../telemetry/singleton.js";
6964
import { createTurnObserver } from "../telemetry/ai-observability.js";
70-
import { terminalProviderFailureMessage } from "../inference-error-message.js";
65+
import {
66+
isResolvedProviderFailureError,
67+
terminalProviderFailureMessage,
68+
} from "../inference-error-message.js";
7169
import { collectToolPlugins, resolveToolPlugins } from "../plugins/tool-plugins.js";
7270
import {
7371
expandExistingPluginMembers,
@@ -136,9 +134,13 @@ export async function refreshSelectedProviderCredential<T>(refresh: () => Promis
136134
export function execUserFailureMessage(
137135
config: Config,
138136
err: unknown,
139-
inferenceStarted: boolean,
137+
providerFailureObserved: boolean,
140138
): string {
141-
if (inferenceStarted || (err instanceof Error && err.name === SELECTED_PROVIDER_FAILURE)) {
139+
if (
140+
providerFailureObserved ||
141+
isResolvedProviderFailureError(err) ||
142+
(err instanceof Error && err.name === SELECTED_PROVIDER_FAILURE)
143+
) {
142144
return terminalProviderFailureMessage(
143145
config.providerName,
144146
config.settings?.providers[config.providerName]?.name,
@@ -324,7 +326,7 @@ export async function runExec(config: Config): Promise<ExecResult> {
324326
let finalized = false;
325327
let turnsUsed = 0;
326328
let runSink: RunSink | null = null;
327-
let inferenceStarted = false;
329+
let providerFailureObserved = false;
328330

329331
const persist = async (
330332
status: "running" | "done" | "failed" | "cancelled",
@@ -614,56 +616,14 @@ export async function runExec(config: Config): Promise<ExecResult> {
614616

615617
const initialCodexProfile = codexProfileFromProviderName(config.providerName);
616618
const initialXaiProfile = xaiProfileFromProviderName(config.providerName);
617-
const initialCodexAccountId = config.providers.find(
618-
(p) => p.name === config.providerName,
619-
)?.codexAccountId;
620-
621-
const buildOpenAICompatibleInitialSource = (): InferenceSource =>
622-
buildOpenAISource({
623-
id: config.providerName,
624-
baseURL: config.baseURL,
625-
apiKey: config.apiKey,
626-
model: config.model,
627-
...(config.reasoningEffort !== undefined
628-
? { reasoningEffort: config.reasoningEffort }
629-
: {}),
630-
});
631-
632-
const buildSessionSources = (): { sources: InferenceSource[]; defaultSource: string } =>
633-
buildSessionSourcesFromConfig(config, sessionId);
634-
635-
const initialBundle = buildSessionSources();
619+
const initialBundle = buildSessionSourcesFromConfig(config, sessionId);
636620
const liveSources = initialBundle.sources;
637621
const liveDefaultSource = initialBundle.defaultSource;
638-
639-
const buildInitialSourceFallback = (): InferenceSource =>
640-
initialCodexProfile !== undefined
641-
? buildCodexSource({
642-
id: config.providerName,
643-
apiKey: config.apiKey,
644-
model: config.model,
645-
sessionId,
646-
...(initialCodexAccountId !== undefined ? { accountId: initialCodexAccountId } : {}),
647-
...(config.reasoningEffort !== undefined
648-
? { reasoningEffort: config.reasoningEffort }
649-
: {}),
650-
})
651-
: initialXaiProfile !== undefined
652-
? buildXaiSource({
653-
id: config.providerName,
654-
apiKey: config.apiKey,
655-
model: config.model,
656-
sessionId,
657-
...(config.reasoningEffort !== undefined
658-
? { reasoningEffort: config.reasoningEffort }
659-
: {}),
660-
})
661-
: buildOpenAICompatibleInitialSource();
662-
663-
let liveSource: InferenceSource =
664-
liveSources.find((s) => s.id === liveDefaultSource) ??
665-
liveSources[0] ??
666-
buildInitialSourceFallback();
622+
const selectedSource = liveSources[0];
623+
if (selectedSource === undefined) {
624+
throw new Error("Selected inference source was not assembled");
625+
}
626+
let liveSource: InferenceSource = selectedSource;
667627

668628
// Refresh pinned Codex instructions before first inference, same as the
669629
// TUI path. Best-effort: a network failure falls back to the disk cache
@@ -794,6 +754,11 @@ export async function runExec(config: Config): Promise<ExecResult> {
794754
// its partial output in partial.jsonl instead of vanishing.
795755
const cycleRecorder = createCycleTextRecorder(() => workdir);
796756
const sink = (event: ReactorEmittedEvent): void => {
757+
if (event.type === "inference.start" || event.type === "inference.done") {
758+
providerFailureObserved = false;
759+
} else if (event.type === "inference.error") {
760+
providerFailureObserved = true;
761+
}
797762
liveSink.sink(event);
798763
cycleRecorder.handleEvent(event);
799764
if (event.type === "inference.text.delta") {
@@ -817,7 +782,6 @@ export async function runExec(config: Config): Promise<ExecResult> {
817782
let runError: string | undefined;
818783
let sinkStatus: ReturnType<typeof liveSink.getStatus> = "cancelled";
819784
try {
820-
inferenceStarted = true;
821785
// Final OAuth refresh immediately before send (token may have aged during MCP).
822786
if (initialCodexProfile !== undefined) {
823787
const { access } = await getValidCodexToken(initialCodexProfile);
@@ -958,7 +922,7 @@ export async function runExec(config: Config): Promise<ExecResult> {
958922
} catch (err) {
959923
const diagnosticMessage = formatCaughtError(err);
960924
logger.error("exec failed: {error}", { error: diagnosticMessage });
961-
const userMessage = execUserFailureMessage(config, err, inferenceStarted);
925+
const userMessage = execUserFailureMessage(config, err, providerFailureObserved);
962926
stderr.write(`Error: ${userMessage}\n`);
963927
await persist("failed", { error: diagnosticMessage });
964928
return {

src/inference-error-message.test.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,12 @@ describe("terminalProviderFailureMessage", () => {
106106
);
107107
});
108108

109+
test("does not duplicate Provider in configured display labels", () => {
110+
expect(terminalProviderFailureMessage("codex/work", "Codex Provider")).toBe(
111+
'Codex Provider failed. Try again or switch with "/model" and select another.',
112+
);
113+
});
114+
109115
test("uses a safe label when the provider id contains only control sequences", () => {
110116
const message = terminalProviderFailureMessage("\u001b[31m\u001b[0m");
111117
expect(message).toBe(

src/inference-error-message.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -95,12 +95,13 @@ function codexUsageLimitLine(error: InferenceErrorLike): string | undefined {
9595
export function terminalProviderFailureMessage(providerId: string, displayLabel?: string): string {
9696
const preferred = displayLabel?.trim() || providerId;
9797
const sanitized = stripTerminalControlSequences(preferred).replace(/\s+/g, " ").trim();
98-
const label = sanitized.length > 0 ? sanitized : "Unknown";
98+
const label = (sanitized.length > 0 ? sanitized : "Unknown").replace(/\s+Provider$/i, "");
9999
return `${label} Provider failed. Try again or switch with "/model" and select another.`;
100100
}
101101

102102
export type ResolvedProviderFailureError = Error & {
103103
readonly name: "ResolvedProviderFailureError";
104+
readonly providerId: string;
104105
readonly diagnosticMessage: string;
105106
};
106107

@@ -111,6 +112,7 @@ export function createResolvedProviderFailureError(
111112
): ResolvedProviderFailureError {
112113
return Object.assign(new Error(terminalProviderFailureMessage(providerId, displayLabel)), {
113114
name: "ResolvedProviderFailureError" as const,
115+
providerId,
114116
diagnosticMessage,
115117
});
116118
}
@@ -121,6 +123,8 @@ export function isResolvedProviderFailureError(
121123
return (
122124
error instanceof Error &&
123125
error.name === "ResolvedProviderFailureError" &&
126+
"providerId" in error &&
127+
typeof error.providerId === "string" &&
124128
"diagnosticMessage" in error &&
125129
typeof error.diagnosticMessage === "string"
126130
);

src/session/run-sink.test.ts

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -223,7 +223,7 @@ describe("createRunSink", () => {
223223
});
224224
});
225225

226-
test("does not leak authoritative source attribution across retry attempts", () => {
226+
test("retains selected provider attribution across a same-provider retry", () => {
227227
const { captured, runSink } = attributionHarness();
228228

229229
runSink.sink(event("inference.start", { model: "model-a" }));
@@ -234,13 +234,13 @@ describe("createRunSink", () => {
234234
}),
235235
);
236236
runSink.sink(event("inference.error", { error: { message: "retry" } }));
237-
runSink.sink(event("inference.start", { model: "model-b" }));
237+
runSink.sink(event("inference.start", { model: "model-a" }));
238238
failMessageRun(runSink);
239239

240240
expect(captured).toHaveLength(1);
241241
expect(captured[0]?.properties).toMatchObject({
242-
$ai_provider: "unknown",
243-
$ai_model: "model-b",
242+
$ai_provider: "provider-a",
243+
$ai_model: "model-a",
244244
$ai_is_error: true,
245245
});
246246
});

src/subagent/agent-fleet.ts

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,7 @@ interface FleetRecord {
9898
status: WaitJSONStatus;
9999
report?: string;
100100
error?: string;
101+
providerFailure?: true;
101102
/** Set once a wait_agents caller has been handed this result. */
102103
collected?: boolean;
103104
/** Set once the payload has been compacted away to bound memory. */
@@ -118,6 +119,7 @@ interface FleetOverlay {
118119
lastWaitStatus?: WaitJSONStatus;
119120
tombstoned?: boolean;
120121
hint?: string;
122+
providerFailure?: true;
121123
}
122124

123125
const RECOVERY_HINT =
@@ -170,6 +172,12 @@ class FleetMailbox {
170172
return record !== undefined && record.status !== "running" && record.collected !== true;
171173
}
172174

175+
markProviderFailure(id: string): void {
176+
const existing = this.records.get(id);
177+
if (existing === undefined) return;
178+
existing.providerFailure = true;
179+
}
180+
173181
/**
174182
* Overlay wait-status override so wait unblocks while the session may still
175183
* be running (send_input interrupt:true followup, close_agent teardown).
@@ -301,6 +309,7 @@ class FleetMailbox {
301309
...(overlay.hint !== undefined ? { hint: overlay.hint } : {}),
302310
...(payload?.report !== undefined ? { report: payload.report } : {}),
303311
...(payload?.error !== undefined && status === "failed" ? { error: payload.error } : {}),
312+
...(overlay.providerFailure === true ? { providerFailure: true } : {}),
304313
};
305314
}
306315

@@ -869,7 +878,13 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
869878
if (!childCtl.signal.aborted) childCtl.abort();
870879
});
871880

881+
let providerFailureObserved = false;
872882
const onEvent = (event: ReactorEmittedEvent): void => {
883+
if (event.type === "inference.start" || event.type === "inference.done") {
884+
providerFailureObserved = false;
885+
} else if (event.type === "inference.error") {
886+
providerFailureObserved = true;
887+
}
873888
deps.sessions.appendEvent(session.id, event);
874889
deps.onEvent?.(event);
875890
};
@@ -1084,10 +1099,14 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
10841099
: err instanceof Error
10851100
? err.message
10861101
: String(err);
1102+
const isProviderFailure = isResolvedProviderFailureError(err);
10871103
const authMessage = formatSubAgentSpawnAuthFailureMessage(description, err);
10881104
const failReason =
10891105
authMessage ??
1090-
(isResolvedProviderFailureError(err) ? err.message : diagnosticMessage);
1106+
(isProviderFailure ? err.message : diagnosticMessage);
1107+
if (isProviderFailure || providerFailureObserved) {
1108+
deps.fleetRecords.markProviderFailure(session.id);
1109+
}
10911110
deps.sessions.fail(session.id, failReason);
10921111
})
10931112
.finally(() => {
@@ -1232,6 +1251,7 @@ export function createWaitAgentsTool(deps: WaitAgentsDeps): AgentTool {
12321251
? { report: taken.report }
12331252
: {}),
12341253
...(taken.error !== undefined ? { error: taken.error } : {}),
1254+
...(taken.providerFailure === true ? { provider_failure: true } : {}),
12351255
...(taken.hint !== undefined ? { hint: taken.hint } : {}),
12361256
};
12371257
});
Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
import type { ReactorEmittedEvent } from "@intx/inference";
2+
3+
export interface PendingRunSettlement {
4+
readonly settled: Promise<void>;
5+
cancel: () => void;
6+
}
7+
8+
export interface RunEventSettlement {
9+
beginSend: () => PendingRunSettlement;
10+
handleEvent: (event: ReactorEmittedEvent) => void;
11+
endStream: () => void;
12+
}
13+
14+
/**
15+
* Coordinates Agent.send() with its streamed connector.reply. Agent.send resolves
16+
* when the connector reply is produced, which can precede consumption of earlier
17+
* inference.error events from the same run.
18+
*/
19+
export function createRunEventSettlement(): RunEventSettlement {
20+
const pending: { resolve: () => void }[] = [];
21+
let streamEnded = false;
22+
23+
return {
24+
beginSend: () => {
25+
const entry: { resolve: () => void } = {
26+
resolve: () => {
27+
throw new Error("Run settlement resolver was not initialized");
28+
},
29+
};
30+
const settled = new Promise<void>((resolve) => {
31+
entry.resolve = resolve;
32+
});
33+
if (streamEnded) entry.resolve();
34+
else pending.push(entry);
35+
return {
36+
settled,
37+
cancel: () => {
38+
const index = pending.indexOf(entry);
39+
if (index >= 0) pending.splice(index, 1);
40+
},
41+
};
42+
},
43+
handleEvent: (event) => {
44+
if (event.type !== "connector.reply") return;
45+
pending.shift()?.resolve();
46+
},
47+
endStream: () => {
48+
streamEnded = true;
49+
for (const entry of pending.splice(0)) entry.resolve();
50+
},
51+
};
52+
}

0 commit comments

Comments
 (0)