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
4 changes: 2 additions & 2 deletions src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ export interface Config {
ccVersion: string;
logLevel: string;
corsOrigin: string;
/** Max wall-clock ms for upstream to send response headers + first byte. */
/** Per-attempt deadline for upstream headers and any non-2xx error body. */
upstreamTimeoutMs: number;
/** Max ms between consecutive chunks during streaming. 0 = disabled. */
idleTimeoutMs: number;
Expand Down Expand Up @@ -92,7 +92,7 @@ export function loadConfig(): Config {
const corsOrigin = process.env.CORS_ORIGIN ?? "*";

// Upstream timeouts. The connection timeout covers the wall-clock time
// until the upstream returns response headers + first byte — bump it for
// until the upstream returns headers (and consumes any error body) — bump it for
// slow reasoning models. The idle timeout catches stalled streams where
// the upstream opened the connection but stopped sending chunks
// mid-response (e.g. tool call hung on the upstream side). Set
Expand Down
61 changes: 34 additions & 27 deletions src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import { getProxyVersion } from "@/version.js";
import {
validateOpenAIChatRequest,
validateAnthropicRequest,
validateCountTokensRequest,
ValidationError,
} from "@/translate/validation.js";
import type { AnthropicRequest, AnthropicSSERecord } from "@/translate/anthropic-types.js";
Expand Down Expand Up @@ -159,24 +160,17 @@ function corsHeaders(): Record<string, string> {
// Helpers
// ──────────────────────────────────────────

function abortOnClientDisconnect(
req: http.IncomingMessage,
res: http.ServerResponse,
): AbortController {
function abortOnClientDisconnect(res: http.ServerResponse): AbortController {
const abort = new AbortController();
req.on("close", () => {
// IncomingMessage.close marks a completed upload, not a lost response client.
const onClose = (): void => {
if (!res.writableEnded) abort.abort();
});
};
res.once("close", onClose);
if (res.destroyed) onClose();
return abort;
}

function destroyStreamOnClientDisconnect(
req: http.IncomingMessage,
stream: NodeJS.ReadableStream,
): void {
req.on("close", () => (stream as Readable).destroy());
}

/**
* Write a chunk to `res`, returning a Promise that resolves once the
* underlying socket has drained (when backpressure applies). Returns
Expand All @@ -186,7 +180,17 @@ function writeSSE(res: http.ServerResponse, chunk: string): Promise<boolean> {
if (res.writableEnded || res.destroyed) return Promise.resolve(false);
if (res.write(chunk)) return Promise.resolve(true);
return new Promise((resolve) => {
res.once("drain", () => resolve(!res.writableEnded && !res.destroyed));
const settle = (writable: boolean): void => {
res.off("drain", onDrain);
res.off("close", onClose);
res.off("error", onClose);
resolve(writable);
};
const onDrain = (): void => settle(!res.writableEnded && !res.destroyed);
const onClose = (): void => settle(false);
res.once("drain", onDrain);
res.once("close", onClose);
res.once("error", onClose);
});
}

Expand Down Expand Up @@ -214,7 +218,7 @@ async function pumpStream(
// Encoder blew up — turn it into a stream error so the catch below
// handles it uniformly instead of crashing the proxy.
(stream as Readable).destroy(err as Error);
break;
throw err;
}
for (const chunk of chunks) {
if (!(await writable(chunk))) return;
Expand Down Expand Up @@ -322,7 +326,7 @@ async function handleChatCompletions(

const ccBody = toCCRequest(openAIReq);

const abort = abortOnClientDisconnect(req, res);
const abort = abortOnClientDisconnect(res);

try {
const result = await sendToCC(
Expand Down Expand Up @@ -366,13 +370,8 @@ async function handleChatCompletions(
res.write(formatSSEDone());
res.end();
}
// No destroyStreamOnClientDisconnect here — by the time pumpStream
// returns the stream has already ended or errored, so the call would
// be a no-op. Mid-stream disconnects are handled by the abort signal
// (see abortOnClientDisconnect + nodeReaderToStream's abortSignal
// listener).
// The response abort signal covers both streaming and JSON clients.
} else {
destroyStreamOnClientDisconnect(req, stream);
const events = await collectEvents(stream);
const response = buildNonStreamingResponse(events, model, encoder.id);
sendJson(res, 200, response);
Expand Down Expand Up @@ -418,7 +417,7 @@ async function handleMessages(req: http.IncomingMessage, res: http.ServerRespons
const encoder = new AnthropicStreamEncoder(model);
const ccBody = anToCCRequest(anthropicReq);

const abort = abortOnClientDisconnect(req, res);
const abort = abortOnClientDisconnect(res);

try {
const result = await sendToCC(
Expand Down Expand Up @@ -462,10 +461,8 @@ async function handleMessages(req: http.IncomingMessage, res: http.ServerRespons
},
);
if (!res.writableEnded && !res.destroyed) res.end();
// No destroyStreamOnClientDisconnect here — see OpenAI streaming path
// for rationale (abort signal already covers mid-stream disconnect).
// The response abort signal already covers mid-stream disconnects.
} else {
destroyStreamOnClientDisconnect(req, stream);
const events = await collectEvents(stream);
const response = buildAnthropicResponse(events, model, encoder.messageId);
res.writeHead(200, { "Content-Type": "application/json", ...corsHeaders() });
Expand Down Expand Up @@ -494,7 +491,17 @@ async function handleCountTokens(
);
}

const body = rawBody as Record<string, unknown>;
let body: Record<string, unknown>;
try {
body = validateCountTokensRequest(rawBody);
} catch (err) {
return sendAnthropicError(
res,
400,
"invalid_request_error",
err instanceof ValidationError ? err.message : "Invalid request body",
);
}

const parts: string[] = [];
if (typeof body.system === "string") parts.push(body.system);
Expand Down
104 changes: 65 additions & 39 deletions src/translate/anthropic.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import crypto from "node:crypto";
import { toolArgumentSuffix } from "./tool-arguments.js";
import type {
AnthropicRequest,
AnthropicContentBlock,
Expand Down Expand Up @@ -245,8 +246,12 @@ export function toCCRequest(
export class AnthropicStreamEncoder {
readonly messageId: string;
private blockIndex = 0;
private currentBlockIndex = 0;
private currentBlockType: "text" | "thinking" | "tool_use" | null = null;
private currentToolCallId: string | null = null;
private readonly toolBlocks = new Map<
string,
{ index: number; arguments: string; closed: boolean }
>();
private pendingStart: CCEvent | null = null;
private started = false;
private pinged = false;
Expand Down Expand Up @@ -281,6 +286,7 @@ export class AnthropicStreamEncoder {
records.push(this.makeMessageStart(0));
}
this.closeCurrentBlock(records);
this.closeToolBlocks(records);
records.push({
event: "error",
data: { type: "error", error: { type: "api_error", message: msg } },
Expand Down Expand Up @@ -328,6 +334,7 @@ export class AnthropicStreamEncoder {
}

this.closeCurrentBlock(records);
this.closeToolBlocks(records);

const finishReason = (event.data.finishReason as string) ?? "stop";
const usage = extractUsage(event.data as Record<string, unknown>);
Expand Down Expand Up @@ -381,21 +388,18 @@ export class AnthropicStreamEncoder {
case "tool-call-delta": {
const tcId = (event.data.toolCallId as string) ?? "";
const tcName = (event.data.name as string) ?? "";
if (this.currentBlockType !== "tool_use" || this.currentToolCallId !== tcId) {
this.closeCurrentBlock(records);
this.ensureBlockOpenWith(records, "tool_use", {
type: "tool_use",
id: tcId,
name: tcName,
input: {},
});
this.currentToolCallId = tcId;
}
const block = this.ensureToolBlock(records, tcId, tcName);
if (block.closed) throw new Error("Inconsistent upstream tool arguments");
const args = (event.data.arguments as string) ?? "";
block.arguments += args;
records.push(
this.makeDelta({
type: "input_json_delta",
partial_json: (event.data.arguments as string) ?? "",
}),
this.makeDelta(
{
type: "input_json_delta",
partial_json: args,
},
block.index,
),
);
break;
}
Expand All @@ -406,33 +410,51 @@ export class AnthropicStreamEncoder {
const input = event.data.input ?? event.data.arguments;
const argsStr =
typeof input === "string" ? input : input != null ? JSON.stringify(input) : "";
// If the upstream streamed deltas for this tool call first and then
// sent the final `tool-call` event with the same id, reuse the block
// it already opened instead of creating a duplicate `tool_use`.
if (this.currentBlockType === "tool_use" && this.currentToolCallId === tcId) {
if (argsStr) {
records.push(this.makeDelta({ type: "input_json_delta", partial_json: argsStr }));
}
break;
}
this.closeCurrentBlock(records);
this.ensureBlockOpenWith(records, "tool_use", {
type: "tool_use",
id: tcId,
name: tcName,
input: {},
});
if (argsStr) {
records.push(this.makeDelta({ type: "input_json_delta", partial_json: argsStr }));
const block = this.ensureToolBlock(records, tcId, tcName);
const suffix = toolArgumentSuffix(block.arguments, argsStr);
if (suffix) {
if (block.closed) throw new Error("Inconsistent upstream tool arguments");
records.push(
this.makeDelta({ type: "input_json_delta", partial_json: suffix }, block.index),
);
block.arguments += suffix;
}
this.closeCurrentBlock(records);
this.closeToolBlock(records, block);
break;
}
}

return records;
}

private ensureToolBlock(records: AnthropicSSERecord[], id: string, name: string) {
const existing = this.toolBlocks.get(id);
if (existing) return existing;
this.closeCurrentBlock(records);
this.ensureBlockOpenWith(records, "tool_use", { type: "tool_use", id, name, input: {} });
const block = { index: this.currentBlockIndex, arguments: "", closed: false };
this.toolBlocks.set(id, block);
// Tool blocks have independent lifetimes: interleaved deltas retain indices.
this.currentBlockType = null;
return block;
}

private closeToolBlock(
records: AnthropicSSERecord[],
block: { index: number; closed: boolean },
): void {
if (block.closed) return;
records.push({
event: "content_block_stop",
data: { type: "content_block_stop", index: block.index },
});
block.closed = true;
}

private closeToolBlocks(records: AnthropicSSERecord[]): void {
for (const block of this.toolBlocks.values()) this.closeToolBlock(records, block);
}

private ensureBlockOpen(
records: AnthropicSSERecord[],
type: "text" | "thinking" | "tool_use",
Expand All @@ -450,11 +472,12 @@ export class AnthropicStreamEncoder {
block: ContentBlockStartShape,
): void {
this.currentBlockType = type;
this.currentBlockIndex = this.blockIndex++;
records.push({
event: "content_block_start",
data: {
type: "content_block_start",
index: this.blockIndex,
index: this.currentBlockIndex,
content_block: block,
},
});
Expand All @@ -477,17 +500,16 @@ export class AnthropicStreamEncoder {

records.push({
event: "content_block_stop",
data: { type: "content_block_stop", index: this.blockIndex },
data: { type: "content_block_stop", index: this.currentBlockIndex },
});

this.blockIndex++;
this.currentBlockType = null;
}

private makeDelta(delta: DeltaShape): AnthropicSSERecord {
private makeDelta(delta: DeltaShape, index = this.currentBlockIndex): AnthropicSSERecord {
return {
event: "content_block_delta",
data: { type: "content_block_delta", index: this.blockIndex, delta },
data: { type: "content_block_delta", index, delta },
};
}

Expand Down Expand Up @@ -529,6 +551,7 @@ export class AnthropicStreamEncoder {
this.started = true;
}
this.closeCurrentBlock(records);
this.closeToolBlocks(records);
records.push({
event: "message_delta",
data: {
Expand All @@ -555,6 +578,9 @@ export function buildAnthropicResponse(

for (const event of events) {
switch (event.type) {
case "error":
// Do not turn failed generations (or private diagnostics) into content.
throw new Error("CC upstream generation failed");
case "text-delta":
textContent += (event.data.text as string) ?? "";
break;
Expand Down
Loading
Loading