From 62066a705d41f483689e94a9340bd0d43150b707 Mon Sep 17 00:00:00 2001 From: xelr233 Date: Tue, 15 Sep 2026 23:06:17 +0800 Subject: [PATCH 1/4] =?UTF-8?q?test:=20=E6=9C=80=E5=B0=8F=E6=B5=8B?= =?UTF-8?q?=E8=AF=95=E8=84=9A=E6=89=8B=E6=9E=B6=EF=BC=88helpers=20+=20npm?= =?UTF-8?q?=20test=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 本 PR 只带这一组测试所需的脚手架:test/helpers.mjs 与 package.json 的 test 脚本。 helpers.mjs 与 PR #34 / #39 中的文件**逐字节一致**(blob 3342bb82), 所以三份先后合并都不会冲突 —— git 对「两侧新增同一路径且内容相同」不视为冲突。 --- package.json | 1 + test/helpers.mjs | 152 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 153 insertions(+) create mode 100644 test/helpers.mjs diff --git a/package.json b/package.json index 8227668..e5bc345 100644 --- a/package.json +++ b/package.json @@ -8,6 +8,7 @@ "scripts": { "start": "node proxy.mjs", "dev": "node --watch proxy.mjs", + "test": "node --test test/*.test.mjs", "docker:build": "docker build -t commandcode-proxy:latest .", "docker:build:multi": "docker buildx build --platform linux/amd64,linux/arm64 -t commandcode-proxy:latest ." }, diff --git a/test/helpers.mjs b/test/helpers.mjs new file mode 100644 index 0000000..3342bb8 --- /dev/null +++ b/test/helpers.mjs @@ -0,0 +1,152 @@ +// 测试用 mock 上游 + 代理进程管理。 +// 全部走 loopback,不需要真 key、不访问 Command Code 或 npm registry。 +import http from 'node:http'; +import { spawn } from 'node:child_process'; +import { setTimeout as sleep } from 'node:timers/promises'; +import { mkdtempSync, copyFileSync, existsSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join, dirname } from 'node:path'; +import { fileURLToPath } from 'node:url'; + +export const REPO = dirname(fileURLToPath(import.meta.url)).replace(/[/\\]test$/, ''); + +// ── 挂起保护 ────────────────────────────────────────────── +// 不能用 --test-timeout:Node 18 没有该选项(20.11 才加入),加了会让 +// engines 下限直接跑不起来。改用启动一个 unref 的定时器,进程若因泄漏 +// 的 socket / 未 await 的句柄而无法退出,到点强制退出并说明原因。 +// 正常结束时定时器被 unref,不阻止退出。 +const HANG_GUARD_MS = Number(process.env.CC_TEST_HANG_GUARD_MS ?? 120000); +const hangGuard = setTimeout(() => { + console.error('[test] 超时未退出:疑似有 server/socket 未关闭(' + + '检查每个测试是否都在 finally 里 await close())。强制退出。'); + process.exit(1); +}, HANG_GUARD_MS); +hangGuard.unref?.(); + +// 取一个当前空闲的端口:让内核分配(listen 0)后立刻释放。 +// 不能用 pid 派生区间 —— node --test 各文件并行,pid 取模会在不同 pid 间 +// 映射到同一区间(如 pid 100 与 pid 600 同桶),进而偶发 EADDRINUSE。 +// 内核分配把冲突面缩到「释放到重新占用」之间的极小窗口,调用方另有重试兜底。 +import net from 'node:net'; +export async function allocPort() { + return await new Promise((resolve, reject) => { + const srv = net.createServer(); + srv.once('error', reject); + srv.listen(0, '127.0.0.1', () => { + const { port } = srv.address(); + srv.close(() => resolve(port)); + }); + }); +} + +/** + * 关闭一个 http server,且**保证有界**。 + * server.close() 只停止接受新连接,会一直等到既有连接结束 —— 若有 keep-alive + * 或未被对端关闭的 socket,它会永远挂着,把 CI 拖到 job 超时。 + * (实际发生过:Fork 测试漏写一个 await 导致 mock 泄漏,三个矩阵 job 全部 + * 空转 10 分钟后被取消。)故先强制断开所有连接,再 close,并叠加兜底超时。 + */ +export async function closeServer(server, timeoutMs = 3000) { + if (!server || !server.listening) return; + try { server.closeAllConnections?.(); } catch {} + await Promise.race([new Promise(r => server.close(r)), sleep(timeoutMs)]); + try { server.closeAllConnections?.(); } catch {} +} + +/** 启动一个 mock 上游。ndjson 为要回给代理的 CC NDJSON 行数组。 */ +export async function startMockUpstream(opts = {}) { + const port = await allocPort(); + const seen = []; + const server = http.createServer((req, res) => { + const chunks = []; + req.on('data', c => chunks.push(c)); + req.on('end', async () => { + const raw = Buffer.concat(chunks).toString('utf8'); + seen.push({ url: req.url, method: req.method, headers: req.headers, raw }); + if (opts.onRequest) await opts.onRequest(req, res, seen[seen.length - 1]); + if (res.writableEnded) return; + const status = opts.status ?? 200; + if (status !== 200) { + res.writeHead(status, { 'Content-Type': 'application/json' }); + res.end(opts.errorBody || JSON.stringify({ error: { message: 'mock error' } })); + return; + } + res.writeHead(200, { 'Content-Type': 'text/event-stream' }); + for (const line of opts.ndjson ?? [ + '{"type":"text-start"}', + '{"type":"text-delta","text":"hello"}', + '{"type":"text-end"}', + '{"type":"finish-step","finishReason":"stop","usage":{"inputTokens":9,"outputTokens":3}}', + '{"type":"finish","finishReason":"stop","totalUsage":{"inputTokens":9,"outputTokens":3,"cachedInputTokens":0}}', + ]) res.write(line + '\n'); + res.end(); + }); + }); + await new Promise(r => server.listen(port, '127.0.0.1', r)); + return { port, seen, close: () => closeServer(server), + // 最后一次 /alpha/generate 的请求体(wire 层断言的主要入口) + lastGenerate: () => { + const g = seen.filter(s => s.url === '/alpha/generate').pop(); + return g ? { raw: g.raw, body: JSON.parse(g.raw), headers: g.headers } : null; + }, + generateCount: () => seen.filter(s => s.url === '/alpha/generate').length }; +} + +/** 在临时 cwd 中启动代理(复刻真实部署:proxy.mjs 与 config.json 同目录)。 */ +export async function startProxy({ upstreamPort, env = {}, cwd } = {}) { + const port = await allocPort(); + const logs = []; + // 自建的临时工作目录用完必须删;调用方传了 cwd 则由调用方负责。 + const ownWorkdir = cwd === undefined; + const workdir = cwd ?? mkdtempSync(join(tmpdir(), 'ccp-test-')); + copyFileSync(join(REPO, 'proxy.mjs'), join(workdir, 'proxy.mjs')); + if (!existsSync(join(workdir, 'config.json'))) { + copyFileSync(join(REPO, 'config.json'), join(workdir, 'config.json')); + } + const child = spawn(process.execPath, ['proxy.mjs'], { + cwd: workdir, + env: { ...process.env, PORT: String(port), HOST: '127.0.0.1', + CC_API_BASE: 'http://127.0.0.1:' + upstreamPort, + CC_USE_PROVIDER_MODELS: 'false', // 不访问 /provider/v1/models + ...env }, + stdio: ['ignore', 'pipe', 'pipe'], + }); + child.stdout.on('data', d => logs.push(d.toString())); + child.stderr.on('data', d => logs.push(d.toString())); + + const base = 'http://127.0.0.1:' + port; + let up = false; + for (let i = 0; i < 80; i++) { + if (child.exitCode !== null) break; + try { const r = await fetch(base + '/health'); if (r.ok) { up = true; break; } } catch {} + await sleep(125); + } + if (!up) { child.kill(); throw new Error('proxy did not start:\n' + logs.join('')); } + + return { + port, base, child, logs: () => logs.join(''), + get: (path, init) => fetch(base + path, init), + post: (path, body, headers = {}) => fetch(base + path, { + method: 'POST', headers: { 'Content-Type': 'application/json', ...headers }, + body: typeof body === 'string' ? body : JSON.stringify(body), + }), + kill: () => new Promise(r => { + child.once('exit', () => { + if (ownWorkdir) { try { rmSync(workdir, { recursive: true, force: true }); } catch {} } + r(); + }); + child.kill(); + setTimeout(() => { + if (ownWorkdir) { try { rmSync(workdir, { recursive: true, force: true }); } catch {} } + r(); + }, 2000); + }), + }; +} + +/** 一次性搭好 mock 上游 + 代理。 */ +export async function setup(opts = {}) { + const mock = await startMockUpstream(opts); + const proxy = await startProxy({ upstreamPort: mock.port, env: opts.env, cwd: opts.cwd }); + return { mock, proxy, async close() { await proxy.kill(); await mock.close(); } }; +} From cfac3b04ed26496b84946145f3992f9da6f132b1 Mon Sep 17 00:00:00 2001 From: xelr233 Date: Tue, 15 Sep 2026 23:08:01 +0800 Subject: [PATCH 2/4] =?UTF-8?q?fix:=20=E9=87=87=E7=BA=B3=E4=B8=8A=E6=B8=B8?= =?UTF-8?q?=20error=20=E4=BA=8B=E4=BB=B6=E8=87=AA=E5=B8=A6=E7=9A=84=20stat?= =?UTF-8?q?usCode=EF=BC=8C=E5=B9=B6=E8=A1=A5=E9=BD=90=E6=97=A0=E5=86=85?= =?UTF-8?q?=E5=AE=B9=E4=BA=8B=E4=BB=B6=E7=9A=84=E9=9D=99=E9=BB=98=E5=88=97?= =?UTF-8?q?=E8=A1=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 线上故障排查中发现的两个可观测性问题。 【一】mapCcEventError 丢掉 error.statusCode CLI 的 readStreamErrorEvent 读的就是这个字段,取值链是 parseEmbeddedErrorJSON(message)?.status ?? error.statusCode ?? null 原实现只看 message 里的 "" 前缀,statusCode 一律被丢掉 → 一律塌成 502。 后果:上游报 429/503(限流、容量)时我们回 502 upstream_error —— 客户端不按限流退避,监控也把它错归类成后端故障。 线上那条 "The request limited providers for this model and they are currently at capacity..." 就可能因此被记成 502 而不是 429 (按 CLI 的 isStreamErrorRetryable,该消息不含任何 terminal 标记 premium_credits_exhausted / model_not_in_plan / insufficient credits, 所以它是可重试的)。 改法:采纳 error.statusCode("" 前缀仍优先,与 CLI 一致), 返回值增加 reportedStatus 以区分"上游报的"与"我们映射后的"; 四个 CC error 日志点改为先映射再记日志,并打出 upstreamStatus / upstreamRetryable / code / mappedTo —— 与作者 78353d9 对 mapCcError 的处理保持一致。 【二】无内容事件的静默列表不全 上游每个响应都会发一串不携带内容的事件(text-start / text-end / start / start-step / reasoning-start / reasoning-end / finish-step / provider-metadata / tool-input-start|delta|end / tool-error)。三条非流式路径里: · OpenAI / Anthropic 非流式:缺 text-start / start / start-step / reasoning-start / finish-step · Responses 非流式:**一个静默列表都没有** 于是 journalctl 被 Unknown CC event type 刷屏,真正的错误被淹没。 改法:三处统一成同一份静默列表;default 仍保留警告,真正没见过的类型照旧留痕。 测试:新增 test/connection-lifecycle.test.mjs 5 条(statusCode 映射 4 条 + 三协议 × 流式/非流式不产生噪音 1 条)。 --- proxy.mjs | 78 ++++++++++++++++++++++++--- test/connection-lifecycle.test.mjs | 87 ++++++++++++++++++++++++++++++ 2 files changed, 157 insertions(+), 8 deletions(-) create mode 100644 test/connection-lifecycle.test.mjs diff --git a/proxy.mjs b/proxy.mjs index e69ffae..eef0b0f 100644 --- a/proxy.mjs +++ b/proxy.mjs @@ -799,8 +799,16 @@ function createSseTranslator(model, completionId, created) { case 'error': { const msg = event.error?.message || event.message || 'Unknown error'; - log('warn', 'CC stream error', { message: msg }); this.upstreamError = mapCcEventError(event); + // 先映射再记日志,并把上游自带的状态/可重试性一并打出 —— + // 排查容量/限流类问题时,真正需要的就是这两个字段 + log('warn', 'CC stream error', { + message: msg, + upstreamStatus: this.upstreamError.reportedStatus, + upstreamRetryable: event.error?.isRetryable, + code: this.upstreamError.code, + mappedTo: this.upstreamError.status, + }); // Don't emit a finish_reason chunk — let the natural stream termination // handle it. Otherwise a subsequent finish(tool_calls) would be ignored // by downstream agent loops that stop at the first finish_reason. @@ -924,8 +932,17 @@ function mapCcError(ccStatus, ccBody) { function mapCcEventError(event) { const message = event.error?.message || event.message || 'Unknown CC error'; const code = event.error?.code || event.code || null; + // 上游 error 事件除了 message 还可能自带 statusCode / isRetryable —— + // CLI 的 readStreamErrorEvent 读的正是这两个字段,取值链是 + // parseEmbeddedErrorJSON(message)?.status ?? error.statusCode ?? null + // 原实现只看 message 里的 "" 前缀,statusCode 一律被丢掉, + // 于是 429 / 503 这类「该退避重试」的信号在代理这一层被抹平成 502「服务端错误」: + // 客户端不再按限流退避,监控也会把它错误归类成后端故障。 const statusMatch = message.match(/^<(\d{3})>/); - const ccStatus = statusMatch ? Number(statusMatch[1]) : 502; + const reportedStatus = statusMatch + ? Number(statusMatch[1]) + : (Number.isInteger(event.error?.statusCode) ? event.error.statusCode : null); + const ccStatus = reportedStatus ?? 502; const mapped = CC_STATUS_MAP[ccStatus] || { status: 502, type: 'upstream_error' }; // 与 mapCcError 保持一致:终态为 429 时带上 retry_after, @@ -934,11 +951,13 @@ function mapCcEventError(event) { return { status: 429, code, + reportedStatus, body: { error: { message, type: 'rate_limit_error', ...(code ? { code } : {}) }, retry_after: 30 }, }; } - return { status: mapped.status, code, body: { error: { message, type: mapped.type, ...(code ? { code } : {}) } } }; + return { status: mapped.status, code, reportedStatus, + body: { error: { message, type: mapped.type, ...(code ? { code } : {}) } } }; } // ── HTTP 请求处理 ────────────────────────────────── @@ -1363,10 +1382,23 @@ async function handleChatCompletions(req, res) { break; case 'error': lastCcEvent = event.type; - log('warn', 'CC stream error (non-stream)', { message: event.error?.message || event.message }); upstreamError = mapCcEventError(event); + log('warn', 'CC stream error (non-stream)', { + message: event.error?.message || event.message, + upstreamStatus: upstreamError.reportedStatus, + upstreamRetryable: event.error?.isRetryable, + code: upstreamError.code, + mappedTo: upstreamError.status, + }); break; - case 'reasoning-end': case 'provider-metadata': case 'tool-input-start': case 'tool-input-delta': case 'tool-input-end': case 'tool-error': case 'text-end': + // 无内容的事件:与流式翻译器的静默列表保持一致。 + // text-start / start / start-step / reasoning-start 原先只在流式路径被识别, + // 非流式路径会掉进 default 打成 'Unknown CC event type' —— 上游每个响应都会发, + // 于是线上刷屏。它们本身不携带内容(内容在 text-delta),纯粹是噪音。 + case 'text-start': case 'text-end': case 'start': case 'start-step': + case 'reasoning-start': case 'reasoning-end': case 'finish-step': + case 'provider-metadata': case 'tool-input-start': case 'tool-input-delta': case 'tool-input-end': + case 'tool-error': // Silent - no user-visible content break; default: @@ -2178,10 +2210,23 @@ async function handleMessages(req, res) { break; case 'error': lastCcEvent = event.type; - log('warn', 'CC error (Anthropic non-stream)', { message: event.error?.message || event.message }); upstreamError = mapCcEventError(event); + log('warn', 'CC error (Anthropic non-stream)', { + message: event.error?.message || event.message, + upstreamStatus: upstreamError.reportedStatus, + upstreamRetryable: event.error?.isRetryable, + code: upstreamError.code, + mappedTo: upstreamError.status, + }); break; - case 'reasoning-end': case 'provider-metadata': case 'tool-input-start': case 'tool-input-delta': case 'tool-input-end': case 'tool-error': case 'text-end': + // 无内容的事件:与流式翻译器的静默列表保持一致。 + // text-start / start / start-step / reasoning-start 原先只在流式路径被识别, + // 非流式路径会掉进 default 打成 'Unknown CC event type' —— 上游每个响应都会发, + // 于是线上刷屏。它们本身不携带内容(内容在 text-delta),纯粹是噪音。 + case 'text-start': case 'text-end': case 'start': case 'start-step': + case 'reasoning-start': case 'reasoning-end': case 'finish-step': + case 'provider-metadata': case 'tool-input-start': case 'tool-input-delta': case 'tool-input-end': + case 'tool-error': // Silent - no user-visible content break; default: @@ -2921,8 +2966,25 @@ async function handleResponses(req, res) { break; case 'error': lastCcEvent = event.type; - log('warn', 'CC stream error (non-stream)', { message: event.error ? event.error.message : event.message }); upstreamError = mapCcEventError(event); + log('warn', 'CC stream error (non-stream)', { + message: event.error ? event.error.message : event.message, + upstreamStatus: upstreamError.reportedStatus, + upstreamRetryable: event.error?.isRetryable, + code: upstreamError.code, + mappedTo: upstreamError.status, + }); + break; + // 无内容的事件:与流式翻译器以及另两条非流式路径保持一致。 + // 这条路径原先**没有静默列表**,于是上游每个响应都会发的一串无内容事件 + //(text-start / text-end / start / start-step / reasoning-start / reasoning-end / + // provider-metadata / tool-input-* / tool-error)全部掉进 default 打成 + // 'Unknown CC event type',线上刷屏、把真正的错误淹掉。 + case 'text-start': case 'text-end': case 'start': case 'start-step': + case 'reasoning-start': case 'reasoning-end': case 'finish-step': + case 'provider-metadata': case 'tool-input-start': case 'tool-input-delta': case 'tool-input-end': + case 'tool-error': + // Silent - no user-visible content break; default: log('warn', 'Unknown CC event type', { type: event.type }); diff --git a/test/connection-lifecycle.test.mjs b/test/connection-lifecycle.test.mjs new file mode 100644 index 0000000..fd3cb15 --- /dev/null +++ b/test/connection-lifecycle.test.mjs @@ -0,0 +1,87 @@ +// 连接生命周期与错误可观测性。 +// 这组用例来自一次真实线上故障的排查(间歇性反代 502 / 客户端 connection error): +// ① 上游 error 事件自带 statusCode,被丢掉后 429/503 一律塌成 502 +// ② 无内容事件的静默列表不全,Response 非流式那条**根本没有**,日志被刷屏 +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { setup } from './helpers.mjs'; + +const AUTH = { Authorization: 'Bearer user_test' }; +const CHAT = { model: 'm', messages: [{ role: 'user', content: 'hi' }] }; + +// ── ① error 事件自带的 statusCode 必须被采纳 ────────────────── +// CLI 的 readStreamErrorEvent 读的就是 error.statusCode / error.isRetryable, +// 取值链是 parseEmbeddedErrorJSON(message)?.status ?? error.statusCode ?? null。 +// 原实现只看 message 里的 "" 前缀,statusCode 全被丢掉 → 一律塌成 502。 + +test('error 事件带 statusCode=429 → 回 429 且带 retry_after(不是 502)', async () => { + const s = await setup({ ndjson: [ + '{"type":"text-start"}', + '{"type":"text-delta","text":"partial"}', + '{"type":"error","error":{"message":"providers are currently at capacity","statusCode":429}}', + ] }); + try { + const r = await s.proxy.post('/v1/chat/completions', CHAT, AUTH); + const j = await r.json(); + assert.equal(r.status, 429, 'statusCode 是上游给的,不能抹成 502「服务端错误」'); + assert.equal(j.error.type, 'rate_limit_error'); + assert.equal(j.retry_after, 30, '429 要带退避提示,否则客户端不知道等多久'); + } finally { await s.close(); } +}); + +test('error 事件带 statusCode=503 → 回 503', async () => { + const s = await setup({ ndjson: [ + '{"type":"text-start"}', + '{"type":"error","error":{"message":"service unavailable","statusCode":503}}', + ] }); + try { + const r = await s.proxy.post('/v1/chat/completions', CHAT, AUTH); + assert.equal(r.status, 503); + } finally { await s.close(); } +}); + +test('error 事件没有 statusCode → 回落 502(保持原行为)', async () => { + const s = await setup({ ndjson: [ + '{"type":"text-start"}', + '{"type":"error","error":{"message":"something broke"}}', + ] }); + try { + const r = await s.proxy.post('/v1/chat/completions', CHAT, AUTH); + assert.equal(r.status, 502); + } finally { await s.close(); } +}); + +test('message 里的 "" 前缀优先于 statusCode(对齐 CLI 的取值链)', async () => { + const s = await setup({ ndjson: [ + '{"type":"text-start"}', + '{"type":"error","error":{"message":"<400> bad request","statusCode":503}}', + ] }); + try { + const r = await s.proxy.post('/v1/chat/completions', CHAT, AUTH); + assert.equal(r.status, 400, ' 前缀是最优先的取值来源'); + } finally { await s.close(); } +}); + +// ── ② 无内容事件不应产生 Unknown CC event type 警告 ──────────── +// 上游每个响应都会发一串不携带内容的事件(text-start / text-end / start / +// start-step / reasoning-start / reasoning-end / provider-metadata / +// tool-input-start|delta|end / tool-error)。Responses 非流式那条路径原先 +// 一个静默列表都没有,每个响应刷十来条 warn,真正的错误被淹没。 + +test('标准 NDJSON 序列不产生任何 Unknown CC event type 警告(三协议 × 流式/非流式)', async () => { + const s = await setup(); + try { + const msg = { model: 'm', max_tokens: 50, messages: [{ role: 'user', content: 'hi' }] }; + await (await s.proxy.post('/v1/chat/completions', { ...CHAT, stream: true }, AUTH)).text(); + await (await s.proxy.post('/v1/chat/completions', CHAT, AUTH)).text(); + await (await s.proxy.post('/v1/messages', { ...msg, stream: true }, { 'x-api-key': 'user_test' })).text(); + await (await s.proxy.post('/v1/messages', msg, { 'x-api-key': 'user_test' })).text(); + await (await s.proxy.post('/v1/responses', { model: 'm', stream: true, input: 'hi' }, AUTH)).text(); + await (await s.proxy.post('/v1/responses', { model: 'm', input: 'hi' }, AUTH)).text(); + + const logs = s.proxy.logs(); + assert.ok(!logs.includes('Unknown CC event type'), + '不应出现 Unknown CC event type 警告,实际命中:\n' + + logs.split('\n').filter(l => l.includes('Unknown CC')).join('\n')); + } finally { await s.close(); } +}); From 1479e5e8e0ab9afe864177a60d964291003649af Mon Sep 17 00:00:00 2001 From: xelr233 Date: Tue, 15 Sep 2026 23:09:45 +0800 Subject: [PATCH 3/4] =?UTF-8?q?fix:=20=E6=B5=81=E7=A9=BA=E9=97=B2=E8=B6=85?= =?UTF-8?q?=E6=97=B6=E6=94=B9=E7=94=A8=20res.end()=20=E6=94=B6=E5=B0=BE=20?= =?UTF-8?q?=E2=80=94=E2=80=94=20res.destroy()=20=E4=BC=9A=E4=B8=A2?= =?UTF-8?q?=E7=BC=93=E5=86=B2=E5=B9=B6=E5=8F=91=20RST?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 三处流式超时路径(OpenAI / Anthropic / Responses)原先都是同一个模式: res.write(\`data: \${JSON.stringify({ error: ... })}\\n\\n\`); res.destroy(); res.write() 是异步的,紧接着 destroy() 会把尚未刷出的缓冲丢掉并发 RST。 反向代理侧看到的就是"上游连接被重置": 响应头尚未转发到客户端 → 502 Bad Gateway 已转发 → 客户端 connection error / 截断的流 也就是说:代理本来是"主动截断并告知错误",实际却变成了"把客户端连接搞断"。 下游 SDK 本来能把 rate_limit_error 当可重试错误处理,现在只能吃一个连接层异常。 线上现场(1c2g + OpenResty 反代,资源指标全部健康:NRestarts=0、 MemoryCurrent=215MB、LimitNOFILE=524288、CPU 1.5%、无 OOM): 反代 error.log: sendfile() failed (32: Broken pipe) while sending request to upstream upstream timed out (110) while connecting to upstream 代理 journal: Stream idle timeout {elapsedMs:147005, bytesReceived:833893, lastCcEvent:"reasoning-delta"} (833KB/147s ≈ 5.5KB/s —— 上游本来就慢,30s 空闲阈值确实会被触发; 问题不在"超时",在超时之后怎么收尾) 改法:三处改为 res.end(errEvent) —— 把错误事件正常写进 SSE 流再发 FIN, 下游按可重试错误处理。下游若已僵死(不读也不断),仍由 CLIENT_DRAIN_TIMEOUT_MS 那条路径强制断开,职责不变(那里的 destroy 故意保留)。 测试:新增 1 条,并**验证过有区分度** —— 把 end() 换回 destroy() 时该用例失败, 客户端拿到 "TypeError: terminated"(连接被重置);换回 end() 通过。 用例自带一个"发一半就挂住"的上游,用 fetch().text() 是否成功即可判别两种收尾方式。 --- proxy.mjs | 16 ++++++---- test/connection-lifecycle.test.mjs | 51 ++++++++++++++++++++++++++++-- 2 files changed, 59 insertions(+), 8 deletions(-) diff --git a/proxy.mjs b/proxy.mjs index eef0b0f..60be813 100644 --- a/proxy.mjs +++ b/proxy.mjs @@ -1321,8 +1321,12 @@ async function handleChatCompletions(req, res) { return; } if (!res.writableEnded) { - try { res.write(`data: ${JSON.stringify({ error: { message: timeoutMsg, type: 'rate_limit_error' }, retry_after: 5 })}\n\n`); } catch {} - try { res.destroy(); } catch {} + // 必须 end() 而不是 destroy():res.write 是异步的,紧接着 destroy 会把尚未 + // 刷出的缓冲丢掉并发 RST。反向代理看到上游连接被重置,要么回 502,要么让 + // 客户端看到 connection error —— 这正是"吐字慢 + 间歇性 502"的成因之一。 + // end() 会把错误事件正常送进 SSE 流再发 FIN,客户端 SDK 能按可重试错误处理。 + // 下游若已僵死(不读也不断),由 CLIENT_DRAIN_TIMEOUT_MS 那条路径负责兜底。 + try { res.end(`data: ${JSON.stringify({ error: { message: timeoutMsg, type: 'rate_limit_error' }, retry_after: 5 })}\n\n`); } catch {} } } else { log('error', 'Stream error', { message: e.message }); @@ -2147,8 +2151,8 @@ async function handleMessages(req, res) { const timeoutMsg = consecutiveTimeouts >= TIMEOUT_REDUCE_CONTEXT_THRESHOLD ? 'Response timeout - try reducing context length (summarize earlier messages)' : 'Response timeout - request timed out'; - try { res.write(`event: error\ndata: ${JSON.stringify({ type: 'error', error: { type: 'rate_limit_error', message: timeoutMsg }, retry_after: 5 })}\n\n`); } catch {} - try { res.destroy(); } catch {} + // end() 而不是 destroy():理由见 handleChatCompletions 流式超时分支 + try { res.end(`event: error\ndata: ${JSON.stringify({ type: 'error', error: { type: 'rate_limit_error', message: timeoutMsg }, retry_after: 5 })}\n\n`); } catch {} } } else { log('error', 'Anthropic stream error', { message: e.message }); @@ -2905,8 +2909,8 @@ async function handleResponses(req, res) { : 'Response timeout - request timed out'; if (!started) { sendResponsesError(res, 429, 'rate_limit_error', timeoutMsg, 5); return; } if (!res.writableEnded) { - try { res.write(translator.errorEvent(timeoutMsg)); } catch (e2) {} - try { res.destroy(); } catch (e2) {} + // end() 而不是 destroy():理由见 handleChatCompletions 流式超时分支 + try { res.end(translator.errorEvent(timeoutMsg)); } catch (e2) {} } } else { log('error', 'Stream error', { message: e.message, path: '/v1/responses' }); diff --git a/test/connection-lifecycle.test.mjs b/test/connection-lifecycle.test.mjs index fd3cb15..f5ea771 100644 --- a/test/connection-lifecycle.test.mjs +++ b/test/connection-lifecycle.test.mjs @@ -1,10 +1,12 @@ // 连接生命周期与错误可观测性。 // 这组用例来自一次真实线上故障的排查(间歇性反代 502 / 客户端 connection error): // ① 上游 error 事件自带 statusCode,被丢掉后 429/503 一律塌成 502 -// ② 无内容事件的静默列表不全,Response 非流式那条**根本没有**,日志被刷屏 +// ② 无内容事件的静默列表不全,Responses 非流式那条**根本没有**,日志被刷屏 +// ③ 流空闲超时用 res.destroy() 收尾,把"代理主动截断"变成"把客户端连接搞断" import { test } from 'node:test'; import assert from 'node:assert/strict'; -import { setup } from './helpers.mjs'; +import http from 'node:http'; +import { setup, startProxy, allocPort, closeServer } from './helpers.mjs'; const AUTH = { Authorization: 'Bearer user_test' }; const CHAT = { model: 'm', messages: [{ role: 'user', content: 'hi' }] }; @@ -85,3 +87,48 @@ test('标准 NDJSON 序列不产生任何 Unknown CC event type 警告(三协 logs.split('\n').filter(l => l.includes('Unknown CC')).join('\n')); } finally { await s.close(); } }); + +// ── ③ 流空闲超时必须用 end() 收尾,不能 destroy() ────────────── +// 原实现是 res.write(errEvent) 紧跟 res.destroy():write 是异步的,destroy 会把 +// 尚未刷出的缓冲丢掉并发 RST。反代侧看到的就是 "upstream prematurely closed +// connection" —— 响应头未转发时回 502,已转发时客户端看到 connection error。 +// +// 判别方式:客户端必须能**完整读到**已产生的 delta 与超时错误事件。 +// destroy 会让这条读挂掉(ECONNRESET / terminated),end 则正常收束。 + +/** 一个只在 /alpha/generate 上"发一半就挂住"的上游;其余路由正常应答。 */ +async function startStallingUpstream() { + const port = await allocPort(); + const server = http.createServer((req, res) => { + req.on('data', () => {}); + req.on('end', () => { + if (req.url !== '/alpha/generate') { + // 预请求(fingerprint / lifecycle)必须正常应答,否则代理会卡在初始化上 + res.writeHead(200, { 'Content-Type': 'application/json' }); + res.end('{}'); + return; + } + res.writeHead(200, { 'Content-Type': 'text/event-stream' }); + res.write('{"type":"text-start"}\n'); + res.write('{"type":"text-delta","text":"partial-content"}\n'); + // 之后不再写任何数据 → 触发代理的流空闲超时 + }); + }); + await new Promise(r => server.listen(port, '127.0.0.1', r)); + return { port, close: () => closeServer(server) }; +} + +test('流空闲超时以 end() 收尾:已产生内容 + 错误事件都能完整送达', async () => { + const upstream = await startStallingUpstream(); + const proxy = await startProxy({ upstreamPort: upstream.port, env: { CC_STREAM_IDLE_MS: '300' } }); + try { + const r = await proxy.post('/v1/chat/completions', { ...CHAT, stream: true }, AUTH); + const text = await r.text(); + assert.equal(r.status, 200); + assert.ok(text.includes('partial-content'), + '已发出的内容不能因为收尾方式而丢失(destroy 会丢缓冲 + RST)'); + assert.ok(text.includes('rate_limit_error'), + '超时错误事件必须完整送进流里,下游 SDK 才能按可重试错误处理'); + } finally { await proxy.kill(); await upstream.close(); } +}); + From 25db0c3b79d7df3fd4ca8fdcef03b92c3efb0b0b Mon Sep 17 00:00:00 2001 From: xelr233 Date: Tue, 15 Sep 2026 23:10:23 +0800 Subject: [PATCH 4/4] =?UTF-8?q?fix:=20=E6=98=BE=E5=BC=8F=E8=AE=BE=E7=BD=AE?= =?UTF-8?q?=20server.keepAliveTimeout=EF=BC=8C=E6=B6=88=E9=99=A4=E5=8F=8D?= =?UTF-8?q?=E4=BB=A3=E5=A4=8D=E7=94=A8=E5=B7=B2=E5=85=B3=E9=97=AD=E8=BF=9E?= =?UTF-8?q?=E6=8E=A5=E7=9A=84=20EPIPE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit proxy.mjs 从未设置过 server.keepAliveTimeout,等于把「反代空闲超时 vs 后端空闲超时」 的时序完全交给 Node 默认值(5s)与反代配置的巧合。而这个项目的部署形态是已知的 (README 里就是 nginx / OpenResty 反代),不该靠巧合。 规则:反代的 upstream keepalive_timeout 必须**小于**后端的 keepAliveTimeout。 一旦反过来的,反代会从缓存里取出一条后端已关闭的连接,把请求体写过去 → EPIPE, 而 POST 是非幂等、nginx 默认不重试 → 客户端直接吃 502。线上 error.log 里那 8 条 sendfile() failed (32: Broken pipe) while sending request to upstream 就是这一类。注意是 sendfile() 而非 writev(),说明这些请求体大到被反代缓冲落盘。 Node 默认 5s 与反代常见的 4s 只差 1 秒余量;而两边的计时基准本就不同 (反代从"读完响应放回缓存"起算,后端从"写完响应"起算)。大响应体(线上是 600~830KB 的流式响应)下这点余量随时会被吃掉。 改法:显式 server.keepAliveTimeout = 65s、server.headersTimeout = 66s (CC_KEEPALIVE_TIMEOUT_MS 可覆盖),与 Node 官方"部署在反向代理之后"的建议一致 (keepAliveTimeout > 前端 idle timeout)。启动横幅打出该值并提示反代侧的对应项, 便于部署方对齐。 测试:新增 1 条,锁定"启动横幅必须打出 keepAliveTimeout 且提示反代对应项"。 --- proxy.mjs | 20 ++++++++++++++++++++ test/connection-lifecycle.test.mjs | 17 +++++++++++++++++ 2 files changed, 37 insertions(+) diff --git a/proxy.mjs b/proxy.mjs index 60be813..afafb9a 100644 --- a/proxy.mjs +++ b/proxy.mjs @@ -3132,6 +3132,25 @@ process.on('unhandledRejection', (reason) => { } }); +// ── keep-alive 时序(放在反向代理后面时是必调项) ────────────── +// 反代(nginx/OpenResty)的 upstream keepalive_timeout 必须**小于**这里的值, +// 否则反代会复用一条后端已经关掉的连接:它把请求体写过去,后端早已 FIN, +// 写这一侧就是 EPIPE —— nginx 侧表现为 +// sendfile() failed (32: Broken pipe) while sending request to upstream +// 而这条请求是 POST(非幂等),nginx 默认不会重试 → 客户端直接吃 502。 +// +// Node 默认 keepAliveTimeout=5s。反代若用常见的 4s,余量只有 1 秒;一旦反代的 +// 空闲判定基准与后端差一点(大响应体读完的时刻 vs 后端写完的时刻),就会踩上。 +// 这里显式抬到 65s,让「谁先关」不再取决于一两秒的抖动 —— 与 Node 官方在 +// 反向代理后部署的建议一致(keepAliveTimeout > 前端 idle timeout)。 +// 反代侧仍建议设 keepalive_timeout 60s 以内。 +const KEEPALIVE_TIMEOUT_MS = (() => { + const ms = Number.parseInt(process.env.CC_KEEPALIVE_TIMEOUT_MS ?? '', 10); + return Number.isFinite(ms) && ms > 0 ? ms : 65000; +})(); +server.keepAliveTimeout = KEEPALIVE_TIMEOUT_MS; +server.headersTimeout = KEEPALIVE_TIMEOUT_MS + 1000; // Node 要求 headersTimeout > keepAliveTimeout + server.listen(CFG.port, CFG.host, () => { log('info', 'CC Proxy started', { url: `http://${CFG.host}:${CFG.port}`, @@ -3142,6 +3161,7 @@ server.listen(CFG.port, CFG.host, () => { emptySystemPlaceholder: CFG.emptySystemPlaceholder ? 'on (space placeholder for requests without system prompt, issue #17)' : 'off', logFile: CFG.logFile || '(console only)', clientDrainTimeout: CLIENT_DRAIN_TIMEOUT_MS > 0 ? `${CLIENT_DRAIN_TIMEOUT_MS}ms` : 'disabled', + keepAliveTimeout: `${KEEPALIVE_TIMEOUT_MS}ms (反代侧 keepalive_timeout 必须小于它)`, idleTimeouts: `stream ${STREAM_IDLE_TIMEOUT_MS}ms / nonstream ${NONSTREAM_IDLE_TIMEOUT_MS}ms`, maxInflight: MAX_INFLIGHT > 0 ? `${MAX_INFLIGHT} (global, /health exempt)` : 'unlimited (CC_MAX_INFLIGHT=0)', }); diff --git a/test/connection-lifecycle.test.mjs b/test/connection-lifecycle.test.mjs index f5ea771..11c3c02 100644 --- a/test/connection-lifecycle.test.mjs +++ b/test/connection-lifecycle.test.mjs @@ -132,3 +132,20 @@ test('流空闲超时以 end() 收尾:已产生内容 + 错误事件都能完 } finally { await proxy.kill(); await upstream.close(); } }); + +// ── ④ keep-alive 时序必须在启动横幅里可见 ────────────────────── +// 反代的 upstream keepalive_timeout 必须小于后端的 keepAliveTimeout,否则反代会 +// 复用一条后端已关闭的连接,写请求体时吃 EPIPE(POST 非幂等、nginx 默认不重试 +// → 客户端直接 502)。Node 默认 5s 与反代常见的 4s 只差 1 秒,太薄。 +test('启动横幅打出 keepAliveTimeout,便于与反代配置对齐', async () => { + const s = await setup(); + try { + const logs = s.proxy.logs(); + assert.ok(/keepAliveTimeout[":\s]+65000ms/.test(logs), + '启动横幅必须打出 keepAliveTimeout。实际:\n' + + logs.split('\n').filter(l => l.includes('CC Proxy started')).join('\n')); + assert.ok(logs.includes('反代侧 keepalive_timeout 必须小于它'), + '必须提示反代侧的对应设置,否则这个值没有可操作性'); + } finally { await s.close(); } +}); +