Skip to content
This repository was archived by the owner on Aug 6, 2026. It is now read-only.

Commit cac7efd

Browse files
authored
fix(agent): continue unattended turns after transient upstream stream failures
Initial and resume turns run unattended in cloud sandboxes — there is no user watching who could retry — yet handleTurnFailure only treats follow-up turns in interactive mode as recoverable. A single transient upstream stream death during the initial or resume turn therefore failed the whole task run and tore down the sandbox. Give those turns a bounded (2 per turn) number of automatic retries when the failure classifies as a transient upstream failure. A stream that died mid-response already delivered the original prompt into the session history, so that case retries with a hidden continuation prompt; failures where the request may never have been processed re-send the original prompt instead. The session is re-read on every attempt so a teardown or re-initialization during the retry delay cannot route the retry to a stale session. Non-upstream errors and exhausted budgets fail exactly as before.
1 parent 59f538b commit cac7efd

3 files changed

Lines changed: 239 additions & 31 deletions

File tree

‎packages/agent/src/server/agent-server.test.ts‎

Lines changed: 117 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -956,6 +956,123 @@ describe("AgentServer HTTP Mode", () => {
956956
},
957957
);
958958

959+
function createRetryTestServer(prompt: ReturnType<typeof vi.fn>) {
960+
const testServer = createFailureTestServer();
961+
testServer.session = {
962+
acpSessionId: "acp-1",
963+
payload: { run_id: "run-1" },
964+
logWriter: { appendRawLine: vi.fn(), flush: vi.fn(async () => {}) },
965+
clientConnection: { prompt },
966+
};
967+
return testServer as unknown as {
968+
promptWithUpstreamRetry(request: {
969+
sessionId: string;
970+
prompt: ContentBlock[];
971+
}): Promise<{ stopReason: string }>;
972+
};
973+
}
974+
975+
it("continues an unattended turn after a transient upstream stream death", async () => {
976+
vi.useFakeTimers();
977+
try {
978+
const prompt = vi
979+
.fn()
980+
.mockRejectedValueOnce(new Error("API Error: terminated"))
981+
.mockResolvedValueOnce({ stopReason: "end_turn" });
982+
const testServer = createRetryTestServer(prompt);
983+
984+
const resultPromise = testServer.promptWithUpstreamRetry({
985+
sessionId: "acp-1",
986+
prompt: [{ type: "text", text: "do the task" }],
987+
});
988+
await vi.advanceTimersByTimeAsync(5_000);
989+
990+
await expect(resultPromise).resolves.toEqual({
991+
stopReason: "end_turn",
992+
});
993+
expect(prompt).toHaveBeenCalledTimes(2);
994+
const retryRequest = prompt.mock.calls[1][0] as {
995+
sessionId: string;
996+
prompt: Array<{ type: string; text: string }>;
997+
};
998+
expect(retryRequest.sessionId).toBe("acp-1");
999+
expect(retryRequest.prompt[0].text).toContain(
1000+
"interrupted by a transient connection error",
1001+
);
1002+
} finally {
1003+
vi.useRealTimers();
1004+
}
1005+
});
1006+
1007+
it("re-sends the original prompt when the failure happened before the stream started", async () => {
1008+
vi.useFakeTimers();
1009+
try {
1010+
const prompt = vi
1011+
.fn()
1012+
.mockRejectedValueOnce(new Error("API Error: Connection error."))
1013+
.mockResolvedValueOnce({ stopReason: "end_turn" });
1014+
const testServer = createRetryTestServer(prompt);
1015+
1016+
const resultPromise = testServer.promptWithUpstreamRetry({
1017+
sessionId: "acp-1",
1018+
prompt: [{ type: "text", text: "do the task" }],
1019+
});
1020+
await vi.advanceTimersByTimeAsync(5_000);
1021+
1022+
await expect(resultPromise).resolves.toEqual({
1023+
stopReason: "end_turn",
1024+
});
1025+
expect(prompt).toHaveBeenCalledTimes(2);
1026+
const retryRequest = prompt.mock.calls[1][0] as {
1027+
sessionId: string;
1028+
prompt: Array<{ type: string; text: string }>;
1029+
};
1030+
expect(retryRequest.sessionId).toBe("acp-1");
1031+
expect(retryRequest.prompt).toEqual([
1032+
{ type: "text", text: "do the task" },
1033+
]);
1034+
} finally {
1035+
vi.useRealTimers();
1036+
}
1037+
});
1038+
1039+
it("does not retry a genuine agent error", async () => {
1040+
const prompt = vi.fn().mockRejectedValue(new Error("boom"));
1041+
const testServer = createRetryTestServer(prompt);
1042+
1043+
await expect(
1044+
testServer.promptWithUpstreamRetry({
1045+
sessionId: "acp-1",
1046+
prompt: [{ type: "text", text: "do the task" }],
1047+
}),
1048+
).rejects.toThrow("boom");
1049+
expect(prompt).toHaveBeenCalledTimes(1);
1050+
});
1051+
1052+
it("stops continuing once the bounded retry budget is exhausted", async () => {
1053+
vi.useFakeTimers();
1054+
try {
1055+
const prompt = vi
1056+
.fn()
1057+
.mockRejectedValue(new Error("API Error: terminated"));
1058+
const testServer = createRetryTestServer(prompt);
1059+
1060+
const resultPromise = testServer.promptWithUpstreamRetry({
1061+
sessionId: "acp-1",
1062+
prompt: [{ type: "text", text: "do the task" }],
1063+
});
1064+
const assertion = expect(resultPromise).rejects.toThrow(
1065+
"API Error: terminated",
1066+
);
1067+
await vi.advanceTimersByTimeAsync(10_000);
1068+
1069+
await assertion;
1070+
expect(prompt).toHaveBeenCalledTimes(3);
1071+
} finally {
1072+
vi.useRealTimers();
1073+
}
1074+
});
1075+
9591076
it("persists structured turn completion notifications", () => {
9601077
const appendRawLine = vi.fn();
9611078
const testServer = new AgentServer({

‎packages/agent/src/server/agent-server.ts‎

Lines changed: 74 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,12 @@ type MessageCallback = (message: unknown) => void;
125125

126126
export const SSE_KEEPALIVE_INTERVAL_MS = 25_000;
127127

128+
// Bounded per-turn retries for unattended (initial/resume) turns that hit a
129+
// transient upstream failure. Two covers a retry whose own attempt also gets
130+
// cut once, without letting a hard upstream outage loop forever.
131+
const MAX_UPSTREAM_TURN_RETRIES = 2;
132+
const UPSTREAM_TURN_RETRY_DELAY_MS = 5_000;
133+
128134
class NdJsonTap {
129135
private decoder = new TextDecoder();
130136
private buffer = "";
@@ -1516,6 +1522,72 @@ export class AgentServer {
15161522
return { classification: classifyAgentError(message), message };
15171523
}
15181524

1525+
/**
1526+
* Send an initial/resume turn prompt, absorbing transient upstream
1527+
* failures with a bounded number of retries. These turns run unattended
1528+
* (no user watching who could retry), so without this a single transient
1529+
* transport cut fails the whole run. A stream that died mid-response has
1530+
* already delivered the original prompt into the session history, so that
1531+
* case retries with a hidden continuation; failures where the request may
1532+
* never have been processed re-send the original prompt instead.
1533+
*/
1534+
private async promptWithUpstreamRetry(request: {
1535+
sessionId: string;
1536+
prompt: ContentBlock[];
1537+
_meta?: Record<string, unknown>;
1538+
}): Promise<PromptResponse> {
1539+
let retries = 0;
1540+
let continueInterruptedTurn = false;
1541+
for (;;) {
1542+
// Re-read the session on every attempt: it can be torn down or
1543+
// replaced while the retry delay is pending.
1544+
const session = this.session;
1545+
if (!session) {
1546+
throw new Error("Agent session ended before the turn could be sent");
1547+
}
1548+
const attempt = continueInterruptedTurn
1549+
? {
1550+
sessionId: session.acpSessionId,
1551+
prompt: [
1552+
hiddenTextBlock(
1553+
"The previous response was interrupted by a transient connection error. " +
1554+
"Continue from where you left off — do not repeat work that already completed.",
1555+
),
1556+
],
1557+
}
1558+
: { ...request, sessionId: session.acpSessionId };
1559+
try {
1560+
return await session.clientConnection.prompt(attempt);
1561+
} catch (error) {
1562+
const { classification, message } =
1563+
this.extractErrorClassification(error);
1564+
if (
1565+
!upstreamProviderFailureClassifications.has(classification) ||
1566+
retries >= MAX_UPSTREAM_TURN_RETRIES
1567+
) {
1568+
throw error;
1569+
}
1570+
retries += 1;
1571+
// Only a mid-response stream death guarantees the prompt reached the
1572+
// model; connection/timeout/status failures re-send the original.
1573+
continueInterruptedTurn =
1574+
classification === "upstream_stream_terminated";
1575+
this.logger.warn(
1576+
"Turn hit a transient upstream failure; retrying after a short delay",
1577+
{
1578+
classification,
1579+
message,
1580+
attempt: retries,
1581+
continueInterruptedTurn,
1582+
},
1583+
);
1584+
await new Promise((resolve) =>
1585+
setTimeout(resolve, UPSTREAM_TURN_RETRY_DELAY_MS),
1586+
);
1587+
}
1588+
}
1589+
}
1590+
15191591
private async handleTurnFailure(
15201592
payload: JwtPayload,
15211593
phase: "initial" | "resume" | "followup",
@@ -1665,7 +1737,7 @@ export class AgentServer {
16651737

16661738
this.session.logWriter.resetTurnMessages(payload.run_id);
16671739

1668-
const result = await this.session.clientConnection.prompt({
1740+
const result = await this.promptWithUpstreamRetry({
16691741
sessionId: this.session.acpSessionId,
16701742
prompt: initialPrompt,
16711743
...(initialPromptMeta ? { _meta: initialPromptMeta } : {}),
@@ -1877,7 +1949,7 @@ export class AgentServer {
18771949

18781950
this.session.logWriter.resetTurnMessages(payload.run_id);
18791951

1880-
const result = await this.session.clientConnection.prompt({
1952+
const result = await this.promptWithUpstreamRetry({
18811953
sessionId: this.session.acpSessionId,
18821954
prompt: builtPrompt.prompt,
18831955
...(builtPrompt.meta ? { _meta: builtPrompt.meta } : {}),

‎packages/agent/src/server/question-relay.test.ts‎

Lines changed: 48 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -730,7 +730,7 @@ describe("Question relay", () => {
730730
expect(promptSpy).not.toHaveBeenCalled();
731731
});
732732

733-
it("does not replay a transient upstream termination before any session activity", async () => {
733+
it("replays a transient upstream termination with a continuation prompt", async () => {
734734
vi.spyOn(server.posthogAPI, "getTask").mockResolvedValue({
735735
id: "test-task-id",
736736
title: "t",
@@ -744,7 +744,8 @@ describe("Question relay", () => {
744744

745745
const promptSpy = vi
746746
.fn()
747-
.mockRejectedValueOnce(createTransientPromptError());
747+
.mockRejectedValueOnce(createTransientPromptError())
748+
.mockResolvedValueOnce({ stopReason: "cancelled" });
748749
const updateTaskRunSpy = vi
749750
.spyOn(server.posthogAPI, "updateTaskRun")
750751
.mockResolvedValue({} as TaskRun);
@@ -762,20 +763,26 @@ describe("Question relay", () => {
762763
},
763764
};
764765

765-
await server.sendInitialTaskMessage(TEST_PAYLOAD);
766-
767-
expect(promptSpy).toHaveBeenCalledTimes(1);
768-
expect(updateTaskRunSpy).toHaveBeenCalledWith(
769-
"test-task-id",
770-
"test-run-id",
771-
{
772-
status: "failed",
773-
error_message: UPSTREAM_PROVIDER_FAILURE_MESSAGE,
774-
},
766+
vi.useFakeTimers();
767+
try {
768+
const sendPromise = server.sendInitialTaskMessage(TEST_PAYLOAD);
769+
await vi.advanceTimersByTimeAsync(5_000);
770+
await sendPromise;
771+
} finally {
772+
vi.useRealTimers();
773+
}
774+
775+
expect(promptSpy).toHaveBeenCalledTimes(2);
776+
const continuation = promptSpy.mock.calls[1][0] as {
777+
prompt: Array<{ type: string; text: string }>;
778+
};
779+
expect(continuation.prompt[0].text).toContain(
780+
"interrupted by a transient connection error",
775781
);
782+
expect(updateTaskRunSpy).not.toHaveBeenCalled();
776783
});
777784

778-
it("surfaces upstream provider failures with a retryable message", async () => {
785+
it("re-sends the original prompt after an upstream provider failure", async () => {
779786
vi.spyOn(server.posthogAPI, "getTask").mockResolvedValue({
780787
id: "test-task-id",
781788
title: "t",
@@ -789,7 +796,8 @@ describe("Question relay", () => {
789796

790797
const promptSpy = vi
791798
.fn()
792-
.mockRejectedValueOnce(createUpstreamProviderFailureError());
799+
.mockRejectedValueOnce(createUpstreamProviderFailureError())
800+
.mockResolvedValueOnce({ stopReason: "cancelled" });
793801
const updateTaskRunSpy = vi
794802
.spyOn(server.posthogAPI, "updateTaskRun")
795803
.mockResolvedValue({} as TaskRun);
@@ -807,20 +815,24 @@ describe("Question relay", () => {
807815
},
808816
};
809817

810-
await server.sendInitialTaskMessage(TEST_PAYLOAD);
811-
812-
expect(promptSpy).toHaveBeenCalledTimes(1);
813-
expect(updateTaskRunSpy).toHaveBeenCalledWith(
814-
"test-task-id",
815-
"test-run-id",
816-
{
817-
status: "failed",
818-
error_message: UPSTREAM_PROVIDER_FAILURE_MESSAGE,
819-
},
820-
);
818+
vi.useFakeTimers();
819+
try {
820+
const sendPromise = server.sendInitialTaskMessage(TEST_PAYLOAD);
821+
await vi.advanceTimersByTimeAsync(5_000);
822+
await sendPromise;
823+
} finally {
824+
vi.useRealTimers();
825+
}
826+
827+
expect(promptSpy).toHaveBeenCalledTimes(2);
828+
const retryRequest = promptSpy.mock.calls[1][0] as {
829+
prompt: Array<{ type: string; text: string }>;
830+
};
831+
expect(retryRequest.prompt[0].text).toBe("original task description");
832+
expect(updateTaskRunSpy).not.toHaveBeenCalled();
821833
});
822834

823-
it("surfaces upstream connection errors with the shared provider failure message", async () => {
835+
it("surfaces the shared provider failure message once upstream retries are exhausted", async () => {
824836
vi.spyOn(server.posthogAPI, "getTask").mockResolvedValue({
825837
id: "test-task-id",
826838
title: "t",
@@ -832,7 +844,7 @@ describe("Question relay", () => {
832844
state: {},
833845
} as unknown as TaskRun);
834846

835-
const promptSpy = vi.fn().mockImplementationOnce(async () => {
847+
const promptSpy = vi.fn().mockImplementation(async () => {
836848
throw createTransientConnectionError();
837849
});
838850
const updateTaskRunSpy = vi
@@ -852,9 +864,16 @@ describe("Question relay", () => {
852864
},
853865
};
854866

855-
await server.sendInitialTaskMessage(TEST_PAYLOAD);
867+
vi.useFakeTimers();
868+
try {
869+
const sendPromise = server.sendInitialTaskMessage(TEST_PAYLOAD);
870+
await vi.advanceTimersByTimeAsync(10_000);
871+
await sendPromise;
872+
} finally {
873+
vi.useRealTimers();
874+
}
856875

857-
expect(promptSpy).toHaveBeenCalledTimes(1);
876+
expect(promptSpy).toHaveBeenCalledTimes(3);
858877
expect(updateTaskRunSpy).toHaveBeenCalledWith(
859878
"test-task-id",
860879
"test-run-id",

0 commit comments

Comments
 (0)