Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
110 changes: 85 additions & 25 deletions apps/gateway/src/chat/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -88,6 +92,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";
Expand Down Expand Up @@ -2491,6 +2496,10 @@ chat.openapi(completions, async (c) => {
requestedProvider,
usedModel,
initialRequestedModel,
unifiedFinishReason: getUnifiedFinishReason(
"upstream_error",
usedProvider,
),
});

// Log the timeout error in the database
Expand Down Expand Up @@ -2768,6 +2777,10 @@ chat.openapi(completions, async (c) => {
requestedProvider,
usedModel,
initialRequestedModel,
unifiedFinishReason: getUnifiedFinishReason(
"upstream_error",
usedProvider,
),
});

// Log the error in the database
Expand Down Expand Up @@ -2929,6 +2942,10 @@ chat.openapi(completions, async (c) => {
organizationId: project.organizationId,
projectId: apiKey.projectId,
apiKeyId: apiKey.id,
unifiedFinishReason: getUnifiedFinishReason(
finishReason,
usedProvider,
),
});
}

Expand Down Expand Up @@ -4282,6 +4299,10 @@ chat.openapi(completions, async (c) => {
requestedProvider,
usedModel,
initialRequestedModel,
unifiedFinishReason: getUnifiedFinishReason(
"upstream_error",
usedProvider,
),
});

try {
Expand Down Expand Up @@ -4324,24 +4345,56 @@ 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,
unifiedFinishReason: getUnifiedFinishReason(
normalizedStreamingError.client.type === "gateway_error"
? "gateway_error"
: "upstream_error",
usedProvider,
),
},
);

// Forward the error to the client with the buffered content that caused the error
try {
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++),
});
Expand All @@ -4360,20 +4413,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
Expand Down Expand Up @@ -4489,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";
Expand Down Expand Up @@ -5211,6 +5255,10 @@ chat.openapi(completions, async (c) => {
requestedProvider,
usedModel,
initialRequestedModel,
unifiedFinishReason: getUnifiedFinishReason(
"upstream_error",
usedProvider,
),
});

// Log the error in the database
Expand Down Expand Up @@ -5507,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) ?? [];
Expand Down Expand Up @@ -5618,6 +5670,10 @@ chat.openapi(completions, async (c) => {
organizationId: project.organizationId,
projectId: apiKey.projectId,
apiKeyId: apiKey.id,
unifiedFinishReason: getUnifiedFinishReason(
finishReason,
usedProvider,
),
});
}

Expand Down Expand Up @@ -5885,6 +5941,10 @@ chat.openapi(completions, async (c) => {
usedModel,
initialRequestedModel,
cause: bodyReadCause,
unifiedFinishReason: getUnifiedFinishReason(
"upstream_error",
usedProvider,
),
});

const bodyTimeoutPluginIds = plugins?.map((p) => p.id) ?? [];
Expand Down
56 changes: 56 additions & 0 deletions apps/gateway/src/chat/tools/normalize-streaming-error.spec.ts
Original file line number Diff line number Diff line change
@@ -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: {");
});
});
Loading
Loading