Skip to content

Commit 487f219

Browse files
committed
fix(anthropic): align streaming behavior with OpenAI endpoint, fixes for #15 #17
- Flush SSE header on first upstream content event (thinking/text/tool_use) instead of waiting for text_delta — reasoning models' thinking phase left clients with zero bytes for 60s, triggering downstream first-byte timeout (context canceled) - Add event: ping keepalive after 15s of stream silence (Anthropic standard) - Preserve local delta token counts when finish event lacks usage (prevents false zero-output 429 on complete responses) - Non-stream: judge zero-output by actual content, not usage field alone - tool_result messages now precede text in converted user messages (fixes #15: Claude Code session resume 502 Tool result is missing); omit tool name when tool_use_id has no matching assistant tool_use - params.system space placeholder when no system prompt, prevents ~7.5K upstream default prompt injection (fixes #17); configurable via config.json emptySystemPlaceholder / CC_EMPTY_SYSTEM_PLACEHOLDER=false - cache_creation_input_tokens: emit 0 instead of null (schema compliance)
1 parent fcdb56a commit 487f219

1 file changed

Lines changed: 81 additions & 50 deletions

File tree

‎proxy.mjs‎

Lines changed: 81 additions & 50 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ function loadConfig() {
2323
useProviderModels: true,
2424
modelRefreshIntervalMs: 5 * 60 * 1000, // 5 minutes
2525
zdr: false,
26+
emptySystemPlaceholder: true, // 无 system prompt 时发空格占位,阻止 CC 上游注入 ~7.5K token 默认提示词(issue #17)
2627
};
2728

2829
const configPath = resolve(__dirname, 'config.json');
@@ -43,6 +44,7 @@ function loadConfig() {
4344
if (process.env.LOG_FILE) defaults.logFile = process.env.LOG_FILE;
4445
if (process.env.CC_USE_PROVIDER_MODELS) defaults.useProviderModels = process.env.CC_USE_PROVIDER_MODELS !== 'false';
4546
if (process.env.CMD_ZDR !== undefined) defaults.zdr = process.env.CMD_ZDR === '1';
47+
if (process.env.CC_EMPTY_SYSTEM_PLACEHOLDER) defaults.emptySystemPlaceholder = process.env.CC_EMPTY_SYSTEM_PLACEHOLDER !== 'false';
4648

4749
return defaults;
4850
}
@@ -505,6 +507,14 @@ function buildCcRequest(openaiReq) {
505507
// 条件字段
506508
if (systemPrompt) {
507509
body.params.system = systemPrompt;
510+
} else if (CFG.emptySystemPlaceholder) {
511+
// CC 上游在 params.system 缺省时会注入自身约 7.5K token 的默认提示词(进入
512+
// 默认上下文/前缀路径),既产生大量 cached tokens 又污染对话(模型会以为
513+
// 自己在 CC 的可执行目录里,见 issue #17)。发一个空格占位即可绕过,
514+
// 真机验证 prompt_tokens 从 7653 降到 85。
515+
// 默认开启;config.json 设 "emptySystemPlaceholder": false 或环境变量
516+
// CC_EMPTY_SYSTEM_PLACEHOLDER=false 可关闭(回到原生的缺省行为)。
517+
body.params.system = ' ';
508518
}
509519
if (temperature !== undefined) {
510520
body.params.temperature = temperature;
@@ -1267,10 +1277,13 @@ function buildAnthropicResponse(model, fullText, toolCalls, finishReason, usage,
12671277
stop_sequence: null,
12681278
usage: (() => {
12691279
normalizeUsage(usage || {});
1280+
// CC 未回报 usage 时按内容长度估算输出 token,避免客户端展示/记账为 0
1281+
const estOut = Math.max(1,
1282+
Math.ceil(((fullText || '').length + (thinkingText || '').length) / 4) + (toolCalls ? toolCalls.length * 20 : 0));
12701283
return {
12711284
input_tokens: usage?.inputTokens ?? 0,
1272-
output_tokens: usage?.outputTokens ?? 0,
1273-
cache_creation_input_tokens: usage?.inputTokenDetails?.cacheWriteTokens ?? null,
1285+
output_tokens: usage?.outputTokens || estOut,
1286+
cache_creation_input_tokens: usage?.inputTokenDetails?.cacheWriteTokens ?? 0,
12741287
cache_read_input_tokens: usage?.cachedInputTokens ?? 0,
12751288
};
12761289
})(),
@@ -1338,18 +1351,22 @@ function convertAnthropicToOpenAI(anthropicReq) {
13381351
}
13391352
}
13401353
if (textContent) {
1341-
openaiMessages.push({ role: 'user', content: textContent });
1354+
// 暂存,tool_result 优先入队:OpenAI 语义要求 tool 消息紧跟 assistant 的
1355+
// tool_calls,同一条 user 消息里的文本要排在 tool 结果之后
13421356
}
13431357
for (const tr of toolResults) {
13441358
const toolContent = typeof tr.content === 'string' ? tr.content
13451359
: Array.isArray(tr.content) ? tr.content.map(c => c.text || '').join('')
13461360
: String(tr.content || '');
1347-
openaiMessages.push({
1348-
role: 'tool',
1349-
tool_call_id: tr.tool_use_id,
1350-
name: toolNameFromId[tr.tool_use_id] || '',
1351-
content: toolContent,
1352-
});
1361+
// OpenAI 语义里 tool 消息的 name 是可选的;会话恢复等场景下 tool_use_id 可能
1362+
// 找不到对应 assistant tool_use(历史被客户端裁剪),此时不硬塞空 name,
1363+
// 避免 CC 上游报 "Tool result is missing"(issue #15)
1364+
const toolMsg = { role: 'tool', tool_call_id: tr.tool_use_id, content: toolContent };
1365+
if (toolNameFromId[tr.tool_use_id]) toolMsg.name = toolNameFromId[tr.tool_use_id];
1366+
openaiMessages.push(toolMsg);
1367+
}
1368+
if (textContent) {
1369+
openaiMessages.push({ role: 'user', content: textContent });
13531370
}
13541371
}
13551372
}
@@ -1566,15 +1583,9 @@ async function* createAnthropicSseTranslator(response, model, messageId, ctx) {
15661583
ctx.inputTokens = inputTokens;
15671584
ctx.outputTokens = outputTokens;
15681585
ctx.cachedInputTokens = cachedInputTokens;
1569-
} else {
1570-
inputTokens = 0;
1571-
outputTokens = 0;
1572-
cachedInputTokens = 0;
1573-
cacheWriteTokens = 0;
1574-
ctx.inputTokens = 0;
1575-
ctx.outputTokens = 0;
1576-
ctx.cachedInputTokens = 0;
15771586
}
1587+
// 上游未回报 usage 时保留本地按 delta 计数的估算值——清零会把有内容的
1588+
// 响应误判成零输出(触发 429)。未知字段保持原值即可。
15781589
break;
15791590
}
15801591

@@ -1596,6 +1607,12 @@ async function* createAnthropicSseTranslator(response, model, messageId, ctx) {
15961607
}
15971608
}
15981609

1610+
// 无论上游是否回报 usage,都把本地计数同步进 ctx(零输出判定与 message_delta 账单依赖它)
1611+
ctx.inputTokens = inputTokens;
1612+
ctx.outputTokens = outputTokens;
1613+
ctx.cachedInputTokens = cachedInputTokens;
1614+
ctx.cacheWriteTokens = cacheWriteTokens;
1615+
15991616
// Finalize — close pending text block, emit message_delta + message_stop
16001617
if (!hasError) {
16011618
const closeBlock = closeTextBlock();
@@ -1608,7 +1625,7 @@ async function* createAnthropicSseTranslator(response, model, messageId, ctx) {
16081625
yield `event: message_delta\ndata: ${JSON.stringify({
16091626
type: 'message_delta',
16101627
delta: { stop_reason: stopReason || 'end_turn' },
1611-
usage: { output_tokens: outputTokens, cache_read_input_tokens: cachedInputTokens, cache_creation_input_tokens: cacheWriteTokens || null, input_tokens: inputTokens },
1628+
usage: { output_tokens: outputTokens, cache_read_input_tokens: cachedInputTokens, cache_creation_input_tokens: cacheWriteTokens || 0, input_tokens: inputTokens },
16121629
})}\n\n`;
16131630

16141631
yield `event: message_stop\ndata: ${JSON.stringify({ type: 'message_stop' })}\n\n`;
@@ -1705,8 +1722,36 @@ async function handleMessages(req, res) {
17051722

17061723
if (stream) {
17071724
// ── 流式 Anthropic SSE ──
1708-
let started = false; // 延迟写 200 header,超时/output=0 时返回 JSON 429/502 让 SDK 自动重试
1725+
// 行为与 /v1/chat/completions 对齐:首个上游事件(thinking/text/tool_use)到达即
1726+
// 发 header——之前扣到 text_delta 才发,推理模型 thinking 阶段客户端收不到任何
1727+
// 字节,触发下游 60s 首字节超时(context canceled)。message_start 仍缓冲:
1728+
// 完全无输出时还能回 JSON 429/502 让 SDK 自动重试(同 chat 端点)。
1729+
let started = false;
17091730
const buf = [];
1731+
const SSE_HEADERS = {
1732+
'Content-Type': 'text/event-stream',
1733+
'Cache-Control': 'no-cache',
1734+
'Connection': 'keep-alive',
1735+
'X-Accel-Buffering': 'no',
1736+
};
1737+
const flushBuf = () => {
1738+
if (!started) {
1739+
res.writeHead(200, SSE_HEADERS);
1740+
started = true;
1741+
}
1742+
for (const ev of buf) { try { res.write(ev); } catch {} }
1743+
buf.length = 0;
1744+
};
1745+
1746+
// 心跳:等价于 chat 端点的 ': keepalive'——chat 在每轮读到静默事件时发注释行,
1747+
// Anthropic 翻译器会吞掉 signal 事件,这里改用空闲计时发 ping(Anthropic 标准
1748+
// 事件,官方 SDK 会忽略),覆盖上游排队/长 thinking 的静默窗口
1749+
let lastSentAt = Date.now();
1750+
const heartbeat = setInterval(() => {
1751+
if (started && !aborted && !res.writableEnded && Date.now() - lastSentAt > 15000) {
1752+
try { res.write('event: ping\ndata: {"type":"ping"}\n\n'); lastSentAt = Date.now(); } catch {}
1753+
}
1754+
}, 5000);
17101755

17111756
let ctx;
17121757
try {
@@ -1715,22 +1760,14 @@ async function handleMessages(req, res) {
17151760
const generator = createAnthropicSseTranslator(ccResponse, model, messageId, ctx);
17161761
for await (const event of generator) {
17171762
if (aborted) break;
1718-
if (!started) {
1719-
buf.push(event);
1720-
// 确认有真实内容后才发 200 header
1721-
if (event.includes('"text_delta"') || event.includes('"tool_use"')) {
1722-
res.writeHead(200, {
1723-
'Content-Type': 'text/event-stream',
1724-
'Cache-Control': 'no-cache',
1725-
'Connection': 'keep-alive',
1726-
'X-Accel-Buffering': 'no',
1727-
});
1728-
started = true;
1729-
for (const ev of buf) res.write(ev);
1730-
buf.length = 0;
1731-
}
1763+
if (!started && !event.startsWith('event: message_start')) {
1764+
flushBuf();
1765+
}
1766+
if (started) {
1767+
try { res.write(event); } catch {}
1768+
lastSentAt = Date.now();
17321769
} else {
1733-
res.write(event);
1770+
buf.push(event);
17341771
}
17351772
}
17361773

@@ -1745,26 +1782,16 @@ async function handleMessages(req, res) {
17451782
ctx.upstreamError.body.error.message,
17461783
);
17471784
}
1785+
// started 时 error 事件已在循环中经 SSE 下发,按规范 error 事件即终结
17481786
} else if (ctx.outputTokens === 0) {
17491787
try { abortController.abort(); } catch {}
17501788
if (!started) {
17511789
sendAnthropicError(res, 429, 'rate_limit_error', 'Empty response from upstream (zero output tokens)', 10);
17521790
return;
17531791
}
1754-
for (const ev of buf) { try { res.write(ev); } catch {} }
1755-
buf.length = 0;
1792+
flushBuf();
17561793
} else {
1757-
if (!started) {
1758-
res.writeHead(200, {
1759-
'Content-Type': 'text/event-stream',
1760-
'Cache-Control': 'no-cache',
1761-
'Connection': 'keep-alive',
1762-
'X-Accel-Buffering': 'no',
1763-
});
1764-
started = true;
1765-
}
1766-
for (const ev of buf) res.write(ev);
1767-
buf.length = 0;
1794+
flushBuf();
17681795
}
17691796
}
17701797
} catch (e) {
@@ -1814,6 +1841,8 @@ async function handleMessages(req, res) {
18141841
} catch {}
18151842
}
18161843
}
1844+
} finally {
1845+
clearInterval(heartbeat);
18171846
}
18181847

18191848
if (!res.writableEnded) res.end();
@@ -1855,7 +1884,7 @@ async function handleMessages(req, res) {
18551884
case 'finish':
18561885
lastCcEvent = event.type;
18571886
finishReason = mapFinishReason(event.finishReason || 'stop');
1858-
if (event.totalUsage) usage = event.totalUsage;
1887+
if (event.totalUsage || event.usage) usage = event.totalUsage || event.usage;
18591888
break;
18601889
case 'error':
18611890
lastCcEvent = event.type;
@@ -1893,8 +1922,9 @@ async function handleMessages(req, res) {
18931922
return;
18941923
}
18951924

1896-
// 输出 token 为 0 时记为错误,避免下游异常计费
1897-
if ((usage?.outputTokens ?? 0) === 0) {
1925+
// 零输出判定改为按实际内容:上游偶发不回 totalUsage 时,旧逻辑(usage?.outputTokens ?? 0 === 0)
1926+
// 会把有完整文本的响应误杀成 429
1927+
if (!fullText && !thinkingText && !toolCalls) {
18981928
try { if (!abortController.signal.aborted) abortController.abort(); } catch {}
18991929
sendAnthropicError(res, 429, 'rate_limit_error', 'Empty response from upstream (zero output tokens)', 10);
19001930
return;
@@ -2052,6 +2082,7 @@ server.listen(CFG.port, CFG.host, () => {
20522082
models: MODELS.length,
20532083
session: '12h + 1h jitter, per API key',
20542084
zdr: CFG.zdr ? 'enabled (x-cmd-zdr: 1 on generation/init requests)' : 'off (CMD_ZDR=1 or per-request x-cmd-zdr: 1 to enable)',
2085+
emptySystemPlaceholder: CFG.emptySystemPlaceholder ? 'on (space placeholder for requests without system prompt, issue #17)' : 'off',
20552086
logFile: CFG.logFile || '(console only)',
20562087
});
20572088
if (!CFG.apiKey) {

0 commit comments

Comments
 (0)