From b457024d67000a4d00bbae8fc8f42faab5698947 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Wed, 19 Aug 2026 23:00:44 +0800 Subject: [PATCH 01/12] feat(api): abort signal support for openai-codex (completePrompt + createMessage) - completePrompt: use a request-local signal built with mergeAbortSignalAndTimeout(options?.abortSignal, options?.timeoutMs) for the fetch call instead of the handler-wide AbortController; re-throw abort errors as-is so cancellation is detectable by the "AbortError" name - createMessage: pass metadata into executeRequest and bridge metadata?.abortSignal into the internal AbortController (Bedrock pattern: pre-aborted guard + { once: true } listener), covering both the OpenAI SDK streaming path and the manual SSE fetch fallback - specs: port the reference completePrompt coverage (request body, timeoutMs=0, abortSignal/timeoutMs merging, error paths) and add pre-aborted and in-flight abort tests rejecting with name === "AbortError"; port the createMessage abort bridge + pre-aborted tests into the native tool calls spec --- .../openai-codex-native-tool-calls.spec.ts | 89 +++++ .../providers/__tests__/openai-codex.spec.ts | 364 ++++++++++++++++++ src/api/providers/openai-codex.ts | 26 +- 3 files changed, 474 insertions(+), 5 deletions(-) diff --git a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts index d9fcdcb967..b1ded8bbf3 100644 --- a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts +++ b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts @@ -518,4 +518,93 @@ describe("OpenAiCodexHandler native tool calls", () => { }), ) }) + + describe("createMessage abort signal", () => { + it("should bridge the external abortSignal into the internal AbortController", async () => { + vi.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vi.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + const mockCreate = vi.fn().mockResolvedValue({ + async *[Symbol.asyncIterator]() { + yield { type: "response.text.delta", delta: "test" } + yield { + type: "response.completed", + response: { + id: "resp_1", + status: "completed", + output: [{ type: "message", content: [{ type: "output_text", text: "test" }] }], + usage: { input_tokens: 1, output_tokens: 1 }, + }, + } + }, + }) + Object.assign(handler, { + client: { + responses: { create: mockCreate }, + }, + }) + + const controller = new AbortController() + const stream = handler.createMessage("system", [{ role: "user", content: "hello" }], { + taskId: "t", + abortSignal: controller.signal, + }) + + // Consume the stream to trigger the request + await collectStream(stream) + + expect(mockCreate).toHaveBeenCalled() + const createCallArgs = mockCreate.mock.calls[0][1] as { signal?: AbortSignal } + expect(createCallArgs.signal).toBeDefined() + expect(createCallArgs.signal).toBeInstanceOf(AbortSignal) + + // Verify the signal is not aborted before we abort the external one + expect(createCallArgs.signal?.aborted).toBe(false) + + // Abort the external signal; the bridge aborts the internal controller + controller.abort() + expect(controller.signal.aborted).toBe(true) + }) + + it("should immediately abort when the external signal is already aborted", async () => { + vi.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vi.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + const mockCreate = vi.fn().mockResolvedValue({ + async *[Symbol.asyncIterator]() { + yield { type: "response.text.delta", delta: "test" } + yield { + type: "response.completed", + response: { + id: "resp_1", + status: "completed", + output: [{ type: "message", content: [{ type: "output_text", text: "test" }] }], + usage: { input_tokens: 1, output_tokens: 1 }, + }, + } + }, + }) + Object.assign(handler, { + client: { + responses: { create: mockCreate }, + }, + }) + + const controller = new AbortController() + controller.abort() // Pre-abort + + const stream = handler.createMessage("system", [{ role: "user", content: "hello" }], { + taskId: "t", + abortSignal: controller.signal, + }) + + // Consume the stream to trigger the request + await collectStream(stream) + + expect(mockCreate).toHaveBeenCalled() + const createCallArgs = mockCreate.mock.calls[0][1] as { signal?: AbortSignal } + // The internal signal should already be aborted since the external one was pre-aborted + expect(createCallArgs.signal?.aborted).toBe(true) + }) + }) }) diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index 9a256535c1..01ff123f0a 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -641,3 +641,367 @@ describe("OpenAiCodexHandler Luna Responses Lite requests", () => { }) }) }) + +describe("OpenAiCodexHandler.completePrompt", () => { + afterEach(() => { + vitest.restoreAllMocks() + vitest.unstubAllGlobals() + }) + + it("should call fetch with correct request body and return text response", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + const mockFetch = vitest.fn().mockResolvedValue({ + ok: true, + json: () => + Promise.resolve({ + output: [ + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "Hello world" }], + }, + ], + }), + text: () => Promise.resolve(""), + }) + vitest.stubGlobal("fetch", mockFetch) + + const result = await handler.completePrompt("test prompt") + + expect(result).toBe("Hello world") + expect(mockFetch).toHaveBeenCalledWith( + expect.stringContaining("/responses"), + expect.objectContaining({ + method: "POST", + headers: expect.objectContaining({ + Authorization: "Bearer test-token", + originator: "zoo-code", + }), + body: expect.stringContaining('"model":"gpt-5.1-codex"'), + }), + ) + }) + + it("should treat timeoutMs=0 as no timeout", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + let fetchInit: RequestInit | undefined + const mockFetch = vitest.fn().mockImplementation(async (_url: string, init?: RequestInit) => { + fetchInit = init + return { + ok: true, + json: () => + Promise.resolve({ + output: [ + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "response" }], + }, + ], + }), + text: () => Promise.resolve(""), + } + }) + vitest.stubGlobal("fetch", mockFetch) + + await handler.completePrompt("test prompt", { timeoutMs: 0 }) + + expect(mockFetch).toHaveBeenCalled() + expect(fetchInit?.signal).toBeDefined() + expect(fetchInit?.signal?.aborted).toBe(false) + }) + + it("should merge abortSignal with local controller", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + let fetchInit: RequestInit | undefined + const mockFetch = vitest.fn().mockImplementation(async (_url: string, init?: RequestInit) => { + fetchInit = init + return { + ok: true, + json: () => + Promise.resolve({ + output: [ + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "response" }], + }, + ], + }), + text: () => Promise.resolve(""), + } + }) + vitest.stubGlobal("fetch", mockFetch) + + const controller = new AbortController() + const promise = handler.completePrompt("test prompt", { abortSignal: controller.signal }) + controller.abort() + await promise + + expect(mockFetch).toHaveBeenCalled() + expect(fetchInit?.signal).toBeDefined() + expect(fetchInit?.signal?.aborted).toBe(true) + }) + + it("should merge abortSignal and timeoutMs together", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + let fetchInit: RequestInit | undefined + const mockFetch = vitest.fn().mockImplementation(async (_url: string, init?: RequestInit) => { + fetchInit = init + return { + ok: true, + json: () => + Promise.resolve({ + output: [ + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "response" }], + }, + ], + }), + text: () => Promise.resolve(""), + } + }) + vitest.stubGlobal("fetch", mockFetch) + + const controller = new AbortController() + const promise = handler.completePrompt("test prompt", { abortSignal: controller.signal, timeoutMs: 5000 }) + controller.abort() + await promise + + expect(mockFetch).toHaveBeenCalled() + expect(fetchInit?.signal).toBeDefined() + expect(fetchInit?.signal?.aborted).toBe(true) + }) + + it("should reject with an AbortError when the signal is already aborted", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + // Emulate fetch semantics: an already-aborted signal rejects immediately + const mockFetch = vitest.fn().mockImplementation((_url: string, init?: RequestInit) => { + if (init?.signal?.aborted) { + const abortError = new Error("The operation was aborted") + abortError.name = "AbortError" + return Promise.reject(abortError) + } + return Promise.resolve({ + ok: true, + json: () => Promise.resolve({}), + text: () => Promise.resolve(""), + }) + }) + vitest.stubGlobal("fetch", mockFetch) + + const controller = new AbortController() + controller.abort() + + await expect(handler.completePrompt("test prompt", { abortSignal: controller.signal })).rejects.toSatisfy( + (error: unknown) => error instanceof Error && error.name === "AbortError", + ) + }) + + it("should abort an in-flight request when the external signal aborts", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + // Emulate fetch semantics: the pending request rejects when the signal aborts + const mockFetch = vitest.fn().mockImplementation((_url: string, init?: RequestInit) => { + return new Promise((_resolve, reject) => { + init?.signal?.addEventListener( + "abort", + () => { + const abortError = new Error("The operation was aborted") + abortError.name = "AbortError" + reject(abortError) + }, + { once: true }, + ) + }) + }) + vitest.stubGlobal("fetch", mockFetch) + + const controller = new AbortController() + const promise = handler.completePrompt("test prompt", { abortSignal: controller.signal }) + await vitest.waitFor(() => expect(mockFetch).toHaveBeenCalled()) + controller.abort() + + await expect(promise).rejects.toSatisfy( + (error: unknown) => error instanceof Error && error.name === "AbortError", + ) + }) + + it("should return empty string when no output text found", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + const mockFetch = vitest.fn().mockResolvedValue({ + ok: true, + json: () => Promise.resolve({ output: [{ type: "message", role: "assistant", content: [] }] }), + text: () => Promise.resolve(""), + }) + vitest.stubGlobal("fetch", mockFetch) + + const result = await handler.completePrompt("test prompt") + + expect(result).toBe("") + }) + + it("should handle responseData.text fallback", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + const mockFetch = vitest.fn().mockResolvedValue({ + ok: true, + json: () => Promise.resolve({ text: "fallback response" }), + text: () => Promise.resolve(""), + }) + vitest.stubGlobal("fetch", mockFetch) + + const result = await handler.completePrompt("test prompt") + + expect(result).toBe("fallback response") + }) + + it("should throw error when not authenticated", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue(null) + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + const controller = new AbortController() + await expect(handler.completePrompt("test prompt", { abortSignal: controller.signal })).rejects.toThrow() + }) + + it("should throw error when fetch returns non-ok response", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + const mockFetch = vitest.fn().mockResolvedValue({ + ok: false, + status: 401, + text: () => Promise.resolve("Unauthorized"), + }) + vitest.stubGlobal("fetch", mockFetch) + + const controller = new AbortController() + await expect(handler.completePrompt("test prompt", { abortSignal: controller.signal })).rejects.toThrow() + }) + + it("should include reasoning config when model has reasoning effort", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + const mockFetch = vitest.fn().mockResolvedValue({ + ok: true, + json: () => + Promise.resolve({ + output: [ + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "response" }], + }, + ], + }), + text: () => Promise.resolve(""), + }) + vitest.stubGlobal("fetch", mockFetch) + + await handler.completePrompt("test prompt") + + expect(mockFetch).toHaveBeenCalled() + const fetchOptions = mockFetch.mock.calls[0][1] + const requestBody = JSON.parse(fetchOptions.body) + expect(requestBody.include).toContain("reasoning.encrypted_content") + expect(requestBody.reasoning).toBeDefined() + expect(requestBody.reasoning.effort).toBe("medium") + }) + + it("should include ChatGPT-Account-Id header when accountId is available", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_12345") + + const mockFetch = vitest.fn().mockResolvedValue({ + ok: true, + json: () => + Promise.resolve({ + output: [ + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "response" }], + }, + ], + }), + text: () => Promise.resolve(""), + }) + vitest.stubGlobal("fetch", mockFetch) + + await handler.completePrompt("test prompt") + + expect(mockFetch).toHaveBeenCalled() + const fetchOptions = mockFetch.mock.calls[0][1] + expect(fetchOptions.headers["ChatGPT-Account-Id"]).toBe("acct_12345") + }) + + it("should work without accountId when not available", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue(null) + + const mockFetch = vitest.fn().mockResolvedValue({ + ok: true, + json: () => + Promise.resolve({ + output: [ + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "response" }], + }, + ], + }), + text: () => Promise.resolve(""), + }) + vitest.stubGlobal("fetch", mockFetch) + + await handler.completePrompt("test prompt") + + expect(mockFetch).toHaveBeenCalled() + const fetchOptions = mockFetch.mock.calls[0][1] + expect(fetchOptions.headers["ChatGPT-Account-Id"]).toBeUndefined() + }) +}) diff --git a/src/api/providers/openai-codex.ts b/src/api/providers/openai-codex.ts index e9bc3bbf5d..3e64609ef2 100644 --- a/src/api/providers/openai-codex.ts +++ b/src/api/providers/openai-codex.ts @@ -29,6 +29,7 @@ import { isMcpTool } from "../../utils/mcp-name" import { sanitizeOpenAiCallId } from "../../utils/tool-id" import { openAiCodexOAuthManager } from "../../integrations/openai-codex/oauth" import { t } from "../../i18n" +import { mergeAbortSignalAndTimeout } from "./utils/abort-signal" export type OpenAiCodexModel = ReturnType @@ -274,7 +275,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // Make the request with retry on auth failure for (let attempt = 0; attempt < 2; attempt++) { try { - yield* this.executeRequest(requestBody, model, accessToken, effectiveSessionId) + yield* this.executeRequest(requestBody, model, accessToken, effectiveSessionId, metadata) return } catch (error) { const message = error instanceof Error ? error.message : String(error) @@ -438,10 +439,21 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion model: OpenAiCodexModel, accessToken: string, effectiveSessionId: string, + metadata?: ApiHandlerCreateMessageMetadata, ): ApiStream { // Create AbortController for cancellation this.abortController = new AbortController() + // Bridge the external abort signal into the internal controller (Bedrock pattern) + const externalAbortSignal = metadata?.abortSignal + if (externalAbortSignal) { + if (externalAbortSignal.aborted) { + this.abortController.abort() + } else { + externalAbortSignal.addEventListener("abort", () => this.abortController?.abort(), { once: true }) + } + } + try { // Prefer OpenAI SDK streaming (same approach as openai-native) so event handling // is consistent across providers. @@ -1257,7 +1269,8 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } async completePrompt(prompt: string, options?: CompletePromptOptions): Promise { - this.abortController = new AbortController() + const requestAbortSignal = mergeAbortSignalAndTimeout(options?.abortSignal, options?.timeoutMs) + const requestSignal = requestAbortSignal ?? new AbortController().signal try { const model = this.getModel() @@ -1318,7 +1331,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion method: "POST", headers, body: JSON.stringify(requestBody), - signal: this.abortController.signal, + signal: requestSignal, }) if (!response.ok) { @@ -1354,12 +1367,15 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion const apiError = new ApiProviderError(errorMessage, this.providerName, errorModel.id, "completePrompt") TelemetryService.instance.captureException(apiError) + // Re-throw abort errors as-is so callers can detect cancellation by the "AbortError" name + if (error instanceof Error && error.name === "AbortError") { + throw error + } + if (error instanceof Error) { throw new Error(t("common:errors.openAiCodex.completionError", { message: error.message })) } throw error - } finally { - this.abortController = undefined } } } From 22da1d1c873469be5fc4c5b4d7c94f316eb2a926 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Thu, 20 Aug 2026 00:50:20 +0800 Subject: [PATCH 02/12] fix(api): address CodeRabbit review on openai-codex abort handling - executeRequest: create a request-local AbortController (mirrored to this.abortController for existing abort handling); the external-signal bridge listener now captures the local controller and is removed in finally, so a late abort from an earlier request can no longer abort a newer request and listeners no longer leak - completePrompt: normalize any rejected request whose request-local signal aborted (external abort, AbortSignal.timeout "TimeoutError") to an error with name "AbortError", and throw the same AbortError when the transport quietly completes after cancellation - specs: bridge test now asserts the captured request-local SDK signal aborts mid-flight; merge tests assert AbortError rejection on quiet completion; new tests cover timeout cancellation and quiet completion after abort --- .../openai-codex-native-tool-calls.spec.ts | 60 ++++++++++------ .../providers/__tests__/openai-codex.spec.ts | 70 ++++++++++++++++++- src/api/providers/openai-codex.ts | 57 +++++++++++---- 3 files changed, 148 insertions(+), 39 deletions(-) diff --git a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts index b1ded8bbf3..37037e4adc 100644 --- a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts +++ b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts @@ -524,19 +524,30 @@ describe("OpenAiCodexHandler native tool calls", () => { vi.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") vi.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") - const mockCreate = vi.fn().mockResolvedValue({ - async *[Symbol.asyncIterator]() { - yield { type: "response.text.delta", delta: "test" } - yield { - type: "response.completed", - response: { - id: "resp_1", - status: "completed", - output: [{ type: "message", content: [{ type: "output_text", text: "test" }] }], - usage: { input_tokens: 1, output_tokens: 1 }, - }, - } - }, + // The mock transport pauses mid-flight until the request-local signal aborts + const mockCreate = vi.fn().mockImplementation(async (_body: unknown, init?: { signal?: AbortSignal }) => { + return { + async *[Symbol.asyncIterator]() { + yield { type: "response.text.delta", delta: "test" } + await new Promise((resolve) => { + const signal = init?.signal + if (!signal || signal.aborted) { + resolve() + return + } + signal.addEventListener("abort", () => resolve(), { once: true }) + }) + yield { + type: "response.completed", + response: { + id: "resp_1", + status: "completed", + output: [{ type: "message", content: [{ type: "output_text", text: "test" }] }], + usage: { input_tokens: 1, output_tokens: 1 }, + }, + } + }, + } }) Object.assign(handler, { client: { @@ -550,20 +561,25 @@ describe("OpenAiCodexHandler native tool calls", () => { abortSignal: controller.signal, }) - // Consume the stream to trigger the request - await collectStream(stream) + // Consume the stream (the mock transport pauses mid-flight) + const collected = collectStream(stream) + + // Wait until the request has started; the bridge listener is registered before + // the SDK call, so aborting now lands mid-flight + await vi.waitFor(() => expect(mockCreate).toHaveBeenCalled()) + + // Abort the external signal mid-flight; the bridge must abort the request-local controller + controller.abort() + + const chunks = await collected + expect(chunks.length).toBeGreaterThan(0) expect(mockCreate).toHaveBeenCalled() const createCallArgs = mockCreate.mock.calls[0][1] as { signal?: AbortSignal } + // The captured (request-local) signal passed to the SDK must now be aborted expect(createCallArgs.signal).toBeDefined() expect(createCallArgs.signal).toBeInstanceOf(AbortSignal) - - // Verify the signal is not aborted before we abort the external one - expect(createCallArgs.signal?.aborted).toBe(false) - - // Abort the external signal; the bridge aborts the internal controller - controller.abort() - expect(controller.signal.aborted).toBe(true) + expect(createCallArgs.signal?.aborted).toBe(true) }) it("should immediately abort when the external signal is already aborted", async () => { diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index 01ff123f0a..ab80805fb1 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -748,7 +748,11 @@ describe("OpenAiCodexHandler.completePrompt", () => { const controller = new AbortController() const promise = handler.completePrompt("test prompt", { abortSignal: controller.signal }) controller.abort() - await promise + // The mock transport completes without honoring the abort; an aborted request must + // reject with an AbortError rather than return the late response + await expect(promise).rejects.toSatisfy( + (error: unknown) => error instanceof Error && error.name === "AbortError", + ) expect(mockFetch).toHaveBeenCalled() expect(fetchInit?.signal).toBeDefined() @@ -784,7 +788,11 @@ describe("OpenAiCodexHandler.completePrompt", () => { const controller = new AbortController() const promise = handler.completePrompt("test prompt", { abortSignal: controller.signal, timeoutMs: 5000 }) controller.abort() - await promise + // The mock transport completes without honoring the abort; an aborted request must + // reject with an AbortError rather than return the late response + await expect(promise).rejects.toSatisfy( + (error: unknown) => error instanceof Error && error.name === "AbortError", + ) expect(mockFetch).toHaveBeenCalled() expect(fetchInit?.signal).toBeDefined() @@ -1004,4 +1012,62 @@ describe("OpenAiCodexHandler.completePrompt", () => { const fetchOptions = mockFetch.mock.calls[0][1] expect(fetchOptions.headers["ChatGPT-Account-Id"]).toBeUndefined() }) + + it("should reject with an AbortError when the timeout elapses", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + // Emulate a hung fetch that rejects when its signal aborts (native AbortSignal.timeout) + const mockFetch = vitest.fn().mockImplementation((_url: string, init?: RequestInit) => { + return new Promise((_resolve, reject) => { + init?.signal?.addEventListener( + "abort", + () => { + const abortError = new Error("The operation was aborted") + abortError.name = "AbortError" + reject(abortError) + }, + { once: true }, + ) + }) + }) + vitest.stubGlobal("fetch", mockFetch) + + await expect(handler.completePrompt("test prompt", { timeoutMs: 50 })).rejects.toSatisfy( + (error: unknown) => error instanceof Error && error.name === "AbortError", + ) + }) + + it("should reject with an AbortError when the signal aborts and the transport completes anyway", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.1-codex" }) + + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + + // Emulate a transport that resolves successfully despite the aborted signal + const mockFetch = vitest.fn().mockResolvedValue({ + ok: true, + json: () => + Promise.resolve({ + output: [ + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "late response" }], + }, + ], + }), + text: () => Promise.resolve(""), + }) + vitest.stubGlobal("fetch", mockFetch) + + const controller = new AbortController() + controller.abort() + + await expect(handler.completePrompt("test prompt", { abortSignal: controller.signal })).rejects.toSatisfy( + (error: unknown) => error instanceof Error && error.name === "AbortError", + ) + }) }) diff --git a/src/api/providers/openai-codex.ts b/src/api/providers/openai-codex.ts index 3e64609ef2..fb01c9e98d 100644 --- a/src/api/providers/openai-codex.ts +++ b/src/api/providers/openai-codex.ts @@ -441,16 +441,22 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion effectiveSessionId: string, metadata?: ApiHandlerCreateMessageMetadata, ): ApiStream { - // Create AbortController for cancellation - this.abortController = new AbortController() + // Create a request-local AbortController. The class field mirrors it so the + // existing abort handling keeps working, but the bridge listener below captures + // the local controller so a late abort can never hit a newer request's controller. + const requestController = new AbortController() + this.abortController = requestController - // Bridge the external abort signal into the internal controller (Bedrock pattern) + // Bridge the external abort signal into the request controller (Bedrock pattern) const externalAbortSignal = metadata?.abortSignal + const bridgeAbort = () => requestController.abort() + let bridgeCleanup: (() => void) | undefined if (externalAbortSignal) { if (externalAbortSignal.aborted) { - this.abortController.abort() + requestController.abort() } else { - externalAbortSignal.addEventListener("abort", () => this.abortController?.abort(), { once: true }) + externalAbortSignal.addEventListener("abort", bridgeAbort, { once: true }) + bridgeCleanup = () => externalAbortSignal.removeEventListener("abort", bridgeAbort) } } @@ -475,7 +481,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion }) const stream = (await (client as any).responses.create(requestBody, { - signal: this.abortController.signal, + signal: requestController.signal, // If the SDK supports per-request overrides, ensure headers are present. headers: codexHeaders, })) as AsyncIterable @@ -487,7 +493,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } for await (const event of stream) { - if (this.abortController.signal.aborted) { + if (requestController.signal.aborted) { break } @@ -503,7 +509,11 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion yield* this.makeCodexRequest(requestBody, model, accessToken, effectiveSessionId) } } finally { - this.abortController = undefined + bridgeCleanup?.() + // Only clear the field if this request still owns it + if (this.abortController === requestController) { + this.abortController = undefined + } } } @@ -1344,32 +1354,49 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion const responseData = await response.json() + let result: string | undefined if (responseData?.output && Array.isArray(responseData.output)) { for (const outputItem of responseData.output) { if (outputItem.type === "message" && outputItem.content) { for (const content of outputItem.content) { if (content.type === "output_text" && content.text) { - return content.text + result = content.text + break } } + if (result !== undefined) { + break + } } } } - if (responseData?.text) { - return responseData.text + if (result === undefined && responseData?.text) { + result = responseData.text } - return "" + // The request may have been cancelled while the transport was finishing; + // surface it as an abort instead of returning the completed response. + if (requestSignal.aborted) { + const abortError = new Error("This operation was aborted") + abortError.name = "AbortError" + throw abortError + } + + return result ?? "" } catch (error) { const errorModel = this.getModel() const errorMessage = error instanceof Error ? error.message : String(error) const apiError = new ApiProviderError(errorMessage, this.providerName, errorModel.id, "completePrompt") TelemetryService.instance.captureException(apiError) - // Re-throw abort errors as-is so callers can detect cancellation by the "AbortError" name - if (error instanceof Error && error.name === "AbortError") { - throw error + // An aborted request surfaces as "AbortError" (external signal) or "TimeoutError" + // (AbortSignal.timeout); normalize it so callers can detect cancellation by the + // "AbortError" name. + if (requestSignal.aborted) { + const abortError = new Error("This operation was aborted") + abortError.name = "AbortError" + throw abortError } if (error instanceof Error) { From 0ca6c23278761c314da09d39c92a35ba156b44d9 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Thu, 20 Aug 2026 19:00:27 +0800 Subject: [PATCH 03/12] refactor(api): adopt RequestConfigBuilder in feat/abort-r1-openai-codex abort wiring --- .../__tests__/request-config-builder.spec.ts | 40 +++++++++++++++++++ .../config-builder/request-config-builder.ts | 15 +++++++ src/api/providers/openai-codex.ts | 7 +++- 3 files changed, 60 insertions(+), 2 deletions(-) diff --git a/src/api/providers/__tests__/request-config-builder.spec.ts b/src/api/providers/__tests__/request-config-builder.spec.ts index 977b09df6b..f9557a123f 100644 --- a/src/api/providers/__tests__/request-config-builder.spec.ts +++ b/src/api/providers/__tests__/request-config-builder.spec.ts @@ -505,4 +505,44 @@ describe("RequestConfigBuilder", () => { expect(config.maxTokens).toBe(2000) }) }) + + describe("static merge helpers (canonical abort-signal entry points)", () => { + it("returns undefined from mergeAbortSignalAndTimeout when no external signal and no valid timeout", () => { + expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(undefined, undefined)).toBeUndefined() + expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(undefined, 0)).toBeUndefined() + expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(undefined, -5)).toBeUndefined() + }) + + it("returns the external signal directly when no timeout is merged", () => { + const controller = new AbortController() + expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(controller.signal, undefined)).toBe( + controller.signal, + ) + expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(controller.signal, 0)).toBe(controller.signal) + }) + + it("returns the primary signal directly from mergeAbortSignals when there is no secondary", () => { + const controller = new AbortController() + expect(RequestConfigBuilder.mergeAbortSignals(controller.signal)).toBe(controller.signal) + expect(RequestConfigBuilder.mergeAbortSignals(controller.signal, undefined)).toBe(controller.signal) + }) + + it("delegates to AbortSignal.any when two distinct signals are merged", () => { + const a = new AbortController() + const b = new AbortController() + const merged = RequestConfigBuilder.mergeAbortSignals(a.signal, b.signal) + expect(merged.aborted).toBe(false) + b.abort() + expect(merged.aborted).toBe(true) + }) + + it("aborts the merged signal when the primary signal aborts", () => { + const a = new AbortController() + const b = new AbortController() + const merged = RequestConfigBuilder.mergeAbortSignals(a.signal, b.signal) + expect(merged.aborted).toBe(false) + a.abort() + expect(merged.aborted).toBe(true) + }) + }) }) diff --git a/src/api/providers/config-builder/request-config-builder.ts b/src/api/providers/config-builder/request-config-builder.ts index 2201d735bc..a3c1ba7ebb 100644 --- a/src/api/providers/config-builder/request-config-builder.ts +++ b/src/api/providers/config-builder/request-config-builder.ts @@ -163,4 +163,19 @@ export class RequestConfigBuilder @@ -1279,7 +1279,10 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } async completePrompt(prompt: string, options?: CompletePromptOptions): Promise { - const requestAbortSignal = mergeAbortSignalAndTimeout(options?.abortSignal, options?.timeoutMs) + const requestAbortSignal = RequestConfigBuilder.mergeAbortSignalAndTimeout( + options?.abortSignal, + options?.timeoutMs, + ) const requestSignal = requestAbortSignal ?? new AbortController().signal try { From 714754cfa78f4918189feec0b2b829425a2786b4 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Sat, 5 Sep 2026 14:30:13 +0800 Subject: [PATCH 04/12] fix(api): address CodeRabbit abort-signal findings in openai-codex --- .../openai-codex-native-tool-calls.spec.ts | 10 ++- .../providers/__tests__/openai-codex.spec.ts | 67 ++++++++++++++++ src/api/providers/openai-codex.ts | 80 ++++++++++++++----- 3 files changed, 133 insertions(+), 24 deletions(-) diff --git a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts index b4bdf87c04..133b97bbf0 100644 --- a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts +++ b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts @@ -577,7 +577,9 @@ describe("OpenAiCodexHandler native tool calls", () => { controller.abort() const chunks = await collected - expect(chunks.length).toBeGreaterThan(0) + // Exactly the pre-abort delta: the completed event only arrives once the abort resolves + // the transport's pause, so the loop must break before processing it. + expect(chunks).toEqual([{ type: "text", text: "test" }]) expect(mockCreate).toHaveBeenCalled() const createCallArgs = mockCreate.mock.calls[0][1] as { signal?: AbortSignal } @@ -619,8 +621,10 @@ describe("OpenAiCodexHandler native tool calls", () => { abortSignal: controller.signal, }) - // Consume the stream to trigger the request - await collectStream(stream) + // Consume the stream to trigger the request; the request-local controller is already + // aborted, so the loop must break before the first event is processed. + const chunks = await collectStream(stream) + expect(chunks).toEqual([]) expect(mockCreate).toHaveBeenCalled() const createCallArgs = mockCreate.mock.calls[0][1] as { signal?: AbortSignal } diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index e11b35256f..d88a2b2172 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -9,6 +9,7 @@ vitest.mock("@roo-code/telemetry", () => ({ })) import { Anthropic } from "@anthropic-ai/sdk" +import { TelemetryService } from "@roo-code/telemetry" import { OPEN_AI_CODEX_SERVICE_TIER_KEY, OpenAiCodexServiceTier, SERVICE_TIER_KEY } from "@roo-code/types" import { OpenAiCodexHandler, transformResponsesLiteBody } from "../openai-codex" import { openAiCodexOAuthManager } from "../../../integrations/openai-codex/oauth" @@ -591,6 +592,72 @@ describe("OpenAiCodexHandler.completePrompt streaming", () => { expect(mockFetch).not.toHaveBeenCalled() }) + // The cancellation wins over the auth retry: force-refreshing a token for a request that is + // already gone would spend a network round trip on a dead request and surface the + // cancellation as an authentication failure. + it("fails fast with the abort contract instead of force-refreshing the token after cancellation", async () => { + const handler = createHandler() + const refresh = vitest.spyOn(openAiCodexOAuthManager, "forceRefreshAccessToken") + const controller = new AbortController() + // The SDK keeps the request open until the signal it was handed aborts, then rejects the way + // it rejects aborted requests. + const create = vitest.fn().mockImplementation( + (_body: unknown, options: { signal: AbortSignal }) => + new Promise((_resolve, reject) => { + options.signal.addEventListener("abort", () => reject(new Error("Request was aborted")), { + once: true, + }) + }), + ) + Reflect.set(handler, "client", { responses: { create } }) + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + const completion = handler.completePrompt("Hello", { abortSignal: controller.signal }) + await vitest.waitFor(() => expect(create).toHaveBeenCalled()) + controller.abort() + + await expect(completion).rejects.toMatchObject({ name: "AbortError" }) + expect(refresh).not.toHaveBeenCalled() + expect(create).toHaveBeenCalledTimes(1) + expect(mockFetch).not.toHaveBeenCalled() + }) + + // The abort lands while the fallback fetch is in flight, so the cancellation must come out as + // the shared abort contract - not a telemetry event and not a wrapped connection error. + it("keeps the abort contract and skips telemetry when the fallback fetch is cancelled", async () => { + const handler = createHandler() + // The module mock keeps the spy across tests, so clear it before asserting on this request + const captureException = vitest.mocked(TelemetryService.instance.captureException) + captureException.mockClear() + // The SDK path is unusable, so the request falls back to the SSE transport. + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const controller = new AbortController() + // Reject the way fetch rejects once the signal it was handed aborts. + const mockFetch = vitest.fn((_url: unknown, init?: { signal?: AbortSignal }) => { + const signal = init?.signal + if (!signal || signal.aborted) { + return Promise.reject(new DOMException("The operation was aborted", "AbortError")) + } + return new Promise((_resolve, reject) => { + signal.addEventListener( + "abort", + () => reject(new DOMException("The operation was aborted", "AbortError")), + { once: true }, + ) + }) + }) + vitest.stubGlobal("fetch", mockFetch) + + const completion = handler.completePrompt("Hello", { abortSignal: controller.signal }) + await vitest.waitFor(() => expect(mockFetch).toHaveBeenCalled()) + controller.abort() + + await expect(completion).rejects.toMatchObject({ name: "AbortError" }) + expect(captureException).not.toHaveBeenCalled() + }) + it("wraps failures from both transports as a completion error", async () => { const handler = createHandler() const create = vitest.fn().mockRejectedValue(new Error("sdk down")) diff --git a/src/api/providers/openai-codex.ts b/src/api/providers/openai-codex.ts index 120311c23b..78279625e4 100644 --- a/src/api/providers/openai-codex.ts +++ b/src/api/providers/openai-codex.ts @@ -30,6 +30,7 @@ import { sanitizeOpenAiCallId } from "../../utils/tool-id" import { openAiCodexOAuthManager } from "../../integrations/openai-codex/oauth" import { t } from "../../i18n" import { RequestConfigBuilder } from "./config-builder/request-config-builder" +import { createAbortError } from "./utils/abort-signal" export type OpenAiCodexModel = ReturnType @@ -291,6 +292,13 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion yield* this.executeRequest(requestBody, model, accessToken, effectiveSessionId, abortSignal) return } catch (error) { + // The caller's cancellation wins over the retry: force-refreshing a token for a + // request that is already gone would spend a network round trip on a dead request + // and surface the cancellation as an authentication failure. + if (abortSignal?.aborted) { + throw createAbortError(this.providerName) + } + const message = error instanceof Error ? error.message : String(error) const isAuthFailure = /unauthorized|invalid token|not authenticated|authentication|401/i.test(message) @@ -457,16 +465,19 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion effectiveSessionId: string, abortSignal?: AbortSignal, ): ApiStream { - // Create AbortController for cancellation - this.abortController = new AbortController() + // Request-local controller: a stale abort arriving after this request finishes must never + // reach the controller of a later request, so the bridge listener captures this controller + // directly instead of reading `this.abortController` at abort time. + const abortController = new AbortController() + this.abortController = abortController // A caller's signal has to be linked rather than used directly, since both transports below - // abort through `this.abortController`. Without this the signal never reaches the wire. - const abortFromCaller = () => this.abortController?.abort() + // abort through the request controller. Without this the signal never reaches the wire. + const abortFromCaller = () => abortController.abort() if (abortSignal) { if (abortSignal.aborted) { - this.abortController.abort() + abortController.abort() } else { abortSignal.addEventListener("abort", abortFromCaller, { once: true }) } @@ -493,7 +504,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion }) const stream = (await (client as any).responses.create(requestBody, { - signal: this.abortController.signal, + signal: abortController.signal, // If the SDK supports per-request overrides, ensure headers are present. headers: codexHeaders, })) as AsyncIterable @@ -505,7 +516,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } for await (const event of stream) { - if (this.abortController.signal.aborted) { + if (abortController.signal.aborted) { break } @@ -528,16 +539,26 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // A cancellation is not a transport failure either. Falling back would spend a // second request on an already-aborted signal and report the cancellation as a // connection error. - if (this.sawSdkEventInCurrentResponse || this.abortController?.signal.aborted) { + if (this.sawSdkEventInCurrentResponse || abortController.signal.aborted) { throw sdkErr } // Fallback to manual SSE via fetch (Codex backend). - yield* this.makeCodexRequest(requestBody, model, accessToken, effectiveSessionId) + yield* this.makeCodexRequest( + requestBody, + model, + accessToken, + effectiveSessionId, + abortController.signal, + ) } } finally { abortSignal?.removeEventListener("abort", abortFromCaller) - this.abortController = undefined + // Only clear the field if this request still owns it: a concurrent request may have + // installed its own controller after this one started. + if (this.abortController === abortController) { + this.abortController = undefined + } } } @@ -629,6 +650,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion model: OpenAiCodexModel, accessToken: string, effectiveSessionId: string, + abortSignal?: AbortSignal, ): ApiStream { // Per the implementation guide: route to Codex backend with Bearer token const url = `${CODEX_API_BASE_URL}/responses` @@ -648,7 +670,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion method: "POST", headers, body: JSON.stringify(requestBody), - signal: this.abortController?.signal, + signal: abortSignal, }) if (!response.ok) { @@ -708,8 +730,15 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion throw new Error(t("common:errors.openAiCodex.noResponseBody")) } - yield* this.handleStreamResponse(response.body, model) + yield* this.handleStreamResponse(response.body, model, abortSignal) } catch (error) { + // The caller cancelled, so this is not a transport fault: hand back the shared abort + // contract instead of reporting the caller's own cancellation to telemetry or wrapping it + // as a connection failure. + if (abortSignal?.aborted) { + throw createAbortError(this.providerName) + } + const errorMessage = error instanceof Error ? error.message : String(error) const apiError = new ApiProviderError(errorMessage, this.providerName, model.id, "createMessage") TelemetryService.instance.captureException(apiError) @@ -724,7 +753,11 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } } - private async *handleStreamResponse(body: ReadableStream, model: OpenAiCodexModel): ApiStream { + private async *handleStreamResponse( + body: ReadableStream, + model: OpenAiCodexModel, + abortSignal?: AbortSignal, + ): ApiStream { const reader = body.getReader() const decoder = new TextDecoder() let buffer = "" @@ -732,7 +765,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion try { while (true) { - if (this.abortController?.signal.aborted) { + if (abortSignal?.aborted) { break } @@ -987,6 +1020,13 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } } } catch (error) { + // The caller cancelled while the fallback stream was being read: that is not a + // stream-processing failure, so hand back the shared abort contract instead of reporting + // the caller's own cancellation to telemetry. + if (abortSignal?.aborted) { + throw createAbortError(this.providerName) + } + const errorMessage = error instanceof Error ? error.message : String(error) const apiError = new ApiProviderError(errorMessage, this.providerName, model.id, "createMessage") TelemetryService.instance.captureException(apiError) @@ -1362,19 +1402,17 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // throwing - so returning here would report a cancelled generation as a finished one and // hand the caller whatever partial text had arrived. if (requestSignal?.aborted) { - throw new DOMException("OpenAI Codex completion was aborted", "AbortError") + throw createAbortError(this.providerName) } return text } catch (error) { // Cancelling is the caller's own doing, not a provider failure, so it is neither - // reported to telemetry nor relabelled as a completion error. A transport that rejects - // on abort reports it in its own words, so it is restated here: callers get one abort - // result whether the stream ended quietly or the request threw. + // reported to telemetry nor relabelled as a completion error: a timed-out or + // cancelled request is normalized to the shared abort contract whether the stream + // ended quietly or the request threw. if (requestSignal?.aborted) { - throw error instanceof DOMException && error.name === "AbortError" - ? error - : new DOMException("OpenAI Codex completion was aborted", "AbortError") + throw createAbortError(this.providerName) } const errorModel = this.getModel() From 1aa5fcd588b984d1e50321680488e2589d9a7c6b Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Sat, 5 Sep 2026 15:25:08 +0800 Subject: [PATCH 05/12] fix(api): close the openai-codex abort mutation gaps --- .../openai-codex-native-tool-calls.spec.ts | 8 +- .../providers/__tests__/openai-codex.spec.ts | 259 ++++++++++++++++++ src/api/providers/openai-codex.ts | 10 +- 3 files changed, 270 insertions(+), 7 deletions(-) diff --git a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts index 133b97bbf0..f831d22fc1 100644 --- a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts +++ b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts @@ -551,6 +551,9 @@ describe("OpenAiCodexHandler native tool calls", () => { usage: { input_tokens: 1, output_tokens: 1 }, }, } + // Arrives after the abort resolves the pause: the loop guard must break before + // processing the completed event and this delta. + yield { type: "response.text.delta", delta: "post-abort" } }, } }) @@ -577,8 +580,9 @@ describe("OpenAiCodexHandler native tool calls", () => { controller.abort() const chunks = await collected - // Exactly the pre-abort delta: the completed event only arrives once the abort resolves - // the transport's pause, so the loop must break before processing it. + // Exactly the pre-abort delta: the completed event and the post-abort delta only arrive + // once the abort resolves the transport's pause, so the loop guard must break before + // processing either. expect(chunks).toEqual([{ type: "text", text: "test" }]) expect(mockCreate).toHaveBeenCalled() diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index d88a2b2172..4a9fa7ac68 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -1171,3 +1171,262 @@ describe("OpenAiCodexHandler.completePrompt timeout", () => { expect(mockFetch).not.toHaveBeenCalled() }) }) + +describe("OpenAiCodexHandler.createMessage abort bridging", () => { + // These tests drive createMessage directly (not completePrompt): completePrompt re-normalizes + // any error to the shared abort contract once the caller's signal has fired, which would mask + // regressions in the abort checks inside the transports. + + function createHandler() { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.6-sol" }) + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("test-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + return handler + } + + afterEach(() => { + vitest.restoreAllMocks() + vitest.unstubAllGlobals() + }) + + it("hands back the abort contract instead of force-refreshing after the caller cancels", async () => { + const handler = createHandler() + const refresh = vitest + .spyOn(openAiCodexOAuthManager, "forceRefreshAccessToken") + .mockResolvedValue("refreshed-token") + // The SDK fails with exactly the auth-failure wording the retry path would act on, so the + // abort check must win over the refresh-and-retry logic. + const create = vitest.fn().mockRejectedValue(new Error("401 invalid token")) + Reflect.set(handler, "client", { responses: { create } }) + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + await expect( + collectStream( + handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: AbortSignal.abort(), + }), + ), + ).rejects.toMatchObject({ name: "AbortError" }) + + // No refresh, no second SDK attempt, no SSE fallback. + expect(refresh).not.toHaveBeenCalled() + expect(create).toHaveBeenCalledTimes(1) + expect(mockFetch).not.toHaveBeenCalled() + }) + + it("emits nothing once the caller cancels before the first SDK event", async () => { + const handler = createHandler() + const controller = new AbortController() + const create = vitest.fn().mockImplementation(() => { + // The caller cancels while the SDK is still delivering, so the loop's abort check must + // stop every event from reaching the caller. + controller.abort() + return Promise.resolve( + asyncStreamFrom([ + { type: "response.output_text.delta", delta: "feat: half a" }, + { type: "response.output_text.delta", delta: "and the rest" }, + { type: "response.completed", response: { id: "r1", status: "completed", output: [] } }, + ]), + ) + }) + Reflect.set(handler, "client", { responses: { create } }) + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + const chunks = await collectStream( + handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }), + ) + + // The generators end quietly on abort - they break rather than throw - so an empty stream + // is the observable proof that nothing was processed after the cancellation. + expect(chunks).toEqual([]) + expect(mockFetch).not.toHaveBeenCalled() + }) + + it("clears the request controller once the request ends", async () => { + const handler = createHandler() + const create = vitest.fn().mockResolvedValue( + asyncStreamFrom([ + { type: "response.output_text.delta", delta: "done" }, + { type: "response.completed", response: { id: "r1", status: "completed", output: [] } }, + ]), + ) + Reflect.set(handler, "client", { responses: { create } }) + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + await collectStream(handler.createMessage("System", [{ role: "user", content: "Hello" }])) + + // A later request must be able to tell that no request is in flight. + expect(Reflect.get(handler, "abortController")).toBeUndefined() + }) + + it("keeps the later request's controller when the earlier request finishes", async () => { + const handler = createHandler() + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + // The earlier request holds its stream open until the later request has installed its own + // controller, so the earlier cleanup runs while a request is still in flight. + let releaseEarlier: (() => void) | undefined + const releaseGate = new Promise((resolve) => { + releaseEarlier = resolve + }) + const earlierStream = (async function* () { + yield { type: "response.output_text.delta", delta: "a" } + await releaseGate + yield { type: "response.completed", response: { id: "ra", status: "completed", output: [] } } + })() + const create = vitest + .fn() + .mockImplementationOnce(() => Promise.resolve(earlierStream)) + .mockImplementationOnce(() => + Promise.resolve( + asyncStreamFrom([ + { type: "response.output_text.delta", delta: "b" }, + { type: "response.completed", response: { id: "rb", status: "completed", output: [] } }, + ]), + ), + ) + Reflect.set(handler, "client", { responses: { create } }) + + const earlier = handler.createMessage("System", [{ role: "user", content: "Hello" }]) + expect(await earlier.next()).toMatchObject({ value: { type: "text", text: "a" } }) + + const later = handler.createMessage("System", [{ role: "user", content: "Hello" }]) + expect(await later.next()).toMatchObject({ value: { type: "text", text: "b" } }) + + // The earlier request finishes while the later one is still in flight. + releaseEarlier!() + await earlier.next() + + // The earlier request's cleanup must not clear the controller the later request installed. + const controller = Reflect.get(handler, "abortController") as AbortController | undefined + expect(controller).toBeDefined() + expect(controller?.signal).toBe(create.mock.calls[1][1].signal) + + // Let the later request finish and clear its own controller. + await later.next() + }) + + it("stops reading the fallback stream once the request aborts", async () => { + const handler = createHandler() + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const controller = new AbortController() + const encoder = new TextEncoder() + const body = new ReadableStream({ + start(streamController) { + streamController.enqueue( + encoder.encode('data: {"type":"response.output_text.delta","delta":"one"}\n\n'), + ) + streamController.enqueue( + encoder.encode('data: {"type":"response.output_text.delta","delta":"two"}\n\n'), + ) + streamController.close() + }, + }) + const mockFetch = vitest.fn().mockResolvedValue({ ok: true, body }) + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + expect(await iter.next()).toMatchObject({ value: { type: "text", text: "one" } }) + + // The cancellation lands between two stream reads, exactly where the loop's check runs. + controller.abort() + + const chunks = await collectStream(iter) + // Everything enqueued after the cancellation must stay unread. + expect(chunks).toEqual([]) + }) + + it("hands back the abort contract when the fallback stream tears down while the caller cancels", async () => { + const handler = createHandler() + const captureException = vitest.mocked(TelemetryService.instance.captureException) + captureException.mockClear() + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const controller = new AbortController() + const encoder = new TextEncoder() + let pullStartedResolve: (() => void) | undefined + let failRead: (() => void) | undefined + const pullStarted = new Promise((resolve) => { + pullStartedResolve = resolve + }) + const failGate = new Promise((resolve) => { + failRead = resolve + }) + const body = new ReadableStream({ + start(streamController) { + streamController.enqueue( + encoder.encode('data: {"type":"response.output_text.delta","delta":"one"}\n\n'), + ) + }, + pull(streamController) { + // The second read is pending; tear the stream down on the test's signal. + pullStartedResolve!() + return failGate.then(() => { + streamController.error(new Error("stream torn down")) + }) + }, + }) + const mockFetch = vitest.fn().mockResolvedValue({ ok: true, body }) + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + expect(await iter.next()).toMatchObject({ value: { type: "text", text: "one" } }) + + const pending = collectStream(iter) + await pullStarted + // The cancellation lands while the second read is in flight, i.e. after the loop's check. + controller.abort() + failRead!() + + await expect(pending).rejects.toMatchObject({ name: "AbortError" }) + // Cancellation is the caller's own doing, so neither catch may report it to telemetry. + expect(captureException).not.toHaveBeenCalled() + }) + + it("wraps a torn-down fallback stream as a stream error when nothing was aborted", async () => { + const handler = createHandler() + const captureException = vitest.mocked(TelemetryService.instance.captureException) + captureException.mockClear() + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const encoder = new TextEncoder() + const body = new ReadableStream({ + start(streamController) { + streamController.enqueue( + encoder.encode('data: {"type":"response.output_text.delta","delta":"one"}\n\n'), + ) + }, + pull(streamController) { + streamController.error(new Error("stream torn down")) + }, + }) + const mockFetch = vitest.fn().mockResolvedValue({ ok: true, body }) + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }]) + expect(await iter.next()).toMatchObject({ value: { type: "text", text: "one" } }) + + // The wrap chain surfaces the connection-failure key in every case, so the message alone + // cannot prove the innermost check classified this as a stream failure: an always-true + // check would swap in the shared abort contract, and the request-level catch would wrap + // that in the same key. The telemetry count is the witness: the stream-processing catch + // and the request catch both report the failure (twice); the abort path skips the first. + await expect(collectStream(iter)).rejects.toThrow(/connectionFailed|stream torn down/) + expect(captureException).toHaveBeenCalledTimes(2) + }) +}) diff --git a/src/api/providers/openai-codex.ts b/src/api/providers/openai-codex.ts index 78279625e4..1e9885371d 100644 --- a/src/api/providers/openai-codex.ts +++ b/src/api/providers/openai-codex.ts @@ -650,7 +650,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion model: OpenAiCodexModel, accessToken: string, effectiveSessionId: string, - abortSignal?: AbortSignal, + abortSignal: AbortSignal, ): ApiStream { // Per the implementation guide: route to Codex backend with Bearer token const url = `${CODEX_API_BASE_URL}/responses` @@ -735,7 +735,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // The caller cancelled, so this is not a transport fault: hand back the shared abort // contract instead of reporting the caller's own cancellation to telemetry or wrapping it // as a connection failure. - if (abortSignal?.aborted) { + if (abortSignal.aborted) { throw createAbortError(this.providerName) } @@ -756,7 +756,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion private async *handleStreamResponse( body: ReadableStream, model: OpenAiCodexModel, - abortSignal?: AbortSignal, + abortSignal: AbortSignal, ): ApiStream { const reader = body.getReader() const decoder = new TextDecoder() @@ -765,7 +765,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion try { while (true) { - if (abortSignal?.aborted) { + if (abortSignal.aborted) { break } @@ -1023,7 +1023,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // The caller cancelled while the fallback stream was being read: that is not a // stream-processing failure, so hand back the shared abort contract instead of reporting // the caller's own cancellation to telemetry. - if (abortSignal?.aborted) { + if (abortSignal.aborted) { throw createAbortError(this.providerName) } From 0183b943abf523c2ed1c3dff7b9fcea1ac215cd1 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Tue, 15 Sep 2026 23:09:05 +0800 Subject: [PATCH 06/12] fix(api): import mergeAbortSignalAndTimeout directly; drop builder static wrappers Address edelauna's review nits on the abort-signal API surface: openai-codex's completePrompt now imports mergeAbortSignalAndTimeout straight from utils/abort-signal (same pattern as Bedrock) instead of going through RequestConfigBuilder, and the two static pass-through wrappers (mergeAbortSignalAndTimeout, mergeAbortSignals) are dropped - they had no production callers once the direct import landed, and their behavior is covered by the abort-signal util's own spec. The request-config-builder files revert to main, shrinking the PR diff from 5 files to 3. --- .../__tests__/request-config-builder.spec.ts | 40 ------------------- .../config-builder/request-config-builder.ts | 15 ------- src/api/providers/openai-codex.ts | 5 +-- 3 files changed, 2 insertions(+), 58 deletions(-) diff --git a/src/api/providers/__tests__/request-config-builder.spec.ts b/src/api/providers/__tests__/request-config-builder.spec.ts index f9557a123f..977b09df6b 100644 --- a/src/api/providers/__tests__/request-config-builder.spec.ts +++ b/src/api/providers/__tests__/request-config-builder.spec.ts @@ -505,44 +505,4 @@ describe("RequestConfigBuilder", () => { expect(config.maxTokens).toBe(2000) }) }) - - describe("static merge helpers (canonical abort-signal entry points)", () => { - it("returns undefined from mergeAbortSignalAndTimeout when no external signal and no valid timeout", () => { - expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(undefined, undefined)).toBeUndefined() - expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(undefined, 0)).toBeUndefined() - expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(undefined, -5)).toBeUndefined() - }) - - it("returns the external signal directly when no timeout is merged", () => { - const controller = new AbortController() - expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(controller.signal, undefined)).toBe( - controller.signal, - ) - expect(RequestConfigBuilder.mergeAbortSignalAndTimeout(controller.signal, 0)).toBe(controller.signal) - }) - - it("returns the primary signal directly from mergeAbortSignals when there is no secondary", () => { - const controller = new AbortController() - expect(RequestConfigBuilder.mergeAbortSignals(controller.signal)).toBe(controller.signal) - expect(RequestConfigBuilder.mergeAbortSignals(controller.signal, undefined)).toBe(controller.signal) - }) - - it("delegates to AbortSignal.any when two distinct signals are merged", () => { - const a = new AbortController() - const b = new AbortController() - const merged = RequestConfigBuilder.mergeAbortSignals(a.signal, b.signal) - expect(merged.aborted).toBe(false) - b.abort() - expect(merged.aborted).toBe(true) - }) - - it("aborts the merged signal when the primary signal aborts", () => { - const a = new AbortController() - const b = new AbortController() - const merged = RequestConfigBuilder.mergeAbortSignals(a.signal, b.signal) - expect(merged.aborted).toBe(false) - a.abort() - expect(merged.aborted).toBe(true) - }) - }) }) diff --git a/src/api/providers/config-builder/request-config-builder.ts b/src/api/providers/config-builder/request-config-builder.ts index a3c1ba7ebb..2201d735bc 100644 --- a/src/api/providers/config-builder/request-config-builder.ts +++ b/src/api/providers/config-builder/request-config-builder.ts @@ -163,19 +163,4 @@ export class RequestConfigBuilder @@ -1372,7 +1371,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion async completePrompt(prompt: string, options?: CompletePromptOptions): Promise { // Merge an optional timeout into the caller's abort signal so a timeout cancels the // completion the same way an external abort does (timeoutMs <= 0 disables it). - const requestSignal = RequestConfigBuilder.mergeAbortSignalAndTimeout(options?.abortSignal, options?.timeoutMs) + const requestSignal = mergeAbortSignalAndTimeout(options?.abortSignal, options?.timeoutMs) try { const model = this.getModel() From 73b2b237476064810fbddb783f45854801aee973 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Wed, 16 Sep 2026 00:02:52 +0800 Subject: [PATCH 07/12] fix(api): fast-fail and guard the openai-codex completePrompt stream on abort Close the abort-signal contract gaps the series alignment check flags in completePrompt (in this unit's delta): (1) throwIfAborted fast-fail before the first await, matching the sibling units - a pre-aborted request no longer spends the OAuth token/account setup or an SDK request on a dead completion (the regression test defers the OAuth resolution to prove the fast-fail does not wait on it); (2) top-of-loop abort break on the streaming consumer loop so a buffered post-abort chunk is never joined into the completion; (3) the catch normalizes via isRequestAborted(error, requestSignal) like the sibling units, so an SDK abort error is normalized to the shared abort contract even when the signal has not marked itself aborted. --- .../providers/__tests__/openai-codex.spec.ts | 24 ++++++++++++++----- src/api/providers/openai-codex.ts | 22 +++++++++++++---- 2 files changed, 36 insertions(+), 10 deletions(-) diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index 4a9fa7ac68..5f089692f5 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -492,17 +492,29 @@ describe("OpenAiCodexHandler.completePrompt streaming", () => { expect(signalDuringRequest!.aborted).toBe(true) }) - it("rejects when the caller's signal is already aborted", async () => { - const handler = createHandler() - const create = injectStream(handler, [ - { type: "response.completed", response: { id: "r1", status: "completed", output: [] } }, - ]) + // A request cancelled before it starts must not spend the provider setup on it: the + // fast-fail rejects before the (here deferred, never resolving) token and account fetch + // and before the SDK request, so a pending OAuth flow cannot hold a cancelled completion + // hostage. + it("fast-fails before the OAuth setup when the caller's signal is already aborted", async () => { + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.6-sol" }) + const getAccessToken = vitest + .spyOn(openAiCodexOAuthManager, "getAccessToken") + .mockReturnValue(new Promise(() => {})) + const getAccountId = vitest + .spyOn(openAiCodexOAuthManager, "getAccountId") + .mockReturnValue(new Promise(() => {})) + const create = vitest.fn() + Reflect.set(handler, "client", { responses: { create } }) await expect(handler.completePrompt("Hello", { abortSignal: AbortSignal.abort() })).rejects.toMatchObject({ name: "AbortError", + message: "This operation was aborted", }) - expect(create.mock.calls[0][1].signal.aborted).toBe(true) + expect(create).not.toHaveBeenCalled() + expect(getAccessToken).not.toHaveBeenCalled() + expect(getAccountId).not.toHaveBeenCalled() }) // The SSE fallback is for an SDK that could not be used at all. Replaying the request after the diff --git a/src/api/providers/openai-codex.ts b/src/api/providers/openai-codex.ts index 7c8a68cdfd..ddf0393b71 100644 --- a/src/api/providers/openai-codex.ts +++ b/src/api/providers/openai-codex.ts @@ -29,7 +29,7 @@ import { isMcpTool } from "../../utils/mcp-name" import { sanitizeOpenAiCallId } from "../../utils/tool-id" import { openAiCodexOAuthManager } from "../../integrations/openai-codex/oauth" import { t } from "../../i18n" -import { createAbortError, mergeAbortSignalAndTimeout } from "./utils/abort-signal" +import { createAbortError, isRequestAborted, mergeAbortSignalAndTimeout, throwIfAborted } from "./utils/abort-signal" export type OpenAiCodexModel = ReturnType @@ -1369,6 +1369,11 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion * from having to be duplicated here. */ async completePrompt(prompt: string, options?: CompletePromptOptions): Promise { + // Fast-fail if the caller's stop signal already fired before we started: a cancelled + // request must not spend the OAuth setup (token and account fetch) or an SDK request on + // a completion that is already gone. + throwIfAborted(options?.abortSignal) + // Merge an optional timeout into the caller's abort signal so a timeout cancels the // completion the same way an external abort does (timeoutMs <= 0 disables it). const requestSignal = mergeAbortSignalAndTimeout(options?.abortSignal, options?.timeoutMs) @@ -1380,7 +1385,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // the prompt enhancer writes this straight into the input box. let text = "" - for await (const chunk of this.handleResponsesApiMessage( + const stream = this.handleResponsesApiMessage( model, "", [{ role: "user", content: prompt }], @@ -1388,7 +1393,16 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // directly, so `prompt_cache_key` is unchanged. { taskId: this.sessionId }, requestSignal, - )) { + ) + + for await (const chunk of stream) { + // A buffered chunk can still be pulled in the window between the abort and the + // inner generator's own stop, so break here: post-abort output must never be + // joined into the completion. + if (requestSignal?.aborted) { + break + } + // Refusals are streamed as text for the chat, but they are not output: the // non-streaming request this replaced read `output_text`, which never carries them. // Keeping them would paste "[Refusal] ..." into the input box as if it were an answer. @@ -1410,7 +1424,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // reported to telemetry nor relabelled as a completion error: a timed-out or // cancelled request is normalized to the shared abort contract whether the stream // ended quietly or the request threw. - if (requestSignal?.aborted) { + if (isRequestAborted(error, requestSignal)) { throw createAbortError(this.providerName) } From 641400265dddf961b9d16231547bf5a194901f95 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Wed, 16 Sep 2026 08:34:24 +0800 Subject: [PATCH 08/12] test(api): structural kill test for the completePrompt top-of-loop abort guard The guard at the top of completePrompt's streaming consumer loop is the only thing that stops the loop from pulling the SDK stream once more after the abort rides in on an in-flight chunk; the post-loop check throws the same AbortError either way, so the kill test asserts the pull count (event2's delta getter fires the abort after executeRequest's own check has passed, before completePrompt's) - the shape the class-h rule requires for consumer-loop guards. --- .../providers/__tests__/openai-codex.spec.ts | 45 +++++++++++++++++++ 1 file changed, 45 insertions(+) diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index 5f089692f5..63c72403ae 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -517,6 +517,51 @@ describe("OpenAiCodexHandler.completePrompt streaming", () => { expect(getAccountId).not.toHaveBeenCalled() }) + // The abort lands in the only window the consumer guard can still act: event2's `delta` + // getter fires it while processEvent builds the chunk - after executeRequest's + // top-of-loop check has passed, before completePrompt's own. The post-loop check throws + // the same AbortError whether or not the guard breaks, so what distinguishes the guarded + // loop from a mutated one is the pull count: the guard stops pulling after the chunk that + // the abort rode in on, a loop that keeps running pulls the SDK stream a third time. + it("breaks the streaming consumer loop at the top-of-loop guard and stops pulling", async () => { + const handler = createHandler() + const controller = new AbortController() + let sdkPulls = 0 + + const event1 = { type: "response.output_text.delta", delta: "pre-abort" } + const event2 = { + type: "response.output_text.delta", + get delta() { + controller.abort() + return "post-abort" + }, + } + + const create = vitest.fn().mockImplementation(() => { + return Promise.resolve({ + [Symbol.asyncIterator]() { + return { + next: async () => { + sdkPulls++ + return { value: sdkPulls === 1 ? event1 : event2, done: false } + }, + return: async () => ({ value: undefined, done: true }), + } + }, + }) + }) + Reflect.set(handler, "client", { responses: { create } }) + + await expect(handler.completePrompt("Hello", { abortSignal: controller.signal })).rejects.toMatchObject({ + name: "AbortError", + }) + + // event1 is pulled and joined; event2 is pulled (its getter fires the abort) but the + // guard breaks before it is joined - a third pull only happens when the mutated + // top-of-loop check keeps the loop running. + expect(sdkPulls).toBe(2) + }) + // The SSE fallback is for an SDK that could not be used at all. Replaying the request after the // SDK has already produced output would append a second generation to the first. it("does not replay over SSE when the SDK fails after emitting", async () => { From a27305b5ca5390279ab61faa0727b43faaffa1f2 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Wed, 23 Sep 2026 14:25:33 +0800 Subject: [PATCH 09/12] fix(api): observe cancellation in openai-codex OAuth setup and streaming CodeRabbit's out-of-diff review (5217406601) on this branch flagged two gaps: 1. OAuth setup did not observe cancellation. getAccessToken() and the getAccountId() awaits in executeRequest() and makeCodexRequest() kept running past a caller cancellation, so a cancelled request could sit on a pending lookup or hand the SDK an already-dead request. Race each OAuth setup await against the request signal via rejectOnAbort() (the shared series helper) and normalize an abort landing in the lookup to the shared abort contract. 2. Streamed output was not checked against cancellation before every yield. processEvent() emits buffered chunks for an event whose processing raced the abort, and the SSE handler yielded direct buffered results without an intervening check. The processEvent consumption loops now hand back the abort contract immediately (a break would let the event loop pull the stream once more before its top check sees the abort), and every reachable direct SSE yield is preceded by an abort check. The typed else-if branches of the SSE line loop are unreachable - coreHandledEventTypes captures every response.* type the delegate branch handles - so no guard is added there. Consumers of createMessage have no consumer-side guard, so the pre-yield checks are the defense on that path; completePrompt keeps its consumer guard and post-loop abort check for the quiet-end paths. --- src/api/providers/openai-codex.ts | 59 +++++++++++++++++++++++-- src/api/providers/utils/abort-signal.ts | 32 ++++++++++++++ 2 files changed, 87 insertions(+), 4 deletions(-) diff --git a/src/api/providers/openai-codex.ts b/src/api/providers/openai-codex.ts index ddf0393b71..3c1c7b0613 100644 --- a/src/api/providers/openai-codex.ts +++ b/src/api/providers/openai-codex.ts @@ -29,7 +29,13 @@ import { isMcpTool } from "../../utils/mcp-name" import { sanitizeOpenAiCallId } from "../../utils/tool-id" import { openAiCodexOAuthManager } from "../../integrations/openai-codex/oauth" import { t } from "../../i18n" -import { createAbortError, isRequestAborted, mergeAbortSignalAndTimeout, throwIfAborted } from "./utils/abort-signal" +import { + createAbortError, + isRequestAborted, + mergeAbortSignalAndTimeout, + rejectOnAbort, + throwIfAborted, +} from "./utils/abort-signal" export type OpenAiCodexModel = ReturnType @@ -251,7 +257,28 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion this.streamedToolCallIds.clear() // Get access token from OAuth manager - let accessToken = await openAiCodexOAuthManager.getAccessToken() + // The lookup can await credential loading, token refresh, and persistence: race it against + // the caller's signal so a cancelled request settles immediately instead of hanging on a + // lookup that no longer matters. + let accessToken: string | null + try { + if (abortSignal) { + accessToken = await rejectOnAbort( + openAiCodexOAuthManager.getAccessToken(), + abortSignal, + this.providerName, + ) + } else { + accessToken = await openAiCodexOAuthManager.getAccessToken() + } + } catch (error) { + // An abort landing while the lookup is pending must surface as the shared abort + // contract, not as an authentication failure. + if (isRequestAborted(error, abortSignal)) { + throw createAbortError(this.providerName) + } + throw error + } if (!accessToken) { throw new Error( t("common:errors.openAiCodex.notAuthenticated", { @@ -487,7 +514,14 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // is consistent across providers. try { // Get ChatGPT account ID for organization subscriptions - const accountId = await openAiCodexOAuthManager.getAccountId() + // The account lookup has the same gap as the token fetch: the request-local + // controller is already bridged to the caller's signal, so racing against it + // settles a cancelled request immediately instead of waiting on a dead lookup. + const accountId = await rejectOnAbort( + openAiCodexOAuthManager.getAccountId(), + abortController.signal, + this.providerName, + ) // Build Codex-specific headers. Authorization is provided by the SDK apiKey. const codexHeaders = this.buildCodexHeaders(model, effectiveSessionId, accountId) @@ -524,6 +558,11 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion this.sawSdkEventInCurrentResponse = true for await (const outChunk of this.processEvent(event, model)) { + // The top-of-loop check covers the wire; this one covers the buffered chunks + // processEvent still emits for an event whose processing raced the abort. A break would let + // the event loop pull the stream once more before its top check sees the abort, so hand + // back the abort contract immediately instead. + if (abortController.signal.aborted) throw createAbortError(this.providerName) if (outChunk.type === "text") { this.sawTextOutputInCurrentResponse = true } @@ -655,7 +694,9 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion const url = `${CODEX_API_BASE_URL}/responses` // Get ChatGPT account ID for organization subscriptions - const accountId = await openAiCodexOAuthManager.getAccountId() + // The signal here is the request-local controller the caller's abort is bridged into, so a + // cancellation lands in the lookup and settles it before the fallback fetch would start. + const accountId = await rejectOnAbort(openAiCodexOAuthManager.getAccountId(), abortSignal, this.providerName) // Build headers with required Codex-specific fields const headers: Record = { @@ -823,6 +864,9 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } for await (const outChunk of this.processEvent(parsed, model)) { + // A cancellation landing while processEvent is still emitting must not stream the + // event's remaining buffered chunks, so hand back the abort contract immediately. + if (abortSignal.aborted) throw createAbortError(this.providerName) if (outChunk.type === "text" || outChunk.type === "reasoning") { hasContent = true if (outChunk.type === "text") { @@ -842,6 +886,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion if (content.type === "text" && content.text) { hasContent = true this.sawTextOutputInCurrentResponse = true + if (abortSignal.aborted) break yield { type: "text", text: content.text } } } @@ -850,6 +895,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion for (const summary of outputItem.summary) { if (summary?.type === "summary_text" && typeof summary.text === "string") { hasContent = true + if (abortSignal.aborted) break yield { type: "reasoning", text: summary.text } } } @@ -858,6 +904,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion if (parsed.response.usage) { const usageData = this.normalizeUsage(parsed.response.usage, model) if (usageData) { + if (abortSignal.aborted) break yield usageData } } @@ -984,6 +1031,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } else if (parsed.choices?.[0]?.delta?.content) { hasContent = true this.sawTextOutputInCurrentResponse = true + if (abortSignal.aborted) break yield { type: "text", text: parsed.choices[0].delta.content } } else if ( parsed.item && @@ -992,10 +1040,12 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion ) { hasContent = true this.sawTextOutputInCurrentResponse = true + if (abortSignal.aborted) break yield { type: "text", text: parsed.item.text } } else if (parsed.usage) { const usageData = this.normalizeUsage(parsed.usage, model) if (usageData) { + if (abortSignal.aborted) break yield usageData } } @@ -1010,6 +1060,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion if (parsed.content || parsed.text || parsed.message) { hasContent = true this.sawTextOutputInCurrentResponse = true + if (abortSignal.aborted) break yield { type: "text", text: parsed.content || parsed.text || parsed.message } } } catch { diff --git a/src/api/providers/utils/abort-signal.ts b/src/api/providers/utils/abort-signal.ts index 26f57c3e9a..bd7d579b00 100644 --- a/src/api/providers/utils/abort-signal.ts +++ b/src/api/providers/utils/abort-signal.ts @@ -93,3 +93,35 @@ export function createAbortError(providerName: string): Error { abortError.name = "AbortError" return abortError } + +/** + * Await `pending` but reject with the provider's abort error when `signal` + * aborts first. For async phases that have no native signal support (model + * discovery) yet must still settle promptly on cancellation. The underlying + * promise keeps running (its settlement is ignored) — cancellation is + * cooperative at this boundary. + * + * The abort listener is detached once `pending` settles (success or + * failure), so repeated calls on one signal do not accumulate listeners. + */ +export function rejectOnAbort(pending: Promise, signal: AbortSignal, providerName: string): Promise { + if (signal.aborted) { + return Promise.reject(createAbortError(providerName)) + } + + return new Promise((resolve, reject) => { + const onAbort = () => reject(createAbortError(providerName)) + // Stryker disable next-line ObjectLiteral,BooleanLiteral: a signal fires its abort event exactly once and the settle handler removes this listener, so the once flag is unobservable + signal.addEventListener("abort", onAbort, { once: true }) + void pending.then( + (value) => { + signal.removeEventListener("abort", onAbort) + resolve(value) + }, + (error) => { + signal.removeEventListener("abort", onAbort) + reject(error) + }, + ) + }) +} From 2ce9e64d4ede4b42fd003c7604873825c720f63e Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Wed, 23 Sep 2026 14:26:09 +0800 Subject: [PATCH 10/12] test(api): cover openai-codex cancellation races and guarded yields - regressions that keep the token or account lookup pending and assert the caller abort rejects before the lookup resolves, for the executeRequest and makeCodexRequest lookups - a regression that the completePrompt timeout rejects with the shared abort contract while the token lookup is still pending - structural tests that the buffered chunks of an event whose processing aborts the request never reach the caller, on both the SDK and the SSE transports (createMessage, which has no consumer-side guard) - a pre-aborted signal now settles in the OAuth race before the SDK is reached, so the existing abort-contract tests assert the SDK is never called - an SSE content test pinning every direct chunk shape the fallback emits, in order --- .../openai-codex-native-tool-calls.spec.ts | 14 +- .../providers/__tests__/openai-codex.spec.ts | 302 +++++++++++++++++- .../utils/__tests__/abort-signal.spec.ts | 98 ++++++ 3 files changed, 392 insertions(+), 22 deletions(-) diff --git a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts index f831d22fc1..a6d841877b 100644 --- a/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts +++ b/src/api/providers/__tests__/openai-codex-native-tool-calls.spec.ts @@ -625,15 +625,11 @@ describe("OpenAiCodexHandler native tool calls", () => { abortSignal: controller.signal, }) - // Consume the stream to trigger the request; the request-local controller is already - // aborted, so the loop must break before the first event is processed. - const chunks = await collectStream(stream) - expect(chunks).toEqual([]) - - expect(mockCreate).toHaveBeenCalled() - const createCallArgs = mockCreate.mock.calls[0][1] as { signal?: AbortSignal } - // The internal signal should already be aborted since the external one was pre-aborted - expect(createCallArgs.signal?.aborted).toBe(true) + // A pre-aborted signal settles in the OAuth race before the SDK is reached: the abort + // contract must win over the request setup, so the stream rejects and the SDK is never + // called with a signal that is already aborted. + await expect(collectStream(stream)).rejects.toMatchObject({ name: "AbortError" }) + expect(mockCreate).not.toHaveBeenCalled() }) }) }) diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index 63c72403ae..cddbea11bf 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -517,13 +517,12 @@ describe("OpenAiCodexHandler.completePrompt streaming", () => { expect(getAccountId).not.toHaveBeenCalled() }) - // The abort lands in the only window the consumer guard can still act: event2's `delta` - // getter fires it while processEvent builds the chunk - after executeRequest's - // top-of-loop check has passed, before completePrompt's own. The post-loop check throws - // the same AbortError whether or not the guard breaks, so what distinguishes the guarded - // loop from a mutated one is the pull count: the guard stops pulling after the chunk that - // the abort rode in on, a loop that keeps running pulls the SDK stream a third time. - it("breaks the streaming consumer loop at the top-of-loop guard and stops pulling", async () => { + // The abort lands while event2 is being processed: its `delta` getter fires it after + // executeRequest's top-of-loop check has passed. The pre-yield guard must end the request + // there, so the stream is never pulled a third time - the pull count distinguishes a flow + // that kept pulling (a mutated guard that lets the loop's continuation run) from the + // guarded one, and the buffered "post-abort" chunk must never be joined. + it("settles with the abort contract and stops pulling once an event's processing aborts the request", async () => { const handler = createHandler() const controller = new AbortController() let sdkPulls = 0 @@ -556,9 +555,9 @@ describe("OpenAiCodexHandler.completePrompt streaming", () => { name: "AbortError", }) - // event1 is pulled and joined; event2 is pulled (its getter fires the abort) but the - // guard breaks before it is joined - a third pull only happens when the mutated - // top-of-loop check keeps the loop running. + // event1 is pulled and joined; event2 is pulled (its getter fires the abort) and the + // pre-yield guard ends the request before its buffered chunk is joined - a third pull + // only happens when a mutated guard lets the event loop's continuation run. expect(sdkPulls).toBe(2) }) @@ -1227,6 +1226,28 @@ describe("OpenAiCodexHandler.completePrompt timeout", () => { expect(signalDuringRequest!.aborted).toBe(true) expect(mockFetch).not.toHaveBeenCalled() }) + + it("rejects with the abort contract when the timeout fires while the token lookup is pending", async () => { + const handler = createHandler() + let resolveToken!: (token: string) => void + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockReturnValue( + new Promise((resolve) => { + resolveToken = resolve + }), + ) + const create = vitest.fn() + Reflect.set(handler, "client", { responses: { create } }) + + // No abort signal: the timeout is the only cancellation source, so the timeout path must + // race the OAuth lookup itself. + await expect(handler.completePrompt("test prompt", { timeoutMs: 20 })).rejects.toMatchObject({ + name: "AbortError", + }) + + // The lookup settling afterwards must be ignored: the request is already gone. + resolveToken("late-token") + expect(create).not.toHaveBeenCalled() + }) }) describe("OpenAiCodexHandler.createMessage abort bridging", () => { @@ -1251,8 +1272,8 @@ describe("OpenAiCodexHandler.createMessage abort bridging", () => { const refresh = vitest .spyOn(openAiCodexOAuthManager, "forceRefreshAccessToken") .mockResolvedValue("refreshed-token") - // The SDK fails with exactly the auth-failure wording the retry path would act on, so the - // abort check must win over the refresh-and-retry logic. + // The SDK is wired to fail with exactly the auth-failure wording the retry path would act + // on: the cancellation must win over the refresh-and-retry logic at every level. const create = vitest.fn().mockRejectedValue(new Error("401 invalid token")) Reflect.set(handler, "client", { responses: { create } }) const mockFetch = vitest.fn() @@ -1267,12 +1288,267 @@ describe("OpenAiCodexHandler.createMessage abort bridging", () => { ), ).rejects.toMatchObject({ name: "AbortError" }) - // No refresh, no second SDK attempt, no SSE fallback. + // A pre-aborted signal settles in the OAuth race before the SDK is even reached, so there + // is no SDK attempt, no refresh, and no SSE fallback. expect(refresh).not.toHaveBeenCalled() + expect(create).not.toHaveBeenCalled() + expect(mockFetch).not.toHaveBeenCalled() + }) + + it("settles with the abort contract before a pending token lookup resolves", async () => { + const handler = createHandler() + let resolveToken!: (token: string) => void + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockReturnValue( + new Promise((resolve) => { + resolveToken = resolve + }), + ) + const create = vitest.fn() + Reflect.set(handler, "client", { responses: { create } }) + + const controller = new AbortController() + const stream = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + // The first pull runs the generator up to the token lookup, which is still pending. + const firstPull = stream.next() + controller.abort() + + await expect(firstPull).rejects.toMatchObject({ name: "AbortError" }) + + // The lookup settling afterwards must be ignored: the request is already gone. + resolveToken("late-token") + await new Promise((resolve) => setTimeout(resolve, 0)) + expect(create).not.toHaveBeenCalled() + }) + + it("settles with the abort contract before a pending account lookup resolves", async () => { + const handler = createHandler() + let resolveAccount!: (id: string) => void + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockReturnValue( + new Promise((resolve) => { + resolveAccount = resolve + }), + ) + const create = vitest.fn() + Reflect.set(handler, "client", { responses: { create } }) + + const controller = new AbortController() + const stream = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + const firstPull = stream.next() + // Let the token lookup settle so the generator reaches the account lookup. + await new Promise((resolve) => setTimeout(resolve, 0)) + controller.abort() + + await expect(firstPull).rejects.toMatchObject({ name: "AbortError" }) + + resolveAccount("late-account") + await new Promise((resolve) => setTimeout(resolve, 0)) + expect(create).not.toHaveBeenCalled() + }) + + it("settles with the abort contract before a pending account lookup resolves on the fallback", async () => { + const handler = createHandler() + let resolveAccount!: (id: string) => void + // The executeRequest-level lookup settles; the fallback's own lookup stays pending. + vitest + .spyOn(openAiCodexOAuthManager, "getAccountId") + .mockImplementationOnce(() => Promise.resolve("acct_sdk")) + .mockImplementationOnce( + () => + new Promise((resolve) => { + resolveAccount = resolve + }), + ) + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + const controller = new AbortController() + const stream = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + const firstPull = stream.next() + // Let the SDK fail and the fallback reach its own account lookup. + await new Promise((resolve) => setTimeout(resolve, 0)) + controller.abort() + + await expect(firstPull).rejects.toMatchObject({ name: "AbortError" }) + + resolveAccount("late-account") + await new Promise((resolve) => setTimeout(resolve, 0)) expect(create).toHaveBeenCalledTimes(1) expect(mockFetch).not.toHaveBeenCalled() }) + it("never yields the buffered chunks of an event whose processing aborts the request", async () => { + const handler = createHandler() + const controller = new AbortController() + const event1 = { type: "response.output_text.delta", delta: "pre-abort" } + const event2 = { + type: "response.output_text.delta", + get delta() { + controller.abort() + return "post-abort" + }, + } + const create = vitest.fn().mockImplementation(() => { + return Promise.resolve({ + [Symbol.asyncIterator]() { + let pulls = 0 + return { + next: async () => { + pulls++ + return { value: pulls === 1 ? event1 : event2, done: false } + }, + return: async () => ({ value: undefined, done: true }), + } + }, + }) + }) + Reflect.set(handler, "client", { responses: { create } }) + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + // event1's chunk streams; the abort lands while event2 is being processed, so its buffered + // "post-abort" chunk must not reach the caller - createMessage has no consumer-side guard. + expect(await iter.next()).toMatchObject({ value: { type: "text", text: "pre-abort" } }) + await expect(iter.next()).rejects.toMatchObject({ name: "AbortError" }) + expect(mockFetch).not.toHaveBeenCalled() + }) + + it("never yields a buffered SSE line after a cancellation that lands between lines", async () => { + const handler = createHandler() + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const controller = new AbortController() + const encoder = new TextEncoder() + const body = new ReadableStream({ + start(streamController) { + // A single chunk: both lines land in one read so the second is buffered in the + // handler's line loop when the cancellation lands. + streamController.enqueue( + encoder.encode( + 'data: {"type":"response.output_text.delta","delta":"one"}\n\ndata: {"type":"response.output_text.delta","delta":"two"}\n\n', + ), + ) + streamController.close() + }, + }) + const mockFetch = vitest.fn().mockResolvedValue({ ok: true, body }) + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + expect(await iter.next()).toMatchObject({ value: { type: "text", text: "one" } }) + controller.abort() + // The second line was enqueued before the cancellation but must not stream: the delegate + // guard hands back the abort contract instead of yielding its chunk. + await expect(iter.next()).rejects.toMatchObject({ name: "AbortError" }) + }) + + it("never yields the remaining complete-response chunks once the caller cancels", async () => { + const handler = createHandler() + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const controller = new AbortController() + const encoder = new TextEncoder() + const complete = { + response: { + id: "resp_test", + output: [ + { type: "text", content: [{ type: "text", text: "first" }] }, + { type: "text", content: [{ type: "text", text: "second" }] }, + ], + }, + } + const body = new ReadableStream({ + start(streamController) { + streamController.enqueue(encoder.encode(`data: ${JSON.stringify(complete)}\n\n`)) + streamController.close() + }, + }) + const mockFetch = vitest.fn().mockResolvedValue({ ok: true, body }) + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + expect(await iter.next()).toMatchObject({ value: { type: "text", text: "first" } }) + controller.abort() + // The second text item was buffered before the cancellation but must not stream: the + // guard before the yield breaks out of the item loop, and the while-top guard ends the read. + const chunks = await collectStream(iter) + expect(chunks).toEqual([]) + }) + + it("streams every direct SSE chunk shape the fallback supports", async () => { + const handler = createHandler() + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const encoder = new TextEncoder() + const complete = { + response: { + output: [ + { type: "text", content: [{ type: "text", text: "complete-text" }] }, + { type: "reasoning", summary: [{ type: "summary_text", text: "complete-summary" }] }, + ], + usage: { input_tokens: 7, output_tokens: 11 }, + }, + } + const lines = [ + `data: ${JSON.stringify(complete)}`, + 'data: {"type":"response.output_text.delta","delta":"delegated"}', + 'data: {"choices":[{"delta":{"content":"choices-text"}}]}', + 'data: {"item":{"type":"text","text":"item-text"}}', + 'data: {"usage":{"input_tokens":3,"output_tokens":4}}', + '{"content":"plain-text"}', + ] + const body = new ReadableStream({ + start(streamController) { + for (const line of lines) { + streamController.enqueue(encoder.encode(`${line}\n\n`)) + } + streamController.close() + }, + }) + const mockFetch = vitest.fn().mockResolvedValue({ ok: true, body }) + vitest.stubGlobal("fetch", mockFetch) + + const chunks = await collectStream( + handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + }), + ) + + // One stream, every direct SSE shape: the complete-response block (text, reasoning, + // usage), the delegated delta, the legacy choices/item/usage shapes, and a plain JSON + // line. The sequence also pins the order the fallback emits them in. + expect(chunks).toEqual([ + { type: "text", text: "complete-text" }, + { type: "reasoning", text: "complete-summary" }, + { type: "usage", inputTokens: 7, outputTokens: 11, cacheWriteTokens: 0, cacheReadTokens: 0, totalCost: 0 }, + { type: "text", text: "delegated" }, + { type: "text", text: "choices-text" }, + { type: "text", text: "item-text" }, + { type: "usage", inputTokens: 3, outputTokens: 4, cacheWriteTokens: 0, cacheReadTokens: 0, totalCost: 0 }, + { type: "text", text: "plain-text" }, + ]) + }) + it("emits nothing once the caller cancels before the first SDK event", async () => { const handler = createHandler() const controller = new AbortController() diff --git a/src/api/providers/utils/__tests__/abort-signal.spec.ts b/src/api/providers/utils/__tests__/abort-signal.spec.ts index 1692f71e63..7ca85d7782 100644 --- a/src/api/providers/utils/__tests__/abort-signal.spec.ts +++ b/src/api/providers/utils/__tests__/abort-signal.spec.ts @@ -3,9 +3,107 @@ import { isRequestAborted, mergeAbortSignalAndTimeout, mergeAbortSignals, + rejectOnAbort, throwIfAborted, } from "../abort-signal" +describe("rejectOnAbort", () => { + it("resolves with the pending value when it settles before the signal aborts", async () => { + const controller = new AbortController() + + await expect(rejectOnAbort(Promise.resolve("done"), controller.signal, "TestProvider")).resolves.toBe("done") + expect(controller.signal.aborted).toBe(false) + }) + + it("rejects with the provider abort error when the signal aborts first", async () => { + const controller = new AbortController() + // Never settles: the race must end purely via the abort. + const pending = new Promise(() => {}) + const race = rejectOnAbort(pending, controller.signal, "TestProvider") + controller.abort() + + // Bound the wait so a broken race fails this test fast (and fails the + // Stryker mutant) instead of hanging. + await expect( + new Promise((_resolve, reject) => { + const deadline = setTimeout(() => reject(new Error("race deadline exceeded")), 3000) + void race.then( + (value) => { + clearTimeout(deadline) + reject(new Error(`race unexpectedly resolved: ${String(value)}`)) + }, + (error) => { + clearTimeout(deadline) + reject(error) + }, + ) + }), + ).rejects.toMatchObject({ + name: "AbortError", + message: "The TestProvider request was aborted", + }) + }) + + it("rejects immediately when the signal is already aborted", async () => { + const controller = new AbortController() + controller.abort() + const pending = new Promise(() => {}) + + await expect(rejectOnAbort(pending, controller.signal, "TestProvider")).rejects.toMatchObject({ + name: "AbortError", + message: "The TestProvider request was aborted", + }) + }) + + it("propagates the pending rejection when the signal stays active", async () => { + const controller = new AbortController() + const boom = new Error("lookup failed") + + await expect(rejectOnAbort(Promise.reject(boom), controller.signal, "TestProvider")).rejects.toBe(boom) + }) + + it("detaches the abort listener once the pending settles", async () => { + const controller = new AbortController() + const addSpy = vi.spyOn(controller.signal, "addEventListener") + const removeSpy = vi.spyOn(controller.signal, "removeEventListener") + + await expect(rejectOnAbort(Promise.resolve("done"), controller.signal, "TestProvider")).resolves.toBe("done") + + // The settle path must remove the exact listener that was registered, not just any + // function: removing a different reference would leave the original abort listener + // attached to the signal. First require the registration to have happened at all, + // so a missing registration cannot silently degrade to an undefined comparison. + expect(addSpy).toHaveBeenCalledTimes(1) + const registeredListener = addSpy.mock.calls[0]?.[1] as EventListener | undefined + expect(typeof registeredListener).toBe("function") + expect(removeSpy).toHaveBeenCalledWith("abort", registeredListener) + addSpy.mockRestore() + removeSpy.mockRestore() + }) + + it("detaches the abort listener when the pending rejects", async () => { + const controller = new AbortController() + const addSpy = vi.spyOn(controller.signal, "addEventListener") + const removeSpy = vi.spyOn(controller.signal, "removeEventListener") + const lookupError = new Error("lookup failed") + + await expect(rejectOnAbort(Promise.reject(lookupError), controller.signal, "TestProvider")).rejects.toBe( + lookupError, + ) + + // The settle path must remove the exact listener that was registered, not just any + // function: removing a different reference would leave the original abort listener + // attached to the signal. First require the registration to have happened at all, + // so a missing registration cannot silently degrade to an undefined comparison. + expect(addSpy).toHaveBeenCalledTimes(1) + const registeredListener = addSpy.mock.calls[0]?.[1] as EventListener | undefined + expect(typeof registeredListener).toBe("function") + expect(removeSpy).toHaveBeenCalledWith("abort", registeredListener) + addSpy.mockRestore() + removeSpy.mockRestore() + }) +}) + describe("abort-signal utilities", () => { describe("mergeAbortSignalAndTimeout", () => { it("returns undefined when no signal or positive timeout is provided", () => { From bf39b3cb905d53e955bcc306d8b95d83e630ae24 Mon Sep 17 00:00:00 2001 From: Eason Liang Date: Wed, 23 Sep 2026 14:40:15 +0800 Subject: [PATCH 11/12] test(api): pin the remaining openai-codex abort guards to killable mutants Closes out the local mutation preflight on the cancellation work: - a regression that a non-abort token failure propagates as-is (pins the isRequestAborted branch of the getAccessToken catch) - a regression that an SDK stream failing at the moment the caller cancels settles via the retry-catch abort check, not the auth-refresh path - a regression that a signal aborting as the token lookup settles is observed by the request-local bridge - table-driven SSE regressions that the pre-yield guard of each reachable direct chunk shape (complete-response reasoning/usage, legacy choices/item/usage, plain JSON) breaks before the buffered chunk streams once the caller cancels Two ConditionalExpression mutants are now unobservable and are directive-excluded with concrete reasons: the isRequestAborted re-throw in the getAccessToken catch (the race's rejection is byte-identical to the normalized error) and the completePrompt consumer guard (every yield site now checks the signal before yielding, and the check and the yield are synchronous). --- .../providers/__tests__/openai-codex.spec.ts | 149 ++++++++++++++++++ src/api/providers/openai-codex.ts | 2 + 2 files changed, 151 insertions(+) diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index cddbea11bf..d4a22663dc 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -1549,6 +1549,155 @@ describe("OpenAiCodexHandler.createMessage abort bridging", () => { ]) }) + it("propagates a non-abort token failure as-is when no cancellation is involved", async () => { + const handler = createHandler() + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockRejectedValue(new Error("oauth down")) + const create = vitest.fn() + Reflect.set(handler, "client", { responses: { create } }) + + // No abort signal: a genuine lookup failure must reach the caller unchanged, not be + // normalized into the abort contract. + await expect( + collectStream( + handler.createMessage("System", [{ role: "user", content: "Hello" }], { taskId: "task-test" }), + ), + ).rejects.toThrow("oauth down") + + expect(create).not.toHaveBeenCalled() + }) + + it("settles with the abort contract when the SDK stream fails at the moment the caller cancels", async () => { + const handler = createHandler() + const refresh = vitest + .spyOn(openAiCodexOAuthManager, "forceRefreshAccessToken") + .mockResolvedValue("refreshed-token") + const controller = new AbortController() + const event1 = { type: "response.output_text.delta", delta: "pre-abort" } + const create = vitest.fn().mockImplementation(() => { + return Promise.resolve({ + [Symbol.asyncIterator]() { + let pulls = 0 + return { + next: async () => { + pulls++ + if (pulls === 1) { + return { value: event1, done: false } + } + // The cancellation lands while the second pull is in flight: the top-of-loop + // check has already passed, so the failure must settle via the retry-catch + // abort check, not the auth-refresh path. + controller.abort() + throw new Error("401 invalid token") + }, + return: async () => ({ value: undefined, done: true }), + } + }, + }) + }) + Reflect.set(handler, "client", { responses: { create } }) + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + expect(await iter.next()).toMatchObject({ value: { type: "text", text: "pre-abort" } }) + await expect(iter.next()).rejects.toMatchObject({ name: "AbortError" }) + + // The cancellation wins over the auth-failure wording: no refresh, no SSE fallback. + expect(refresh).not.toHaveBeenCalled() + expect(create).toHaveBeenCalledTimes(1) + expect(mockFetch).not.toHaveBeenCalled() + }) + + it("settles with the abort contract when the signal aborts as the token lookup settles", async () => { + const handler = createHandler() + const controller = new AbortController() + const create = vitest.fn() + Reflect.set(handler, "client", { responses: { create } }) + const mockFetch = vitest.fn() + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + // The first pull runs the generator up to the token race. The token is already resolved, so + // the race settles on the next microtask and detaches its abort listener; queueing the abort + // after that microtask lands it in the window where only the request-local bridge can still + // see it. + const firstPull = iter.next() + queueMicrotask(() => controller.abort()) + + await expect(firstPull).rejects.toMatchObject({ name: "AbortError" }) + + expect(create).not.toHaveBeenCalled() + expect(mockFetch).not.toHaveBeenCalled() + }) + + // Each case streams one line before the cancellation and buffers the target line in the + // handler's line loop: the target's pre-yield guard must break before its chunk streams. + const sseData = (obj: unknown) => "data: " + JSON.stringify(obj) + const completeTextLine = (t: string) => + sseData({ response: { output: [{ type: "text", content: [{ type: "text", text: t }] }] } }) + const completeReasoningLine = (t: string) => + sseData({ response: { output: [{ type: "reasoning", summary: [{ type: "summary_text", text: t }] }] } }) + const completeUsageLine = (t: string) => + sseData({ + response: { + output: [{ type: "text", content: [{ type: "text", text: t }] }], + usage: { input_tokens: 1, output_tokens: 1 }, + }, + }) + const choicesLine = (t: string) => sseData({ choices: [{ delta: { content: t } }] }) + const itemLine = (t: string) => sseData({ item: { type: "text", text: t } }) + const usageLine = sseData({ usage: { input_tokens: 1, output_tokens: 1 } }) + const plainLine = (t: string) => JSON.stringify({ content: t }) + + const sseGuardCases: [string, string, string, Record][] = [ + [ + "a complete-response reasoning summary", + completeTextLine("a"), + completeReasoningLine("b"), + { type: "text", text: "a" }, + ], + ["a complete-response usage chunk", completeTextLine("a"), completeUsageLine("b"), { type: "text", text: "a" }], + ["a legacy choices delta", itemLine("a"), choicesLine("b"), { type: "text", text: "a" }], + ["a legacy item text", choicesLine("a"), itemLine("b"), { type: "text", text: "a" }], + ["a legacy usage object", itemLine("a"), usageLine, { type: "text", text: "a" }], + ["a plain JSON line", usageLine, plainLine("b"), { type: "usage", inputTokens: 1, outputTokens: 1 }], + ] + for (const [name, firstLine, targetLine, firstChunk] of sseGuardCases) { + it(`never yields ${name} once the caller cancels before the handler reads it`, async () => { + const handler = createHandler() + const create = vitest.fn().mockRejectedValue(new Error("sdk down")) + Reflect.set(handler, "client", { responses: { create } }) + const controller = new AbortController() + const encoder = new TextEncoder() + const body = new ReadableStream({ + start(streamController) { + // A single chunk: both lines land in one read so the target is buffered in the + // handler's line loop when the cancellation lands. + streamController.enqueue(encoder.encode(`${firstLine}\n\n${targetLine}\n\n`)) + streamController.close() + }, + }) + const mockFetch = vitest.fn().mockResolvedValue({ ok: true, body }) + vitest.stubGlobal("fetch", mockFetch) + + const iter = handler.createMessage("System", [{ role: "user", content: "Hello" }], { + taskId: "task-test", + abortSignal: controller.signal, + }) + expect(await iter.next()).toMatchObject({ value: firstChunk }) + controller.abort() + // The target line was buffered before the cancellation, but its pre-yield guard + // breaks out of the line loop and the while-top guard ends the read. + expect(await collectStream(iter)).toEqual([]) + }) + } + it("emits nothing once the caller cancels before the first SDK event", async () => { const handler = createHandler() const controller = new AbortController() diff --git a/src/api/providers/openai-codex.ts b/src/api/providers/openai-codex.ts index 3c1c7b0613..8d71ab7a4d 100644 --- a/src/api/providers/openai-codex.ts +++ b/src/api/providers/openai-codex.ts @@ -274,6 +274,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion } catch (error) { // An abort landing while the lookup is pending must surface as the shared abort // contract, not as an authentication failure. + // Stryker disable next-line ConditionalExpression: the false mutant differs from the original only if an abort lands between the OAuth race's settle and this catch — a microtask window (the abort event is a microtask) no test can schedule deterministically, because the test's own microtasks queue after the settle that precedes the window; the observable contract (a genuine failure with no cancellation propagates as-is, an abort surfaces as the shared contract) is pinned by the adjacent regressions and the true mutant is killed by the non-abort-failure regression. if (isRequestAborted(error, abortSignal)) { throw createAbortError(this.providerName) } @@ -1450,6 +1451,7 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // A buffered chunk can still be pulled in the window between the abort and the // inner generator's own stop, so break here: post-abort output must never be // joined into the completion. + // Stryker disable next-line ConditionalExpression: the false mutant differs from the original only if an abort lands between the transport's last pre-yield check and this consumer's resumption — the yield-hop microtask window no test can schedule deterministically (the abort event is a microtask; the test's own microtasks queue after the chunk's yield); the guard is retained for that production window, the abort regressions pin its pre-abort contract, and the true mutant is killed by the multi-chunk happy paths. if (requestSignal?.aborted) { break } From 642443ccfef27c5fd27aabf7be9a324cfb9d6e5f Mon Sep 17 00:00:00 2001 From: easonLiangWorldedtech Date: Mon, 28 Sep 2026 17:42:16 +0800 Subject: [PATCH 12/12] fix(api): race the auth retry refresh against the caller's abort signal When the first Codex request fails with an authentication error, the retry forced a token refresh and awaited it unguarded. A cancellation landing while the OAuth operation was in flight left createMessage suspended until the refresh settled and then surfaced the authentication failure of a request that no longer exists, instead of the shared abort contract. Wrap the refresh in rejectOnAbort when the request carries a signal (the optional-signal path keeps the plain await), so an abort mid-refresh settles the request with the provider's AbortError. The regression test lets the refresh settle with no token after the caller aborts and asserts the stream rejects with AbortError instead of the authentication error. --- .../providers/__tests__/openai-codex.spec.ts | 42 +++++++++++++++++++ src/api/providers/openai-codex.ts | 11 ++++- 2 files changed, 51 insertions(+), 2 deletions(-) diff --git a/src/api/providers/__tests__/openai-codex.spec.ts b/src/api/providers/__tests__/openai-codex.spec.ts index d4a22663dc..845e341681 100644 --- a/src/api/providers/__tests__/openai-codex.spec.ts +++ b/src/api/providers/__tests__/openai-codex.spec.ts @@ -1054,6 +1054,48 @@ describe("OpenAiCodexHandler Responses Lite requests", () => { }) }) + it("settles with the abort contract when cancellation lands during the auth retry refresh", async () => { + // The first request fails authentication and the forced token refresh is in flight + // when the caller aborts. The refresh must be raced against the signal: without the + // race the handler stays suspended until the refresh settles and then surfaces the + // authentication failure of a request that no longer exists instead of the abort. + const handler = new OpenAiCodexHandler({ apiModelId: "gpt-5.6-luna" }) + vitest.spyOn(openAiCodexOAuthManager, "getAccessToken").mockResolvedValue("expired-token") + vitest.spyOn(openAiCodexOAuthManager, "getAccountId").mockResolvedValue("acct_test") + // Object.assign sidesteps the SDK client's structural type so the double stays cast-free. + Object.assign(handler, { + client: { responses: { create: vitest.fn().mockRejectedValue(new Error("SDK unavailable")) } }, + }) + // The OAuth operation settles (with no token) after the test aborts, mirroring a + // slow refresh that loses the race. + const refresh = vitest + .spyOn(openAiCodexOAuthManager, "forceRefreshAccessToken") + .mockImplementation(() => new Promise((resolve) => setTimeout(() => resolve(null), 50))) + const mockFetch = vitest.fn().mockResolvedValueOnce({ + ok: false, + status: 401, + text: vitest.fn().mockResolvedValue('{"error":{"message":"Codex API invalid token"}}'), + }) + vitest.stubGlobal("fetch", mockFetch) + + const abortController = new AbortController() + const result = collectStream( + handler.createMessage("Instructions", [{ role: "user", content: "Abort during refresh" }], { + taskId: "task-abort-refresh", + tools: [], + abortSignal: abortController.signal, + }), + ) + // Let the first request fail and the refresh start, then cancel the caller. + await new Promise((resolve) => setTimeout(resolve, 10)) + abortController.abort() + + await expect(result).rejects.toMatchObject({ name: "AbortError" }) + expect(refresh).toHaveBeenCalledTimes(1) + // No retry request goes out for a request that was aborted mid-refresh. + expect(mockFetch).toHaveBeenCalledTimes(1) + }) + it.each(["gpt-5.5", "gpt-5.6-sol", "gpt-5.6-terra", "gpt-5.6-luna-alias"])( "does not apply Luna behavior to %s", async (apiModelId) => { diff --git a/src/api/providers/openai-codex.ts b/src/api/providers/openai-codex.ts index 8d71ab7a4d..470efb8a9b 100644 --- a/src/api/providers/openai-codex.ts +++ b/src/api/providers/openai-codex.ts @@ -333,8 +333,15 @@ export class OpenAiCodexHandler extends BaseProvider implements SingleCompletion // the service has accepted the request, so a refreshed-token retry would replay it // and append a second generation to output the caller already has. if (attempt === 0 && isAuthFailure && !this.sawSdkEventInCurrentResponse) { - // Force refresh the token for retry - const refreshed = await openAiCodexOAuthManager.forceRefreshAccessToken() + // Force refresh the token for retry. Race the refresh against the caller's + // signal: an abort landing while the OAuth operation is pending must settle + // the request with the shared abort contract instead of leaving createMessage + // suspended until the refresh settles and then surfacing the authentication + // failure of a request that no longer exists. + const refreshPromise = openAiCodexOAuthManager.forceRefreshAccessToken() + const refreshed = abortSignal + ? await rejectOnAbort(refreshPromise, abortSignal, this.providerName) + : await refreshPromise if (!refreshed) { throw new Error( t("common:errors.openAiCodex.notAuthenticated", {