Skip to content

Commit 1ad4b55

Browse files
committed
Harden provider failure handling
1 parent dce3df6 commit 1ad4b55

13 files changed

Lines changed: 318 additions & 208 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,
@@ -63,7 +58,10 @@ import { createAgentToolset, type AgentToolset, type OperatorResult } from "../a
6358
import { createAgentWithLiveToolDispatch } from "../agent/live-tool-dispatch.js";
6459
import { liveTelemetry } from "../telemetry/singleton.js";
6560
import { createTurnObserver } from "../telemetry/ai-observability.js";
66-
import { terminalProviderFailureMessage } from "../inference-error-message.js";
61+
import {
62+
isResolvedProviderFailureError,
63+
terminalProviderFailureMessage,
64+
} from "../inference-error-message.js";
6765
import { collectToolPlugins, resolveToolPlugins } from "../plugins/tool-plugins.js";
6866
import {
6967
expandExistingPluginMembers,
@@ -132,9 +130,13 @@ export async function refreshSelectedProviderCredential<T>(refresh: () => Promis
132130
export function execUserFailureMessage(
133131
config: Config,
134132
err: unknown,
135-
inferenceStarted: boolean,
133+
providerFailureObserved: boolean,
136134
): string {
137-
if (inferenceStarted || (err instanceof Error && err.name === SELECTED_PROVIDER_FAILURE)) {
135+
if (
136+
providerFailureObserved ||
137+
isResolvedProviderFailureError(err) ||
138+
(err instanceof Error && err.name === SELECTED_PROVIDER_FAILURE)
139+
) {
138140
return terminalProviderFailureMessage(
139141
config.providerName,
140142
config.settings?.providers[config.providerName]?.name,
@@ -279,7 +281,7 @@ export async function runExec(config: Config): Promise<ExecResult> {
279281
let finalized = false;
280282
let turnsUsed = 0;
281283
let runSink: RunSink | null = null;
282-
let inferenceStarted = false;
284+
let providerFailureObserved = false;
283285

284286
const persist = async (
285287
status: "running" | "done" | "failed" | "cancelled",
@@ -568,56 +570,14 @@ export async function runExec(config: Config): Promise<ExecResult> {
568570

569571
const initialCodexProfile = codexProfileFromProviderName(config.providerName);
570572
const initialXaiProfile = xaiProfileFromProviderName(config.providerName);
571-
const initialCodexAccountId = config.providers.find(
572-
(p) => p.name === config.providerName,
573-
)?.codexAccountId;
574-
575-
const buildOpenAICompatibleInitialSource = (): InferenceSource =>
576-
buildOpenAISource({
577-
id: config.providerName,
578-
baseURL: config.baseURL,
579-
apiKey: config.apiKey,
580-
model: config.model,
581-
...(config.reasoningEffort !== undefined
582-
? { reasoningEffort: config.reasoningEffort }
583-
: {}),
584-
});
585-
586-
const buildSessionSources = (): { sources: InferenceSource[]; defaultSource: string } =>
587-
buildSessionSourcesFromConfig(config, sessionId);
588-
589-
const initialBundle = buildSessionSources();
573+
const initialBundle = buildSessionSourcesFromConfig(config, sessionId);
590574
const liveSources = initialBundle.sources;
591575
const liveDefaultSource = initialBundle.defaultSource;
592-
593-
const buildInitialSourceFallback = (): InferenceSource =>
594-
initialCodexProfile !== undefined
595-
? buildCodexSource({
596-
id: config.providerName,
597-
apiKey: config.apiKey,
598-
model: config.model,
599-
sessionId,
600-
...(initialCodexAccountId !== undefined ? { accountId: initialCodexAccountId } : {}),
601-
...(config.reasoningEffort !== undefined
602-
? { reasoningEffort: config.reasoningEffort }
603-
: {}),
604-
})
605-
: initialXaiProfile !== undefined
606-
? buildXaiSource({
607-
id: config.providerName,
608-
apiKey: config.apiKey,
609-
model: config.model,
610-
sessionId,
611-
...(config.reasoningEffort !== undefined
612-
? { reasoningEffort: config.reasoningEffort }
613-
: {}),
614-
})
615-
: buildOpenAICompatibleInitialSource();
616-
617-
let liveSource: InferenceSource =
618-
liveSources.find((s) => s.id === liveDefaultSource) ??
619-
liveSources[0] ??
620-
buildInitialSourceFallback();
576+
const selectedSource = liveSources[0];
577+
if (selectedSource === undefined) {
578+
throw new Error("Selected inference source was not assembled");
579+
}
580+
let liveSource: InferenceSource = selectedSource;
621581

622582
// Refresh pinned Codex instructions before first inference, same as the
623583
// TUI path. Best-effort: a network failure falls back to the disk cache
@@ -748,6 +708,11 @@ export async function runExec(config: Config): Promise<ExecResult> {
748708
// its partial output in partial.jsonl instead of vanishing.
749709
const cycleRecorder = createCycleTextRecorder(() => workdir);
750710
const sink = (event: ReactorEmittedEvent): void => {
711+
if (event.type === "inference.start" || event.type === "inference.done") {
712+
providerFailureObserved = false;
713+
} else if (event.type === "inference.error") {
714+
providerFailureObserved = true;
715+
}
751716
liveSink.sink(event);
752717
cycleRecorder.handleEvent(event);
753718
if (event.type === "inference.text.delta") {
@@ -771,7 +736,6 @@ export async function runExec(config: Config): Promise<ExecResult> {
771736
let runError: string | undefined;
772737
let sinkStatus: ReturnType<typeof liveSink.getStatus> = "cancelled";
773738
try {
774-
inferenceStarted = true;
775739
// Final OAuth refresh immediately before send (token may have aged during MCP).
776740
if (initialCodexProfile !== undefined) {
777741
const { access } = await getValidCodexToken(initialCodexProfile);
@@ -912,7 +876,7 @@ export async function runExec(config: Config): Promise<ExecResult> {
912876
} catch (err) {
913877
const diagnosticMessage = formatCaughtError(err);
914878
logger.error("exec failed: {error}", { error: diagnosticMessage });
915-
const userMessage = execUserFailureMessage(config, err, inferenceStarted);
879+
const userMessage = execUserFailureMessage(config, err, providerFailureObserved);
916880
stderr.write(`Error: ${userMessage}\n`);
917881
await persist("failed", { error: diagnosticMessage });
918882
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: 26 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,7 @@ interface FleetRecord {
9494
status: "running" | "done" | "failed" | "interrupted";
9595
report?: string;
9696
error?: string;
97+
providerFailure?: true;
9798
/** Set once a wait_agents caller has been handed this result. */
9899
collected?: boolean;
99100
/** Set once the payload has been compacted away to bound memory. */
@@ -141,10 +142,14 @@ class FleetRecords {
141142
this.notify();
142143
}
143144

144-
reject(id: string, error: string): void {
145+
reject(id: string, error: string, providerFailure = false): void {
145146
const existing = this.records.get(id);
146147
if (existing !== undefined && existing.status !== "running") return;
147-
this.records.set(id, { status: "failed", error });
148+
this.records.set(id, {
149+
status: "failed",
150+
error,
151+
...(providerFailure ? { providerFailure: true } : {}),
152+
});
148153
this.enforceCap();
149154
this.notify();
150155
}
@@ -586,7 +591,13 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
586591
if (!childCtl.signal.aborted) childCtl.abort();
587592
});
588593

594+
let providerFailureObserved = false;
589595
const onEvent = (event: ReactorEmittedEvent): void => {
596+
if (event.type === "inference.start" || event.type === "inference.done") {
597+
providerFailureObserved = false;
598+
} else if (event.type === "inference.error") {
599+
providerFailureObserved = true;
600+
}
590601
deps.sessions.appendEvent(session.id, event);
591602
deps.onEvent?.(event);
592603
};
@@ -786,13 +797,18 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
786797
: err instanceof Error
787798
? err.message
788799
: String(err);
789-
const parentMessage = isResolvedProviderFailureError(err)
790-
? err.message
791-
: classifySubAgentInferenceAuthFailure(err) !== null
792-
? terminalProviderFailureMessage(provider.providerName)
793-
: diagnosticMessage;
794-
deps.fleetRecords.reject(session.id, parentMessage);
795-
deps.sessions.fail(session.id, diagnosticMessage);
800+
const providerId = isResolvedProviderFailureError(err)
801+
? err.providerId
802+
: provider.providerName;
803+
const isProviderFailure =
804+
providerFailureObserved ||
805+
isResolvedProviderFailureError(err) ||
806+
classifySubAgentInferenceAuthFailure(err) !== null;
807+
const parentMessage = isProviderFailure
808+
? terminalProviderFailureMessage(providerId, settings?.providers[providerId]?.name)
809+
: diagnosticMessage;
810+
deps.fleetRecords.reject(session.id, parentMessage, isProviderFailure);
811+
deps.sessions.fail(session.id, parentMessage);
796812
})
797813
.finally(() => {
798814
finalizeEnd();
@@ -919,6 +935,7 @@ export function createWaitAgentsTool(deps: WaitAgentsDeps): AgentTool {
919935
status: taken.status,
920936
...(taken.report !== undefined ? { report: taken.report } : {}),
921937
...(taken.error !== undefined ? { error: taken.error } : {}),
938+
...(taken.providerFailure === true ? { provider_failure: true } : {}),
922939
...(taken.hint !== undefined ? { hint: taken.hint } : {}),
923940
};
924941
}

src/subagent/index.test.ts

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1160,7 +1160,9 @@ describe("createTaskTool", () => {
11601160
expect(String(result.content)).not.toContain("401 Unauthorized");
11611161
const row = sessions.list().find((s) => s.description === "auth probe");
11621162
expect(row?.status).toBe("failed");
1163-
expect(row?.error).toBe("401 Unauthorized");
1163+
expect(row?.error).toBe(
1164+
'test-provider Provider failed. Try again or switch with "/model" and select another.',
1165+
);
11641166
});
11651167

11661168
test("terminal child provider failure hides raw diagnostics from the parent result", async () => {
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)