Skip to content
Merged
9 changes: 7 additions & 2 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,9 @@ SERVICE_PORT=3000
# OpenAI Agent quota/WAF failures during stream consumption share this attempt budget.
# Before any content/reasoning is delivered, quota failures may switch to an untried
# healthy account; WAF allows only one switch. Full text is re-externalized on the new
# account, never compacted. Requests with user-uploaded media are not replayed across
# accounts. Exhausted pools stop; token refresh does not clear quota/WAF cooldowns.
# account (a WAF switch reuses the history prefix already uploaded), never compacted. Requests with user-uploaded media are not replayed across
# accounts. Exhausted pools stop; token refresh does not clear quota cooldowns. A WAF
# challenge cools no account: it follows Qwen's load, not the account.
# AGENT_TURN_ALLOW_PROSE_WITH_TOOLS: let prose coexist with tool calls (default false).
# The Anthropic path allows it unconditionally, so enabling it aligns both paths.
# AGENT_TURN_ACCEPT_BARE_FINAL: accept bare prose with no completion wrapper (default
Expand Down Expand Up @@ -129,6 +130,10 @@ AGENT_CONTEXT_FALLBACK_PROMPT_BYTES=86016
# After 3 consecutive WAF challenges on the attachment parse, stop uploading for this many seconds (529 + Retry-After); 0 disables.
AGENT_PARSE_BREAKER_SECONDS=300

# 聊天生成被 WAF 连续拦截(被挤爆啦 / captcha)3 次后,在此秒数内不再发送聊天请求(529/503 带 Retry-After),之后只放行一个探测请求;0 关闭。
# After 3 consecutive WAF challenges on chat generation ("被挤爆啦" / captcha), send nothing for this many seconds (529/503 + Retry-After), then let a single probe through; 0 disables.
CHAT_CHALLENGE_BREAKER_SECONDS=60

# 附件解析速率上限:每个进程在 WINDOW 秒内最多 MAX 次 upload+parse,超出即返回短 Retry-After 的 529,避免触发按 IP 计数的 WAF。MAX=0 关闭。
# Attachment parse rate limit: at most MAX upload+parse per WINDOW seconds per process; beyond that a 529 with a short Retry-After, before the per-IP WAF starts challenging. MAX=0 disables.
AGENT_PARSE_MAX_PER_WINDOW=6
Expand Down
7 changes: 7 additions & 0 deletions src/config/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,13 @@ const config = {
const raw = parseInt(process.env.AGENT_PARSE_BREAKER_SECONDS, 10)
return Number.isFinite(raw) && raw >= 0 ? raw : 300
})(),
// Cortacircuitos del chat challenge (src/utils/upstream-error.js): tras 3 desafios
// seguidos a la generacion no se envia nada durante estos segundos; luego sale una
// sola sonda. 0 lo desactiva.
chatChallengeBreakerSeconds: (() => {
const raw = parseInt(process.env.CHAT_CHALLENGE_BREAKER_SECONDS, 10)
return Number.isFinite(raw) && raw >= 0 ? raw : 60
})(),
// Limitador de ritmo del parse (src/utils/upload.js): como maximo MAX upload+parse por
// ventana de WINDOW segundos por proceso; el resto recibe 529 con Retry-After corto
// ANTES de que el WAF (que cuenta por IP) empiece a desafiar. Medido 2026-09-10:
Expand Down
102 changes: 64 additions & 38 deletions src/controllers/anthropic.js
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ const {
describeUpstreamFailure,
isRateLimitError,
isTransportInterruption,
isWafChallengeError,
noteRateLimitedAccount,
RATE_LIMIT_ANTHROPIC_TYPE,
UpstreamResponseError
Expand Down Expand Up @@ -1188,17 +1189,20 @@ const pingIntervalMs = () => require('../config/index.js').anthropicPingInterval
* 反向代理的空闲计时器,但 SDK 会在读取行时直接丢弃以 `:` 开头的行,客户端因此
* 什么都收不到。ccproxy 网桥当初正是靠改发真正的 ping 事件才消除同样的假死。
*
* 只能在 message_start 之后调用——此时响应头已提交,ping 是合法的流内事件。
* ping 是流内事件,必须跟在 message_start 之后:handleAnthropicStream 延迟提交响应,
* 通过 beforePing 在第一个 ping 之前补发 message_start(也就是延迟提交的上限)。
* @param {object} res - Express 响应
* @param {Function} work - 被包裹的异步任务
* @param {number} [intervalMs] - 发送间隔,缺省取 config.anthropicPingIntervalMs
* @param {Function} [beforePing] - 每次 ping 之前调用
* @returns {Promise<*>} work 的返回值
*/
const runWithAnthropicPing = async (res, work, intervalMs) => {
const runWithAnthropicPing = async (res, work, intervalMs, beforePing) => {
const everyMs = Math.max(1, Number(intervalMs) || pingIntervalMs());
const timer = setInterval(() => {
if (res.writableEnded || res.destroyed) return;
try {
if (typeof beforePing === 'function') beforePing();
writeAnthropicEvent(res, 'ping', { type: 'ping' });
if (typeof res.flush === 'function') res.flush();
} catch (_) {
Expand Down Expand Up @@ -1270,35 +1274,44 @@ const handleAnthropicStream = async (res, ctx, upstream) => {
upstreamOptions = {}
} = ctx;

res.set({
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
});

const createdAt = new Date().toISOString();

// message_start
writeAnthropicEvent(res, 'message_start', {
type: 'message_start',
message: {
id: message_id,
type: 'message',
role: 'assistant',
model,
content: [],
stop_reason: null,
stop_sequence: null,
created_at: createdAt,
metadata: {},
usage: {
input_tokens: 0,
output_tokens: 0,
cache_creation_input_tokens: 0,
cache_read_input_tokens: 0
// Compromiso perezoso: las cabeceras SSE y message_start esperan al primer frame de Qwen que
// pasa assertNoUpstreamFailure (o al primer ping, o al final de una ronda). Si ese primer
// frame es un chat challenge o la cuota agotada, la respuesta sigue libre y el catch de
// handleAnthropicMessages contesta un 529/429 real con Retry-After: la cabecera es lo unico
// que los SDK respetan; un evento de error dentro del stream lo reintentan a los 5 s.
let messageStarted = false;
const ensureMessageStart = () => {
if (messageStarted) return;
messageStarted = true;
res.set({
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive'
});
writeAnthropicEvent(res, 'message_start', {
type: 'message_start',
message: {
id: message_id,
type: 'message',
role: 'assistant',
model,
content: [],
stop_reason: null,
stop_sequence: null,
created_at: createdAt,
metadata: {},
usage: {
input_tokens: 0,
output_tokens: 0,
cache_creation_input_tokens: 0,
cache_read_input_tokens: 0
}
}
}
});
});
};
const withPing = (work) => runWithAnthropicPing(res, work, undefined, ensureMessageStart);

let blockIndex = -1;
let textBlockOpen = false;
Expand Down Expand Up @@ -1754,9 +1767,11 @@ const handleAnthropicStream = async (res, ctx, upstream) => {
startAttempt();

try {
const result = await runWithAnthropicPing(
res,
() => consumeUpstream(currentUpstream, onUpstreamDelta, { shouldStop: () => stopRequested })
const result = await withPing(
() => consumeUpstream(currentUpstream, (json) => {
ensureMessageStart();
return onUpstreamDelta(json);
}, { shouldStop: () => stopRequested })
);
upstreamCompleted = result.completed;
upstreamEventCount = result.eventCount;
Expand All @@ -1773,6 +1788,13 @@ const handleAnthropicStream = async (res, ctx, upstream) => {
action: canFailover ? 'failover' : 'deliver_error'
});
}
// Solo el chat challenge y la cuota salen como status real si aun no se comprometio nada;
// cualquier otro fallo conserva el camino de siempre (evento de error dentro del stream),
// para no convertirlo en un 5xx que los SDK reintentan contra upload/parse.
const rethrow = (error) => {
if (!isWafChallengeError(error) && !isRateLimitError(error)) ensureMessageStart();
throw error;
};
if (!canFailover) {
// Quien servia cuando se cayo, para que el catch del handler marque ESA cuenta si el
// error es de cuota (tras un failover ya no es la del sorteo inicial). Solo el email:
Expand All @@ -1781,7 +1803,7 @@ const handleAnthropicStream = async (res, ctx, upstream) => {
e.failedAccountEmail = ctx.currentAccount?.email || null;
}
logger.error('Anthropic 流式心跳包装失败', 'ANTHROPIC', '', e);
throw e;
rethrow(e);
}

failoverRetries += 1;
Expand All @@ -1791,22 +1813,24 @@ const handleAnthropicStream = async (res, ctx, upstream) => {
recordFailedAccount(e, ctx.currentAccount);
let retryResp = null;
try {
await runWithAnthropicPing(res, async () => {
await withPing(async () => {
retryResp = await sendRequest(requestBody, {
...upstreamOptions,
excludeEmails: failedEmail ? [failedEmail] : []
});
});
} catch (retryError) {
logger.error('Anthropic 流式 failover 重试失败', 'ANTHROPIC', '', retryError);
throw retryError.publicMessage ? retryError : e;
rethrow(retryError.publicMessage ? retryError : e);
}
if (!retryResp?.status || !retryResp.response) throw e;
if (!retryResp?.status || !retryResp.response) rethrow(e);
currentUpstream = retryResp.response;
// Stats y un eventual 429 posterior se atribuyen a quien sirvio de verdad.
if (retryResp.currentAccount) ctx.currentAccount = retryResp.currentAccount;
continue;
}
// Una ronda sin frames JSON tambien compromete: lo que sigue ya escribe bloques.
ensureMessageStart();
attemptsMade += 1;

// 本轮收尾。解析器的尾巴属于这一轮,必须在判定之前放出来。文本通道截断之后例外:
Expand Down Expand Up @@ -1955,7 +1979,7 @@ const handleAnthropicStream = async (res, ctx, upstream) => {

let retryResp = null;
try {
await runWithAnthropicPing(res, async () => {
await withPing(async () => {
retryResp = await sendRequest(appendRetryHint(requestBody, retryHintFor(retryReason)), upstreamOptions);
});
} catch (e) {
Expand Down Expand Up @@ -2802,8 +2826,10 @@ const handleAnthropicMessages = async (req, res) => {
}
// Un prefijo de historial reutilizado pudo ser la causa (file_id que Qwen ya no
// reconoce): se olvida y el reintento del cliente hornea uno nuevo. Un 529 por
// ContextExternalizationError nunca llega aqui con contextPrefixReused.
if (upstreamResp?.contextPrefixReused && error instanceof UpstreamResponseError) {
// ContextExternalizationError nunca llega aqui con contextPrefixReused. Un chat
// challenge tampoco culpa al prefijo: Qwen rechazo antes de leerlo, y olvidarlo haria
// que cada reintento del cliente volviera a subir y parsear el historial entero.
if (upstreamResp?.contextPrefixReused && error instanceof UpstreamResponseError && !isWafChallengeError(error)) {
invalidateContextPrefix(contextPrefixKey);
}
if (!res.headersSent) {
Expand Down
Loading