diff --git a/packages/ai/src/transports/openai-responses-contracts.ts b/packages/ai/src/transports/openai-responses-contracts.ts index 1d16e65f0c3c..8c8fca623dab 100644 --- a/packages/ai/src/transports/openai-responses-contracts.ts +++ b/packages/ai/src/transports/openai-responses-contracts.ts @@ -1,5 +1,4 @@ import { - PROVIDER_FAILURE_WITH_OUTPUT_ERROR_CODE, PROVIDER_POST_DISPATCH_AMBIGUITY_ERROR_CODE, type Api, type ProviderReplayState, @@ -37,18 +36,6 @@ export const OPENAI_RESPONSES_APIS: ReadonlySet = new Set([ "openclaw-azure-openai-responses-transport", ]); -export class OpenAIResponsesWebSocketResponseFailedError extends Error { - readonly code: string; - - constructor(hasOutput: boolean) { - super("OpenAI Responses WebSocket returned response.failed"); - this.name = "OpenAIResponsesWebSocketResponseFailedError"; - this.code = hasOutput - ? PROVIDER_FAILURE_WITH_OUTPUT_ERROR_CODE - : PROVIDER_POST_DISPATCH_AMBIGUITY_ERROR_CODE; - } -} - export class OpenAIResponsesWebSocketPreDispatchError extends Error { constructor(cause: unknown) { super("OpenAI Responses WebSocket failed before request dispatch", { cause }); diff --git a/packages/ai/src/transports/openai-responses-websocket-client.test.ts b/packages/ai/src/transports/openai-responses-websocket-client.test.ts index 8a82c1ce60ff..a1a61f687656 100644 --- a/packages/ai/src/transports/openai-responses-websocket-client.test.ts +++ b/packages/ai/src/transports/openai-responses-websocket-client.test.ts @@ -1,5 +1,4 @@ import { - PROVIDER_FAILURE_WITH_OUTPUT_ERROR_CODE, PROVIDER_POST_DISPATCH_AMBIGUITY_ERROR_CODE, type AssistantMessage, type Context, @@ -7,6 +6,7 @@ import { } from "@openclaw/llm-core"; import { WebSocketError } from "openai/resources/responses/internal-base.js"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { isRetryableAssistantError } from "../../../../src/llm/utils/retry.js"; import { configureAiTransportHost, getAiTransportHost } from "../host.js"; import { cleanupSessionResources } from "../session-resources.js"; import { @@ -246,9 +246,13 @@ function toolCallResponse(responseId: string): StreamMessage[] { } function sdkCompletion(responseId: string): SdkResponse { + return sdkEvent(completedEvent(responseId)); +} + +function sdkEvent(event: Record): SdkResponse { return { data: (async function* () { - yield completedEvent(responseId); + yield event; })(), response: new Response(null, { status: 200 }), }; @@ -727,45 +731,86 @@ describe("native OpenAI Responses WebSocket client integration", () => { expect(transportState.sdkRequests).toHaveLength(1); }); - it("does not replay after an explicit failed response with no output", async () => { - transportState.responseBatches.push([ - message({ type: "response.created", response: { id: "resp_failed" } }), - message({ - type: "response.failed", - response: { id: "resp_failed", status: "failed", output: [] }, - }), - ]); - const result = await run({ messages: [userMessage("hello", 1)], tools: [] }); + it("preserves failed terminal semantics across WebSocket and SSE without degradation", async () => { + const failedEvent = { + type: "response.failed", + response: { + id: "resp_failed", + status: "failed", + model: "gpt-5.6-luna-2026-08-01", + service_tier: "priority", + error: { code: "server_error", message: "503 temporary provider response" }, + output: [], + usage: { + input_tokens: 21, + output_tokens: 4, + total_tokens: 25, + input_tokens_details: { cached_tokens: 6, cache_write_tokens: 2 }, + output_tokens_details: { reasoning_tokens: 3 }, + }, + }, + }; + const pricedModel = { + ...model, + cost: { input: 5, output: 30, cacheRead: 0.5, cacheWrite: 6.25 }, + } satisfies Model<"openai-responses">; + transportState.responseBatches.push( + [message(failedEvent)], + [message(completedEvent("resp_next"))], + ); + transportState.sdkOutcomes.push(sdkEvent(failedEvent)); - expect(result.stopReason).toBe("error"); - expect(result.errorCode).toBe(PROVIDER_POST_DISPATCH_AMBIGUITY_ERROR_CODE); - expect(transportState.websocketRequests).toHaveLength(1); - expect(transportState.sdkRequests).toEqual([]); + const websocket = await run( + { messages: [userMessage("websocket", 1)], tools: [] }, + { model: pricedModel }, + ); + const sse = await run( + { messages: [userMessage("sse", 2)], tools: [] }, + { model: pricedModel, transport: "sse" }, + ); - transportState.sdkOutcomes.push(sdkCompletion("resp_sse")); - const next = await run({ messages: [userMessage("next", 2)], tools: [] }); + const terminalFacts = { + stopReason: "error", + errorMessage: "server_error: 503 temporary provider response", + responseId: "resp_failed", + responseModel: "gpt-5.6-luna-2026-08-01", + usage: { + input: 13, + output: 4, + cacheRead: 6, + cacheWrite: 2, + reasoningTokens: 3, + totalTokens: 25, + }, + }; + expect(websocket).toMatchObject(terminalFacts); + expect(sse).toMatchObject(terminalFacts); + expect(websocket.usage.cost.total).toBeCloseTo(0.000401, 10); + expect(sse.usage.cost.total).toBeCloseTo(0.000401, 10); + expect(isRetryableAssistantError(websocket)).toBe(true); + expect(isRetryableAssistantError(sse)).toBe(true); + + const next = await run( + { messages: [userMessage("next", 3)], tools: [] }, + { model: pricedModel }, + ); expect(next.stopReason).toBe("stop"); - expect(transportState.websocketOptions).toHaveLength(1); + expect(transportState.websocketOptions).toHaveLength(2); + expect(transportState.websocketRequests[1]).not.toHaveProperty("previous_response_id"); expect(transportState.sdkRequests).toHaveLength(1); }); - it("does not fall back after a failed response contains terminal output", async () => { - transportState.responseBatches.push([ - message({ type: "response.created", response: { id: "resp_failed" } }), - message({ - type: "response.failed", - response: { - id: "resp_failed", - status: "failed", - output: [{ type: "message", role: "assistant", content: [] }], - }, - }), - ]); + it("keeps a true post-dispatch connection loss replay-unsafe", async () => { + transportState.responseBatches.push([{ type: "close", code: 1006 }]); const result = await run({ messages: [userMessage("hello", 1)], tools: [] }); - expect(result.stopReason).toBe("error"); - expect(result.errorCode).toBe(PROVIDER_FAILURE_WITH_OUTPUT_ERROR_CODE); + expect(result).toMatchObject({ + stopReason: "error", + errorCode: PROVIDER_POST_DISPATCH_AMBIGUITY_ERROR_CODE, + }); + expect(isRetryableAssistantError(result)).toBe(false); + expect(transportState.websocketRequests).toHaveLength(1); expect(transportState.sdkRequests).toEqual([]); }); diff --git a/packages/ai/src/transports/openai-responses-websocket.ts b/packages/ai/src/transports/openai-responses-websocket.ts index 93640668d2cb..b6bb09794f17 100644 --- a/packages/ai/src/transports/openai-responses-websocket.ts +++ b/packages/ai/src/transports/openai-responses-websocket.ts @@ -16,7 +16,6 @@ import { parseOpenAIResponsesWebSocketServerError, OpenAIResponsesWebSocketPostDispatchError, OpenAIResponsesWebSocketPreDispatchError, - OpenAIResponsesWebSocketResponseFailedError, OpenAIResponsesWebSocketSafeRetryError, } from "./openai-responses-contracts.js"; import { transportAbortError } from "./transport-stream-shared.js"; @@ -425,14 +424,13 @@ export function createOpenAIResponsesWebSocketStream(params: { if (!event) { continue; } - if (event.type === "response.failed") { - throw new OpenAIResponsesWebSocketResponseFailedError(event.response.output.length > 0); - } if (event.type === "response.completed") { terminalResponse = event.response; } terminalReceived = - event.type === "response.completed" || event.type === "response.incomplete"; + event.type === "response.completed" || + event.type === "response.incomplete" || + event.type === "response.failed"; yield event; if (terminalReceived) { degradedWebSocketConnections.delete(degradationKey); @@ -450,12 +448,7 @@ export function createOpenAIResponsesWebSocketStream(params: { if (!requestDispatched && !params.signal?.aborted) { throw new OpenAIResponsesWebSocketPreDispatchError(error); } - if ( - !requestDispatched || - params.callerSignal?.aborted || - error instanceof OpenAIResponsesWebSocketResponseFailedError || - safeRetry - ) { + if (!requestDispatched || params.callerSignal?.aborted || safeRetry) { throw error; } throw new OpenAIResponsesWebSocketPostDispatchError(error);