From 1df40142633d16a1f66e80166a1d5b8abf2d9cc9 Mon Sep 17 00:00:00 2001 From: steebchen Date: Sat, 21 Mar 2026 09:01:15 +0000 Subject: [PATCH 1/3] fix: improve streaming error diagnostics and logger serialization --- apps/gateway/src/chat/chat.ts | 62 ++++--- .../tools/normalize-streaming-error.spec.ts | 56 +++++++ .../chat/tools/normalize-streaming-error.ts | 157 ++++++++++++++++++ packages/logger/src/index.ts | 31 +++- packages/logger/src/logger.spec.ts | 20 +++ 5 files changed, 294 insertions(+), 32 deletions(-) create mode 100644 apps/gateway/src/chat/tools/normalize-streaming-error.spec.ts create mode 100644 apps/gateway/src/chat/tools/normalize-streaming-error.ts diff --git a/apps/gateway/src/chat/chat.ts b/apps/gateway/src/chat/chat.ts index 5a28f411a2..933f205b2c 100644 --- a/apps/gateway/src/chat/chat.ts +++ b/apps/gateway/src/chat/chat.ts @@ -88,6 +88,7 @@ import { healJsonResponse } from "./tools/heal-json-response.js"; import { isModelTrulyFree } from "./tools/is-model-truly-free.js"; import { messagesContainImages } from "./tools/messages-contain-images.js"; import { mightBeCompleteJson } from "./tools/might-be-complete-json.js"; +import { normalizeStreamingError } from "./tools/normalize-streaming-error.js"; import { convertAwsEventStreamToSSE } from "./tools/parse-aws-eventstream.js"; import { parseModelInput } from "./tools/parse-model-input.js"; import { parseProviderResponse } from "./tools/parse-provider-response.js"; @@ -4324,9 +4325,42 @@ chat.openapi(completions, async (c) => { }, }; } else { - logger.warn( - "Error reading stream", + const normalizedStreamingError = normalizeStreamingError({ + error, + provider: usedProvider, + model: usedModel, + bufferSnapshot: buffer ? buffer.substring(0, 5000) : undefined, + phase: "upstream_read", + }); + + logger.error( + "Error reading upstream stream", error instanceof Error ? error : new Error(String(error)), + { + requestId, + usedProvider, + requestedProvider, + usedModel, + initialRequestedModel, + upstreamStatus: res?.status ?? null, + upstreamStatusText: res?.statusText ?? null, + upstreamHeaders: res + ? { + contentType: res.headers.get("content-type"), + contentLength: res.headers.get("content-length"), + transferEncoding: res.headers.get("transfer-encoding"), + requestId: + res.headers.get("x-request-id") ?? + res.headers.get("request-id") ?? + res.headers.get("openai-request-id"), + } + : null, + streamingDiagnostics: normalizedStreamingError.log.details, + timeToFirstToken, + timeToFirstReasoningToken, + firstTokenReceived, + firstReasoningTokenReceived, + }, ); // Forward the error to the client with the buffered content that caused the error @@ -4334,14 +4368,7 @@ chat.openapi(completions, async (c) => { await stream.writeSSE({ event: "error", data: JSON.stringify({ - error: { - message: `Streaming error: ${error instanceof Error ? error.message : String(error)}`, - type: "gateway_error", - param: null, - code: "streaming_error", - // Include the buffer content that caused the parsing error - responseText: buffer.substring(0, 5000), // Limit to 5000 chars to avoid too large error messages - }, + error: normalizedStreamingError.client, }), id: String(eventId++), }); @@ -4360,20 +4387,7 @@ chat.openapi(completions, async (c) => { ); } - // Create structured error object for logging - streamingError = { - message: error instanceof Error ? error.message : String(error), - type: "streaming_error", - code: "streaming_error", - details: { - name: error instanceof Error ? error.name : "UnknownError", - stack: error instanceof Error ? error.stack : undefined, - timestamp: new Date().toISOString(), - provider: usedProvider, - model: usedModel, - bufferSnapshot: buffer ? buffer.substring(0, 5000) : undefined, - }, - }; + streamingError = normalizedStreamingError.log; } } finally { // Clean up the reader to prevent file descriptor leaks diff --git a/apps/gateway/src/chat/tools/normalize-streaming-error.spec.ts b/apps/gateway/src/chat/tools/normalize-streaming-error.spec.ts new file mode 100644 index 0000000000..e8a8dbc8b8 --- /dev/null +++ b/apps/gateway/src/chat/tools/normalize-streaming-error.spec.ts @@ -0,0 +1,56 @@ +import { describe, expect, it } from "vitest"; + +import { normalizeStreamingError } from "./normalize-streaming-error.js"; + +describe("normalizeStreamingError", () => { + it("classifies terminated undici stream reads as upstream termination", () => { + const socketCloseError = new Error("other side closed") as Error & { + code?: string; + }; + socketCloseError.name = "SocketError"; + socketCloseError.code = "UND_ERR_SOCKET"; + + const error = new TypeError("terminated", { + cause: socketCloseError, + }); + + const normalized = normalizeStreamingError({ + error, + provider: "canopywave", + model: "deepseek/deepseek-chat-v3.2", + bufferSnapshot: "\n\n", + phase: "upstream_read", + }); + + expect(normalized.client.message).toBe( + "Upstream stream terminated unexpectedly before completion", + ); + expect(normalized.client.details.statusCode).toBe(502); + expect(normalized.client.details.statusText).toBe( + "Upstream Stream Terminated", + ); + expect(normalized.client.details.errorCode).toBe("UND_ERR_SOCKET"); + expect(normalized.log.details.responseText).toContain("terminated"); + expect(normalized.log.details.cause).toContain("UND_ERR_SOCKET"); + }); + + it("preserves generic streaming read errors with 500 classification", () => { + const error = new SyntaxError("Unexpected end of JSON input"); + + const normalized = normalizeStreamingError({ + error, + provider: "openai", + model: "gpt-4.1-mini", + bufferSnapshot: "data: {", + phase: "upstream_read", + }); + + expect(normalized.client.message).toBe( + "Streaming error: Unexpected end of JSON input", + ); + expect(normalized.client.details.statusCode).toBe(500); + expect(normalized.client.details.statusText).toBe("Streaming Read Error"); + expect(normalized.log.details.name).toBe("SyntaxError"); + expect(normalized.log.details.bufferSnapshot).toBe("data: {"); + }); +}); diff --git a/apps/gateway/src/chat/tools/normalize-streaming-error.ts b/apps/gateway/src/chat/tools/normalize-streaming-error.ts new file mode 100644 index 0000000000..7ae693f083 --- /dev/null +++ b/apps/gateway/src/chat/tools/normalize-streaming-error.ts @@ -0,0 +1,157 @@ +import { extractErrorCause } from "./extract-error-cause.js"; + +interface ErrorWithCode extends Error { + code?: string; + cause?: unknown; +} + +export interface NormalizeStreamingErrorOptions { + error: unknown; + provider: string; + model: string; + bufferSnapshot?: string; + phase: "upstream_connect" | "upstream_read"; +} + +export interface NormalizedStreamingError { + client: { + message: string; + type: "gateway_error"; + param: null; + code: "streaming_error"; + responseText?: string; + details: { + statusCode: number; + statusText: string; + errorName: string; + errorCode?: string; + cause?: string; + }; + }; + log: { + message: string; + type: "streaming_error"; + code: "streaming_error"; + details: { + statusCode: number; + statusText: string; + responseText: string; + cause?: string; + name: string; + errorCode?: string; + timestamp: string; + provider: string; + model: string; + phase: NormalizeStreamingErrorOptions["phase"]; + bufferSnapshot?: string; + stack?: string; + }; + }; +} + +function getErrorCode(error: unknown): string | undefined { + if (!(error instanceof Error)) { + return undefined; + } + + const directCode = + typeof (error as ErrorWithCode).code === "string" + ? (error as ErrorWithCode).code + : undefined; + if (directCode) { + return directCode; + } + + let current = (error as ErrorWithCode).cause; + for (let depth = 0; depth < 5; depth++) { + if (!(current instanceof Error)) { + return undefined; + } + + if (typeof (current as ErrorWithCode).code === "string") { + return (current as ErrorWithCode).code; + } + + current = (current as ErrorWithCode).cause; + } + + return undefined; +} + +function isUpstreamTermination(error: unknown, cause?: string): boolean { + if (!(error instanceof Error)) { + return false; + } + + const normalizedMessage = error.message.trim().toLowerCase(); + const normalizedCause = cause?.toLowerCase() ?? ""; + + return ( + (error.name === "TypeError" && normalizedMessage === "terminated") || + normalizedCause.includes("onhttpsocketclose") || + normalizedCause.includes("socket") || + normalizedCause.includes("other side closed") || + normalizedCause.includes("und_err") || + normalizedCause.includes("econnreset") + ); +} + +export function normalizeStreamingError( + options: NormalizeStreamingErrorOptions, +): NormalizedStreamingError { + const { error, provider, model, bufferSnapshot, phase } = options; + + const errorName = error instanceof Error ? error.name : "UnknownError"; + const rawMessage = + error instanceof Error ? error.message : String(error ?? "Unknown error"); + const cause = extractErrorCause(error); + const errorCode = getErrorCode(error); + + const terminated = isUpstreamTermination(error, cause); + const statusCode = terminated ? 502 : 500; + const statusText = terminated + ? "Upstream Stream Terminated" + : "Streaming Read Error"; + const message = terminated + ? "Upstream stream terminated unexpectedly before completion" + : `Streaming error: ${rawMessage}`; + const responseText = cause ? `${rawMessage} | cause: ${cause}` : rawMessage; + + return { + client: { + message, + type: "gateway_error", + param: null, + code: "streaming_error", + responseText: bufferSnapshot, + details: { + statusCode, + statusText, + errorName, + ...(errorCode ? { errorCode } : {}), + ...(cause ? { cause } : {}), + }, + }, + log: { + message: rawMessage, + type: "streaming_error", + code: "streaming_error", + details: { + statusCode, + statusText, + responseText, + ...(cause ? { cause } : {}), + name: errorName, + ...(errorCode ? { errorCode } : {}), + timestamp: new Date().toISOString(), + provider, + model, + phase, + ...(bufferSnapshot ? { bufferSnapshot } : {}), + ...(error instanceof Error && error.stack + ? { stack: error.stack } + : {}), + }, + }, + }; +} diff --git a/packages/logger/src/index.ts b/packages/logger/src/index.ts index ac1312c2a9..ac63384be4 100644 --- a/packages/logger/src/index.ts +++ b/packages/logger/src/index.ts @@ -130,24 +130,24 @@ class LLMGatewayLogger { } // Core logging methods - public trace(message: string, extra?: object): void { + public trace(message: string, extra?: object | Error): void { const traceContext = this.getTraceContext(); - this.logger.trace({ ...traceContext, ...extra }, message); + this.logger.trace(this.mergeOptionalArg(traceContext, extra), message); } - public debug(message: string, extra?: object): void { + public debug(message: string, extra?: object | Error): void { const traceContext = this.getTraceContext(); - this.logger.debug({ ...traceContext, ...extra }, message); + this.logger.debug(this.mergeOptionalArg(traceContext, extra), message); } - public info(message: string, extra?: object): void { + public info(message: string, extra?: object | Error): void { const traceContext = this.getTraceContext(); - this.logger.info({ ...traceContext, ...extra }, message); + this.logger.info(this.mergeOptionalArg(traceContext, extra), message); } - public warn(message: string, extra?: object): void { + public warn(message: string, extra?: object | Error): void { const traceContext = this.getTraceContext(); - this.logger.warn({ ...traceContext, ...extra }, message); + this.logger.warn(this.mergeOptionalArg(traceContext, extra), message); } public error(message: string, ...args: unknown[]): void { @@ -179,6 +179,21 @@ class LLMGatewayLogger { return result; } + private mergeOptionalArg( + traceContext: object, + extra?: object | Error, + ): Record { + if (extra instanceof Error) { + return { ...traceContext, err: extra }; + } + + if (extra && typeof extra === "object") { + return { ...traceContext, ...extra }; + } + + return { ...traceContext }; + } + // Create child logger with additional context public child(bindings: object): LLMGatewayLogger { const childPino = this.logger.child(bindings); diff --git a/packages/logger/src/logger.spec.ts b/packages/logger/src/logger.spec.ts index cf004d295b..2b7007b835 100644 --- a/packages/logger/src/logger.spec.ts +++ b/packages/logger/src/logger.spec.ts @@ -47,4 +47,24 @@ describe("LLMGateway Logger", () => { expect(childLogger).toBeDefined(); expect(typeof childLogger.info).toBe("function"); }); + + it("should serialize Error instances passed to warn", () => { + const customLogger = createLogger({ + name: "test-logger", + level: "warn", + prettyPrint: false, + }) as any; + + const warnSpy = vi.fn(); + customLogger.logger = { warn: warnSpy }; + + const error = new Error("stream terminated"); + customLogger.warn("Error reading stream", error); + + expect(warnSpy).toHaveBeenCalledTimes(1); + expect(warnSpy.mock.calls[0][0]).toMatchObject({ + err: error, + }); + expect(warnSpy.mock.calls[0][1]).toBe("Error reading stream"); + }); }); From fb58c2fcf5e5b14c69e7ef669eecdadc6695b4e8 Mon Sep 17 00:00:00 2001 From: steebchen Date: Sat, 21 Mar 2026 09:05:43 +0000 Subject: [PATCH 2/3] fix: include unified finish reason in chat errors --- apps/gateway/src/chat/chat.ts | 431 +++++++++++++++++++++++----------- 1 file changed, 288 insertions(+), 143 deletions(-) diff --git a/apps/gateway/src/chat/chat.ts b/apps/gateway/src/chat/chat.ts index 933f205b2c..fc058b974b 100644 --- a/apps/gateway/src/chat/chat.ts +++ b/apps/gateway/src/chat/chat.ts @@ -17,7 +17,11 @@ import { import { isCodingModel } from "@/lib/coding-models.js"; import { calculateCosts, shouldBillCancelledRequests } from "@/lib/costs.js"; import { throwIamException, validateModelAccess } from "@/lib/iam.js"; -import { calculateDataStorageCost, insertLog } from "@/lib/logs.js"; +import { + calculateDataStorageCost, + getUnifiedFinishReason, + insertLog, +} from "@/lib/logs.js"; import { createCombinedSignal, createStreamingCombinedSignal, @@ -44,6 +48,7 @@ import { type InferSelectModel, isCachingEnabled, shortid, + UnifiedFinishReason, type tables, } from "@llmgateway/db"; import { @@ -119,6 +124,70 @@ const SSE_FIELD_PATTERN = /^[a-zA-Z_-]+:\s*/; // Reusable TextDecoder to avoid per-chunk allocation in the streaming hot path const sharedTextDecoder = new TextDecoder(); +function getUnifiedFinishReasonForError( + errorType: string | undefined, + status?: number, +) { + if (!errorType) { + return UnifiedFinishReason.UNKNOWN; + } + + if (errorType === "invalid_request_error") { + return UnifiedFinishReason.CLIENT_ERROR; + } + + if (errorType === "upstream_timeout") { + return UnifiedFinishReason.UPSTREAM_ERROR; + } + + if (errorType === "request_canceled" || errorType === "canceled") { + return UnifiedFinishReason.CANCELED; + } + + const unifiedFinishReason = getUnifiedFinishReason(errorType, undefined); + if (unifiedFinishReason !== UnifiedFinishReason.UNKNOWN) { + return unifiedFinishReason; + } + + if (status !== undefined) { + if (status === 401 || status === 403) { + return UnifiedFinishReason.GATEWAY_ERROR; + } + if (status >= 400 && status < 500) { + return UnifiedFinishReason.CLIENT_ERROR; + } + if (status >= 500) { + return UnifiedFinishReason.UPSTREAM_ERROR; + } + } + + return UnifiedFinishReason.UNKNOWN; +} + +function withUnifiedFinishReason }>( + payload: T, + status?: number, +): T { + if (!payload.error || typeof payload.error !== "object") { + return payload; + } + + if (typeof payload.error.unified_finish_reason === "string") { + return payload; + } + + const errorType = + typeof payload.error.type === "string" ? payload.error.type : undefined; + + return { + ...payload, + error: { + ...payload.error, + unified_finish_reason: getUnifiedFinishReasonForError(errorType, status), + }, + }; +} + export const chat = new OpenAPIHono(); const completions = createRoute({ @@ -234,6 +303,7 @@ const completions = createRoute({ type: z.string(), param: z.string().nullable(), code: z.string(), + unified_finish_reason: z.string().optional(), }), }), }, @@ -256,14 +326,17 @@ chat.openapi(completions, async (c) => { rawBody = await c.req.json(); } catch { return c.json( - { - error: { - message: "Invalid JSON in request body", - type: "invalid_request_error", - param: null, - code: "invalid_json", + withUnifiedFinishReason( + { + error: { + message: "Invalid JSON in request body", + type: "invalid_request_error", + param: null, + code: "invalid_json", + }, }, - }, + 400, + ), 400, ); } @@ -272,14 +345,17 @@ chat.openapi(completions, async (c) => { const validationResult = completionsRequestSchema.safeParse(rawBody); if (!validationResult.success) { return c.json( - { - error: { - message: "Invalid request parameters", - type: "invalid_request_error", - param: null, - code: "invalid_parameters", + withUnifiedFinishReason( + { + error: { + message: "Invalid request parameters", + type: "invalid_request_error", + param: null, + code: "invalid_parameters", + }, }, - }, + 400, + ), 400, ); } @@ -343,14 +419,17 @@ chat.openapi(completions, async (c) => { reasoning_object_effort !== undefined ) { return c.json( - { - error: { - message: - "Cannot specify both reasoning_effort and reasoning.effort. Use one or the other.", - type: "invalid_request_error", - code: "invalid_request", + withUnifiedFinishReason( + { + error: { + message: + "Cannot specify both reasoning_effort and reasoning.effort. Use one or the other.", + type: "invalid_request_error", + code: "invalid_request", + }, }, - }, + 400, + ), 400, ); } @@ -2598,13 +2677,18 @@ chat.openapi(completions, async (c) => { await stream.writeSSE({ event: "error", - data: JSON.stringify({ - error: { - message: `Upstream provider timeout: ${errorMessage}`, - type: "upstream_timeout", - code: "timeout", - }, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { + error: { + message: `Upstream provider timeout: ${errorMessage}`, + type: "upstream_timeout", + code: "timeout", + }, + }, + 504, + ), + ), id: String(eventId++), }); return; @@ -2882,13 +2966,18 @@ chat.openapi(completions, async (c) => { // Send error event to the client await writeSSEAndCache({ event: "error", - data: JSON.stringify({ - error: { - message: `Failed to connect to provider: ${errorMessage}`, - type: "upstream_error", - code: "fetch_failed", - }, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { + error: { + message: `Failed to connect to provider: ${errorMessage}`, + type: "upstream_error", + code: "fetch_failed", + }, + }, + 502, + ), + ), id: String(eventId++), }); await writeSSEAndCache({ @@ -3140,7 +3229,9 @@ chat.openapi(completions, async (c) => { await writeSSEAndCache({ event: "error", - data: JSON.stringify(errorData), + data: JSON.stringify( + withUnifiedFinishReason(errorData, res.status), + ), id: String(eventId++), }); await writeSSEAndCache({ @@ -3198,13 +3289,18 @@ chat.openapi(completions, async (c) => { if (!res || !res.ok) { await writeSSEAndCache({ event: "error", - data: JSON.stringify({ - error: { - message: "All provider attempts failed", - type: "upstream_error", - code: "all_providers_failed", - }, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { + error: { + message: "All provider attempts failed", + type: "upstream_error", + code: "all_providers_failed", + }, + }, + 502, + ), + ), id: String(eventId++), }); await writeSSEAndCache({ @@ -3230,14 +3326,19 @@ chat.openapi(completions, async (c) => { if (!res.body) { await writeSSEAndCache({ event: "error", - data: JSON.stringify({ - error: { - message: "No response body from provider", - type: "gateway_error", - param: null, - code: "gateway_error", - }, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { + error: { + message: "No response body from provider", + type: "gateway_error", + param: null, + code: "gateway_error", + }, + }, + 500, + ), + ), id: String(eventId++), }); await writeSSEAndCache({ @@ -3352,14 +3453,19 @@ chat.openapi(completions, async (c) => { try { await stream.writeSSE({ event: "error", - data: JSON.stringify({ - error: { - message: `Streaming buffer exceeded ${bufferSizeMB}MB limit`, - type: "gateway_error", - param: null, - code: "buffer_overflow", - }, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { + error: { + message: `Streaming buffer exceeded ${bufferSizeMB}MB limit`, + type: "gateway_error", + param: null, + code: "buffer_overflow", + }, + }, + 500, + ), + ), id: String(eventId++), }); await stream.writeSSE({ @@ -3793,15 +3899,20 @@ chat.openapi(completions, async (c) => { await writeSSEAndCache({ event: "error", - data: JSON.stringify({ - error: { - message: awsBedrockStreamError.message, - type: errorType, - code: awsBedrockStreamError.eventType, - param: null, - responseText: awsBedrockStreamError.responseText, - }, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { + error: { + message: awsBedrockStreamError.message, + type: errorType, + code: awsBedrockStreamError.eventType, + param: null, + responseText: awsBedrockStreamError.responseText, + }, + }, + awsBedrockStreamError.statusCode, + ), + ), id: String(eventId++), }); await writeSSEAndCache({ @@ -4288,14 +4399,19 @@ chat.openapi(completions, async (c) => { try { await stream.writeSSE({ event: "error", - data: JSON.stringify({ - error: { - message: `Upstream provider timeout: ${errorMessage}`, - type: "upstream_timeout", - param: null, - code: "timeout", - }, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { + error: { + message: `Upstream provider timeout: ${errorMessage}`, + type: "upstream_timeout", + param: null, + code: "timeout", + }, + }, + 504, + ), + ), id: String(eventId++), }); await stream.writeSSE({ @@ -4367,9 +4483,12 @@ chat.openapi(completions, async (c) => { try { await stream.writeSSE({ event: "error", - data: JSON.stringify({ - error: normalizedStreamingError.client, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { error: normalizedStreamingError.client }, + normalizedStreamingError.client.details.statusCode, + ), + ), id: String(eventId++), }); await stream.writeSSE({ @@ -4513,15 +4632,20 @@ chat.openapi(completions, async (c) => { try { await writeSSEAndCache({ event: "error", - data: JSON.stringify({ - error: { - message: errorMessage, - type: "upstream_error", - code: "upstream_error", - param: null, - responseText: errorMessage, - }, - }), + data: JSON.stringify( + withUnifiedFinishReason( + { + error: { + message: errorMessage, + type: "upstream_error", + code: "upstream_error", + param: null, + responseText: errorMessage, + }, + }, + 502, + ), + ), id: String(eventId++), }); await writeSSEAndCache({ @@ -5338,20 +5462,23 @@ chat.openapi(completions, async (c) => { // Return error response - use 504 for timeouts, 502 for other connection failures return c.json( - { - error: { - message: isTimeoutFetchError - ? `Upstream provider timeout: ${errorMessage}` - : `Failed to connect to provider: ${errorMessage}`, - type: isTimeoutFetchError ? "upstream_timeout" : "upstream_error", - param: null, - code: isTimeoutFetchError ? "timeout" : "fetch_failed", - requestedProvider, - usedProvider, - requestedModel: initialRequestedModel, - usedModel, + withUnifiedFinishReason( + { + error: { + message: isTimeoutFetchError + ? `Upstream provider timeout: ${errorMessage}` + : `Failed to connect to provider: ${errorMessage}`, + type: isTimeoutFetchError ? "upstream_timeout" : "upstream_error", + param: null, + code: isTimeoutFetchError ? "timeout" : "fetch_failed", + requestedProvider, + usedProvider, + requestedModel: initialRequestedModel, + usedModel, + }, }, - }, + isTimeoutFetchError ? 504 : 502, + ), isTimeoutFetchError ? 504 : 502, ); } @@ -5487,14 +5614,17 @@ chat.openapi(completions, async (c) => { }); return c.json( - { - error: { - message: "Request canceled by client", - type: "canceled", - param: null, - code: "request_canceled", + withUnifiedFinishReason( + { + error: { + message: "Request canceled by client", + type: "canceled", + param: null, + code: "request_canceled", + }, }, - }, + 400, + ), 400, ); // Using 400 status code for client closed request } @@ -5598,14 +5728,17 @@ chat.openapi(completions, async (c) => { }); return c.json( - { - error: { - message: `Upstream provider timeout: ${errorMessage}`, - type: "upstream_timeout", - param: null, - code: "timeout", + withUnifiedFinishReason( + { + error: { + message: `Upstream provider timeout: ${errorMessage}`, + type: "upstream_timeout", + param: null, + code: "timeout", + }, }, - }, + 504, + ), 504, ); } @@ -5801,7 +5934,10 @@ chat.openapi(completions, async (c) => { if (finishReason === "client_error") { try { const originalError = JSON.parse(errorResponseText); - return c.json(originalError, res.status as 400); + return c.json( + withUnifiedFinishReason(originalError, res.status), + res.status as 400, + ); } catch { // If we can't parse the original error, fall back to our format } @@ -5809,19 +5945,22 @@ chat.openapi(completions, async (c) => { // Return our wrapped error response for non-client errors return c.json( - { - error: { - message: `Error from provider: ${res.status} ${res.statusText} ${errorResponseText}`, - type: finishReason, - param: null, - code: finishReason, - requestedProvider, - usedProvider, - requestedModel: initialRequestedModel, - usedModel, - responseText: errorResponseText, + withUnifiedFinishReason( + { + error: { + message: `Error from provider: ${res.status} ${res.statusText} ${errorResponseText}`, + type: finishReason, + param: null, + code: finishReason, + requestedProvider, + usedProvider, + requestedModel: initialRequestedModel, + usedModel, + responseText: errorResponseText, + }, }, - }, + 500, + ), 500, ); } @@ -5867,14 +6006,17 @@ chat.openapi(completions, async (c) => { if (!res || !res.ok) { // All retries exhausted return c.json( - { - error: { - message: "All provider attempts failed", - type: "upstream_error", - param: null, - code: "all_providers_failed", + withUnifiedFinishReason( + { + error: { + message: "All provider attempts failed", + type: "upstream_error", + param: null, + code: "all_providers_failed", + }, }, - }, + 502, + ), 502, ); } @@ -5976,14 +6118,17 @@ chat.openapi(completions, async (c) => { }); return c.json( - { - error: { - message: `Upstream provider timeout: ${errorMessage}`, - type: "upstream_timeout", - param: null, - code: "timeout", + withUnifiedFinishReason( + { + error: { + message: `Upstream provider timeout: ${errorMessage}`, + type: "upstream_timeout", + param: null, + code: "timeout", + }, }, - }, + 504, + ), 504, ); } From 50a84b0119d6390734e1f3681012f1251f4f186d Mon Sep 17 00:00:00 2001 From: steebchen Date: Sat, 21 Mar 2026 09:11:32 +0000 Subject: [PATCH 3/3] fix: keep unified finish reason in logs only --- apps/gateway/src/chat/chat.ts | 467 ++++++++++++++-------------------- 1 file changed, 184 insertions(+), 283 deletions(-) diff --git a/apps/gateway/src/chat/chat.ts b/apps/gateway/src/chat/chat.ts index fc058b974b..d57edd43cb 100644 --- a/apps/gateway/src/chat/chat.ts +++ b/apps/gateway/src/chat/chat.ts @@ -48,7 +48,6 @@ import { type InferSelectModel, isCachingEnabled, shortid, - UnifiedFinishReason, type tables, } from "@llmgateway/db"; import { @@ -124,70 +123,6 @@ const SSE_FIELD_PATTERN = /^[a-zA-Z_-]+:\s*/; // Reusable TextDecoder to avoid per-chunk allocation in the streaming hot path const sharedTextDecoder = new TextDecoder(); -function getUnifiedFinishReasonForError( - errorType: string | undefined, - status?: number, -) { - if (!errorType) { - return UnifiedFinishReason.UNKNOWN; - } - - if (errorType === "invalid_request_error") { - return UnifiedFinishReason.CLIENT_ERROR; - } - - if (errorType === "upstream_timeout") { - return UnifiedFinishReason.UPSTREAM_ERROR; - } - - if (errorType === "request_canceled" || errorType === "canceled") { - return UnifiedFinishReason.CANCELED; - } - - const unifiedFinishReason = getUnifiedFinishReason(errorType, undefined); - if (unifiedFinishReason !== UnifiedFinishReason.UNKNOWN) { - return unifiedFinishReason; - } - - if (status !== undefined) { - if (status === 401 || status === 403) { - return UnifiedFinishReason.GATEWAY_ERROR; - } - if (status >= 400 && status < 500) { - return UnifiedFinishReason.CLIENT_ERROR; - } - if (status >= 500) { - return UnifiedFinishReason.UPSTREAM_ERROR; - } - } - - return UnifiedFinishReason.UNKNOWN; -} - -function withUnifiedFinishReason }>( - payload: T, - status?: number, -): T { - if (!payload.error || typeof payload.error !== "object") { - return payload; - } - - if (typeof payload.error.unified_finish_reason === "string") { - return payload; - } - - const errorType = - typeof payload.error.type === "string" ? payload.error.type : undefined; - - return { - ...payload, - error: { - ...payload.error, - unified_finish_reason: getUnifiedFinishReasonForError(errorType, status), - }, - }; -} - export const chat = new OpenAPIHono(); const completions = createRoute({ @@ -303,7 +238,6 @@ const completions = createRoute({ type: z.string(), param: z.string().nullable(), code: z.string(), - unified_finish_reason: z.string().optional(), }), }), }, @@ -326,17 +260,14 @@ chat.openapi(completions, async (c) => { rawBody = await c.req.json(); } catch { return c.json( - withUnifiedFinishReason( - { - error: { - message: "Invalid JSON in request body", - type: "invalid_request_error", - param: null, - code: "invalid_json", - }, + { + error: { + message: "Invalid JSON in request body", + type: "invalid_request_error", + param: null, + code: "invalid_json", }, - 400, - ), + }, 400, ); } @@ -345,17 +276,14 @@ chat.openapi(completions, async (c) => { const validationResult = completionsRequestSchema.safeParse(rawBody); if (!validationResult.success) { return c.json( - withUnifiedFinishReason( - { - error: { - message: "Invalid request parameters", - type: "invalid_request_error", - param: null, - code: "invalid_parameters", - }, + { + error: { + message: "Invalid request parameters", + type: "invalid_request_error", + param: null, + code: "invalid_parameters", }, - 400, - ), + }, 400, ); } @@ -419,17 +347,14 @@ chat.openapi(completions, async (c) => { reasoning_object_effort !== undefined ) { return c.json( - withUnifiedFinishReason( - { - error: { - message: - "Cannot specify both reasoning_effort and reasoning.effort. Use one or the other.", - type: "invalid_request_error", - code: "invalid_request", - }, + { + error: { + message: + "Cannot specify both reasoning_effort and reasoning.effort. Use one or the other.", + type: "invalid_request_error", + code: "invalid_request", }, - 400, - ), + }, 400, ); } @@ -2571,6 +2496,10 @@ chat.openapi(completions, async (c) => { requestedProvider, usedModel, initialRequestedModel, + unifiedFinishReason: getUnifiedFinishReason( + "upstream_error", + usedProvider, + ), }); // Log the timeout error in the database @@ -2677,18 +2606,13 @@ chat.openapi(completions, async (c) => { await stream.writeSSE({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { - error: { - message: `Upstream provider timeout: ${errorMessage}`, - type: "upstream_timeout", - code: "timeout", - }, - }, - 504, - ), - ), + data: JSON.stringify({ + error: { + message: `Upstream provider timeout: ${errorMessage}`, + type: "upstream_timeout", + code: "timeout", + }, + }), id: String(eventId++), }); return; @@ -2853,6 +2777,10 @@ chat.openapi(completions, async (c) => { requestedProvider, usedModel, initialRequestedModel, + unifiedFinishReason: getUnifiedFinishReason( + "upstream_error", + usedProvider, + ), }); // Log the error in the database @@ -2966,18 +2894,13 @@ chat.openapi(completions, async (c) => { // Send error event to the client await writeSSEAndCache({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { - error: { - message: `Failed to connect to provider: ${errorMessage}`, - type: "upstream_error", - code: "fetch_failed", - }, - }, - 502, - ), - ), + data: JSON.stringify({ + error: { + message: `Failed to connect to provider: ${errorMessage}`, + type: "upstream_error", + code: "fetch_failed", + }, + }), id: String(eventId++), }); await writeSSEAndCache({ @@ -3019,6 +2942,10 @@ chat.openapi(completions, async (c) => { organizationId: project.organizationId, projectId: apiKey.projectId, apiKeyId: apiKey.id, + unifiedFinishReason: getUnifiedFinishReason( + finishReason, + usedProvider, + ), }); } @@ -3229,9 +3156,7 @@ chat.openapi(completions, async (c) => { await writeSSEAndCache({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason(errorData, res.status), - ), + data: JSON.stringify(errorData), id: String(eventId++), }); await writeSSEAndCache({ @@ -3289,18 +3214,13 @@ chat.openapi(completions, async (c) => { if (!res || !res.ok) { await writeSSEAndCache({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { - error: { - message: "All provider attempts failed", - type: "upstream_error", - code: "all_providers_failed", - }, - }, - 502, - ), - ), + data: JSON.stringify({ + error: { + message: "All provider attempts failed", + type: "upstream_error", + code: "all_providers_failed", + }, + }), id: String(eventId++), }); await writeSSEAndCache({ @@ -3326,19 +3246,14 @@ chat.openapi(completions, async (c) => { if (!res.body) { await writeSSEAndCache({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { - error: { - message: "No response body from provider", - type: "gateway_error", - param: null, - code: "gateway_error", - }, - }, - 500, - ), - ), + data: JSON.stringify({ + error: { + message: "No response body from provider", + type: "gateway_error", + param: null, + code: "gateway_error", + }, + }), id: String(eventId++), }); await writeSSEAndCache({ @@ -3453,19 +3368,14 @@ chat.openapi(completions, async (c) => { try { await stream.writeSSE({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { - error: { - message: `Streaming buffer exceeded ${bufferSizeMB}MB limit`, - type: "gateway_error", - param: null, - code: "buffer_overflow", - }, - }, - 500, - ), - ), + data: JSON.stringify({ + error: { + message: `Streaming buffer exceeded ${bufferSizeMB}MB limit`, + type: "gateway_error", + param: null, + code: "buffer_overflow", + }, + }), id: String(eventId++), }); await stream.writeSSE({ @@ -3899,20 +3809,15 @@ chat.openapi(completions, async (c) => { await writeSSEAndCache({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { - error: { - message: awsBedrockStreamError.message, - type: errorType, - code: awsBedrockStreamError.eventType, - param: null, - responseText: awsBedrockStreamError.responseText, - }, - }, - awsBedrockStreamError.statusCode, - ), - ), + data: JSON.stringify({ + error: { + message: awsBedrockStreamError.message, + type: errorType, + code: awsBedrockStreamError.eventType, + param: null, + responseText: awsBedrockStreamError.responseText, + }, + }), id: String(eventId++), }); await writeSSEAndCache({ @@ -4394,24 +4299,23 @@ chat.openapi(completions, async (c) => { requestedProvider, usedModel, initialRequestedModel, + unifiedFinishReason: getUnifiedFinishReason( + "upstream_error", + usedProvider, + ), }); try { await stream.writeSSE({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { - error: { - message: `Upstream provider timeout: ${errorMessage}`, - type: "upstream_timeout", - param: null, - code: "timeout", - }, - }, - 504, - ), - ), + data: JSON.stringify({ + error: { + message: `Upstream provider timeout: ${errorMessage}`, + type: "upstream_timeout", + param: null, + code: "timeout", + }, + }), id: String(eventId++), }); await stream.writeSSE({ @@ -4476,6 +4380,12 @@ chat.openapi(completions, async (c) => { timeToFirstReasoningToken, firstTokenReceived, firstReasoningTokenReceived, + unifiedFinishReason: getUnifiedFinishReason( + normalizedStreamingError.client.type === "gateway_error" + ? "gateway_error" + : "upstream_error", + usedProvider, + ), }, ); @@ -4483,12 +4393,9 @@ chat.openapi(completions, async (c) => { try { await stream.writeSSE({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { error: normalizedStreamingError.client }, - normalizedStreamingError.client.details.statusCode, - ), - ), + data: JSON.stringify({ + error: normalizedStreamingError.client, + }), id: String(eventId++), }); await stream.writeSSE({ @@ -4622,6 +4529,10 @@ chat.openapi(completions, async (c) => { completionTokens, totalTokens, reasoningTokens, + unifiedFinishReason: getUnifiedFinishReason( + "upstream_error", + usedProvider, + ), }); const errorMessage = "Response finished successfully but returned no content or tool calls"; @@ -4632,20 +4543,15 @@ chat.openapi(completions, async (c) => { try { await writeSSEAndCache({ event: "error", - data: JSON.stringify( - withUnifiedFinishReason( - { - error: { - message: errorMessage, - type: "upstream_error", - code: "upstream_error", - param: null, - responseText: errorMessage, - }, - }, - 502, - ), - ), + data: JSON.stringify({ + error: { + message: errorMessage, + type: "upstream_error", + code: "upstream_error", + param: null, + responseText: errorMessage, + }, + }), id: String(eventId++), }); await writeSSEAndCache({ @@ -5349,6 +5255,10 @@ chat.openapi(completions, async (c) => { requestedProvider, usedModel, initialRequestedModel, + unifiedFinishReason: getUnifiedFinishReason( + "upstream_error", + usedProvider, + ), }); // Log the error in the database @@ -5462,23 +5372,20 @@ chat.openapi(completions, async (c) => { // Return error response - use 504 for timeouts, 502 for other connection failures return c.json( - withUnifiedFinishReason( - { - error: { - message: isTimeoutFetchError - ? `Upstream provider timeout: ${errorMessage}` - : `Failed to connect to provider: ${errorMessage}`, - type: isTimeoutFetchError ? "upstream_timeout" : "upstream_error", - param: null, - code: isTimeoutFetchError ? "timeout" : "fetch_failed", - requestedProvider, - usedProvider, - requestedModel: initialRequestedModel, - usedModel, - }, + { + error: { + message: isTimeoutFetchError + ? `Upstream provider timeout: ${errorMessage}` + : `Failed to connect to provider: ${errorMessage}`, + type: isTimeoutFetchError ? "upstream_timeout" : "upstream_error", + param: null, + code: isTimeoutFetchError ? "timeout" : "fetch_failed", + requestedProvider, + usedProvider, + requestedModel: initialRequestedModel, + usedModel, }, - isTimeoutFetchError ? 504 : 502, - ), + }, isTimeoutFetchError ? 504 : 502, ); } @@ -5614,17 +5521,14 @@ chat.openapi(completions, async (c) => { }); return c.json( - withUnifiedFinishReason( - { - error: { - message: "Request canceled by client", - type: "canceled", - param: null, - code: "request_canceled", - }, + { + error: { + message: "Request canceled by client", + type: "canceled", + param: null, + code: "request_canceled", }, - 400, - ), + }, 400, ); // Using 400 status code for client closed request } @@ -5651,6 +5555,10 @@ chat.openapi(completions, async (c) => { usedModel, status: res.status, cause: bodyErrorCause, + unifiedFinishReason: getUnifiedFinishReason( + "upstream_error", + usedProvider, + ), }); const bodyTimeoutPluginIds = plugins?.map((p) => p.id) ?? []; @@ -5728,17 +5636,14 @@ chat.openapi(completions, async (c) => { }); return c.json( - withUnifiedFinishReason( - { - error: { - message: `Upstream provider timeout: ${errorMessage}`, - type: "upstream_timeout", - param: null, - code: "timeout", - }, + { + error: { + message: `Upstream provider timeout: ${errorMessage}`, + type: "upstream_timeout", + param: null, + code: "timeout", }, - 504, - ), + }, 504, ); } @@ -5765,6 +5670,10 @@ chat.openapi(completions, async (c) => { organizationId: project.organizationId, projectId: apiKey.projectId, apiKeyId: apiKey.id, + unifiedFinishReason: getUnifiedFinishReason( + finishReason, + usedProvider, + ), }); } @@ -5934,10 +5843,7 @@ chat.openapi(completions, async (c) => { if (finishReason === "client_error") { try { const originalError = JSON.parse(errorResponseText); - return c.json( - withUnifiedFinishReason(originalError, res.status), - res.status as 400, - ); + return c.json(originalError, res.status as 400); } catch { // If we can't parse the original error, fall back to our format } @@ -5945,22 +5851,19 @@ chat.openapi(completions, async (c) => { // Return our wrapped error response for non-client errors return c.json( - withUnifiedFinishReason( - { - error: { - message: `Error from provider: ${res.status} ${res.statusText} ${errorResponseText}`, - type: finishReason, - param: null, - code: finishReason, - requestedProvider, - usedProvider, - requestedModel: initialRequestedModel, - usedModel, - responseText: errorResponseText, - }, + { + error: { + message: `Error from provider: ${res.status} ${res.statusText} ${errorResponseText}`, + type: finishReason, + param: null, + code: finishReason, + requestedProvider, + usedProvider, + requestedModel: initialRequestedModel, + usedModel, + responseText: errorResponseText, }, - 500, - ), + }, 500, ); } @@ -6006,17 +5909,14 @@ chat.openapi(completions, async (c) => { if (!res || !res.ok) { // All retries exhausted return c.json( - withUnifiedFinishReason( - { - error: { - message: "All provider attempts failed", - type: "upstream_error", - param: null, - code: "all_providers_failed", - }, + { + error: { + message: "All provider attempts failed", + type: "upstream_error", + param: null, + code: "all_providers_failed", }, - 502, - ), + }, 502, ); } @@ -6041,6 +5941,10 @@ chat.openapi(completions, async (c) => { usedModel, initialRequestedModel, cause: bodyReadCause, + unifiedFinishReason: getUnifiedFinishReason( + "upstream_error", + usedProvider, + ), }); const bodyTimeoutPluginIds = plugins?.map((p) => p.id) ?? []; @@ -6118,17 +6022,14 @@ chat.openapi(completions, async (c) => { }); return c.json( - withUnifiedFinishReason( - { - error: { - message: `Upstream provider timeout: ${errorMessage}`, - type: "upstream_timeout", - param: null, - code: "timeout", - }, + { + error: { + message: `Upstream provider timeout: ${errorMessage}`, + type: "upstream_timeout", + param: null, + code: "timeout", }, - 504, - ), + }, 504, ); }