diff --git a/src/agents/embedded-agent-runner/thinking.test.ts b/src/agents/embedded-agent-runner/thinking.test.ts index 7f2f666d27c3..bb77481cc767 100644 --- a/src/agents/embedded-agent-runner/thinking.test.ts +++ b/src/agents/embedded-agent-runner/thinking.test.ts @@ -1021,6 +1021,64 @@ describe("wrapAnthropicStreamWithRecovery", () => { await expect(response.result()).resolves.toEqual(finalMessage); expect(events).toHaveLength(2); }); + + it("recovers an error event from a Promise-resolved stream without changing Promise timing", async () => { + const recovered = vi.fn(); + let callCount = 0; + let resolveFirstStream!: (stream: ReturnType) => void; + const firstStreamPromise = new Promise>( + (resolve) => { + resolveFirstStream = resolve; + }, + ); + const finalMessage = createTestAssistantMessage({ + content: [{ type: "text", text: "recovered answer" }], + stopReason: "stop", + }); + const wrapped = wrapAnthropicStreamWithRecovery( + (() => { + const attempt = ++callCount; + if (attempt === 1) { + return firstStreamPromise; + } + const stream = createAssistantMessageEventStream(); + queueMicrotask(() => { + stream.push({ type: "done", reason: "stop", message: finalMessage }); + stream.end(); + }); + return stream; + }) as Parameters[0], + { id: "test-session", onRecoveredAnthropicThinking: recovered }, + ); + + const responsePromise = wrapped({} as never, { messages: [] } as never, {} as never); + expect(responsePromise).toBeInstanceOf(Promise); + let resolved = false; + void Promise.resolve(responsePromise).then(() => { + resolved = true; + }); + await Promise.resolve(); + expect(resolved).toBe(false); + + const firstStream = createAssistantMessageEventStream(); + resolveFirstStream(firstStream); + const response = await responsePromise; + queueMicrotask(() => { + firstStream.push({ + type: "error", + reason: "error", + error: createTestStreamErrorMessage(terminalThinkingSignatureError), + }); + firstStream.end(); + }); + for await (const event of response) { + void event; + } + + await expect(response.result()).resolves.toEqual(finalMessage); + expect(callCount).toBe(2); + expect(recovered).toHaveBeenCalledTimes(1); + }); }); describe("stripStaleThinkingSignaturesForCompactionReplay", () => { diff --git a/src/agents/embedded-agent-runner/thinking.ts b/src/agents/embedded-agent-runner/thinking.ts index 07d5e3b072f4..99fd9a468156 100644 --- a/src/agents/embedded-agent-runner/thinking.ts +++ b/src/agents/embedded-agent-runner/thinking.ts @@ -698,6 +698,26 @@ async function pumpStreamWithRecovery( } } +function createRecoveryStream( + stream: Awaited>, + sessionMeta: RecoverySessionMeta, + retry: () => ReturnType, + notify: () => Promise, +): Awaited> { + const outer = createAssistantMessageEventStream(); + const finalResultPromise = pumpStreamWithRecovery( + outer, + stream, + sessionMeta, + retry, + notify, + ).finally(() => { + outer.end(); + }); + outer.result = () => finalResultPromise; + return outer; +} + export function wrapAnthropicStreamWithRecovery( innerStreamFn: StreamFn, sessionMeta: RecoverySessionMeta, @@ -727,28 +747,20 @@ export function wrapAnthropicStreamWithRecovery( const stream = innerStreamFn(model, context, options); if (stream instanceof Promise) { - return stream.catch((error: unknown) => { - if (!shouldRecoverAnthropicThinkingError(error, requestMeta)) { - throw error; - } - requestMeta.recoveredAnthropicThinking = true; - log.warn( - `[session-recovery] Anthropic thinking request rejected; retrying once without thinking blocks: sessionId=${requestMeta.id}`, - ); - return wrapRetryStreamWithRecoveryNotification(retry(), notify); - }) as ReturnType; + return stream.then( + (resolved) => createRecoveryStream(resolved, requestMeta, retry, notify), + (error: unknown) => { + if (!shouldRecoverAnthropicThinkingError(error, requestMeta)) { + throw error; + } + requestMeta.recoveredAnthropicThinking = true; + log.warn( + `[session-recovery] Anthropic thinking request rejected; retrying once without thinking blocks: sessionId=${requestMeta.id}`, + ); + return wrapRetryStreamWithRecoveryNotification(retry(), notify); + }, + ) as ReturnType; } - const outer = createAssistantMessageEventStream(); - const finalResultPromise = pumpStreamWithRecovery( - outer, - stream, - requestMeta, - retry, - notify, - ).finally(() => { - outer.end(); - }); - outer.result = () => finalResultPromise; - return outer as unknown as ReturnType; + return createRecoveryStream(stream, requestMeta, retry, notify); }; }