diff --git a/.env.example b/.env.example index c00d2dbc..3bf9a26a 100644 --- a/.env.example +++ b/.env.example @@ -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 @@ -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 diff --git a/src/config/index.js b/src/config/index.js index 5c3bf4c6..44d3a77a 100644 --- a/src/config/index.js +++ b/src/config/index.js @@ -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: diff --git a/src/controllers/anthropic.js b/src/controllers/anthropic.js index 8ce770ac..2a4100ba 100644 --- a/src/controllers/anthropic.js +++ b/src/controllers/anthropic.js @@ -59,6 +59,7 @@ const { describeUpstreamFailure, isRateLimitError, isTransportInterruption, + isWafChallengeError, noteRateLimitedAccount, RATE_LIMIT_ANTHROPIC_TYPE, UpstreamResponseError @@ -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 (_) { @@ -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; @@ -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; @@ -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: @@ -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; @@ -1791,7 +1813,7 @@ 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] : [] @@ -1799,14 +1821,16 @@ const handleAnthropicStream = async (res, ctx, upstream) => { }); } 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; // 本轮收尾。解析器的尾巴属于这一轮,必须在判定之前放出来。文本通道截断之后例外: @@ -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) { @@ -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) { diff --git a/src/controllers/chat.image.video.js b/src/controllers/chat.image.video.js index 2ce6aaec..810974b4 100644 --- a/src/controllers/chat.image.video.js +++ b/src/controllers/chat.image.video.js @@ -10,7 +10,7 @@ const { getDefaultModelByChatType } = require('../models/models-map.js') const { getSsxmodForAccount } = require('../utils/ssxmod-manager') const { applyProxyToAxiosConfig, getChatBaseUrl } = require('../utils/proxy-helper'); const { buildRequestHeaders } = require('../utils/header-profile') -const { detectWafChallenge } = require('../utils/upstream-error.js') +const { assertChatChallengeBreakerClosed, chatChallengeFrom, noteChatAnswer, releaseChatProbe } = require('../utils/upstream-error.js') const DATA_URI_REGEX = /^data:(.+);base64,(.*)$/i const HTTP_URL_REGEX = /^https?:\/\//i @@ -66,6 +66,17 @@ const buildAxiosErrorLog = (error) => ({ data: formatPayloadForLog(error?.response?.data) }) +/** + * 图片/视频也走 /api/v2/chat/completions,同样会被 chat challenge 拦截: + * 转成可重试的 503(带 retry_after),与聊天接口一致。 + */ +const imageChallengeError = (challenge) => ({ + error: challenge.publicMessage, + code: challenge.code, + status: 503, + retry_after: challenge.retryAfter +}) + const parseUpstreamImageError = (data) => { try { const rawPayload = formatPayloadForLog(data) @@ -79,20 +90,10 @@ const parseUpstreamImageError = (data) => { payload = JSON.parse(payload) } - // WAF/captcha llega como HTTP 200 con una forma propia (`ret`), NO como - // `success:false`. Sin esta rama el paquete se colaba entero: la generación - // "tenía éxito" y se devolvía la imagen de relleno de Qwen como resultado. - const wafChallenge = detectWafChallenge(payload) - if (wafChallenge) { - logger.error('图片/视频上游触发 WAF/captcha,需人工验证', 'CHAT', '', { - parsed_error: wafChallenge.details, - raw_response_body: rawPayload - }) - return { - error: wafChallenge.publicMessage, - code: wafChallenge.code, - status: 502 - } + const challenge = chatChallengeFrom(payload) + if (challenge) { + logger.error('图片/视频请求被上游 WAF 拦截 (chat challenge)', 'CHAT', '', { raw_response_body: rawPayload }) + return imageChallengeError(challenge) } // 只有明确 success=false 且带错误码时,才按上游错误包处理,避免误伤正常业务响应 @@ -164,16 +165,12 @@ const parseUpstreamErrorFromRawText = (text) => { // La página de captcha del WAF llega como HTML y no se deja ver por ninguna de las // ramas de abajo (no hay `data:`, no es JSON). Se comprueba primero. - const htmlChallenge = detectWafChallenge(trimmed) + const htmlChallenge = chatChallengeFrom(trimmed) if (htmlChallenge) { - logger.error('图片/视频上游返回 WAF/captcha HTML 页面,需人工验证', 'CHAT', '', { + logger.error('图片/视频上游返回 WAF/captcha HTML 页面 (chat challenge)', 'CHAT', '', { raw_preview: trimmed.slice(0, 200) }) - return { - error: htmlChallenge.publicMessage, - code: htmlChallenge.code, - status: 502 - } + return imageChallengeError(htmlChallenge) } const fromWholeText = parseUpstreamImageErrorFromText(trimmed) @@ -191,6 +188,15 @@ const parseUpstreamErrorFromRawText = (text) => { return null } +/** + * t2v pide responseType 'json': un cuerpo que no es JSON (la pagina de captcha, o un + * `data: {...}`) llega como string, y parseUpstreamImageError lo descartaba al fallar el + * JSON.parse; el "video" terminaba siendo la primera URL de la pagina de captcha. + */ +const parseUpstreamBody = (data) => ( + typeof data === 'string' ? parseUpstreamErrorFromRawText(data) : parseUpstreamImageError(data) +) + /** * 收集对象中的所有值 * @param {*} payload - 任意负载 @@ -598,6 +604,7 @@ const isRetryableUpstreamError = (upstreamError) => { */ const sendUpstreamError = (res, upstreamError) => { const { status, ...payload } = upstreamError + if (Number(payload.retry_after) > 0) res.set({ 'Retry-After': String(payload.retry_after) }) return res.status(status || 500).json(payload) } @@ -682,6 +689,7 @@ const normalizeOpenAIImageVideoSize = (size) => { const sendOpenAIErrorResponse = (res, error) => { const status = error?.status || 500 const message = error?.error || error?.message || 'Service error, please try again later' + if (Number(error?.retry_after) > 0) res.set({ 'Retry-After': String(error.retry_after) }) return res.status(status).json({ error: { @@ -877,7 +885,7 @@ const readVideoUpstreamResult = async (responseStream) => { const videoTaskCandidates = extractVideoTaskIdentifiersFromPayload(responseStream) const responseIDs = extractResponseIDsFromPayload(responseStream) return { - upstreamError: parseUpstreamImageError(responseStream), + upstreamError: parseUpstreamBody(responseStream), contentUrl: extractResourceUrlFromPayload(responseStream), videoTaskID: videoTaskCandidates[0] || null, videoTaskCandidates, @@ -976,7 +984,7 @@ const readVideoUpstreamResult = async (responseStream) => { const readImageUpstreamResult = async (responseStream) => { if (!responseStream || typeof responseStream.on !== 'function') { return { - upstreamError: parseUpstreamImageError(responseStream), + upstreamError: parseUpstreamBody(responseStream), contentUrl: extractResourceUrlFromPayload(responseStream), responseIDs: extractResponseIDsFromPayload(responseStream), rawPreview: typeof responseStream === 'string' ? responseStream.slice(0, 400) : '' @@ -1275,6 +1283,7 @@ const generateImageVideoResult = async (payload) => { // 一次取出账户对象,确保 token 与 proxy 走同一个账号 const account = accountManager.getAccount() const token = account ? account.token : null + let probe = false try { const reqBody = { @@ -1296,14 +1305,7 @@ const generateImageVideoResult = async (payload) => { ] } - const chatID = await generateChatID(token, model, account, chat_type) - - if (!chatID) { - throw new Error('生成 chat_id 失败') - } - - reqBody.chat_id = chatID - + // Antes del breaker: una peticion invalida no debe quedarse con la unica sonda. const userPrompt = messages?.[messages.length - 1]?.content if (!userPrompt) { throw { @@ -1312,6 +1314,20 @@ const generateImageVideoResult = async (payload) => { } } + try { + probe = assertChatChallengeBreakerClosed() + } catch (challenge) { + throw imageChallengeError(challenge) + } + + const chatID = await generateChatID(token, model, account, chat_type) + + if (!chatID) { + throw new Error('生成 chat_id 失败') + } + + reqBody.chat_id = chatID + const messagesHistory = messages.filter(item => item.role === 'user' || item.role === 'assistant') const selectedImageList = [] @@ -1434,7 +1450,9 @@ const generateImageVideoResult = async (payload) => { try { responseData = await axios.post(`${chatBaseUrl}/api/v2/chat/completions?chat_id=${chatID}`, reqBody, requestConfig) - const inlineUpstreamError = parseUpstreamImageError(responseData.data) + const inlineUpstreamError = parseUpstreamBody(responseData.data) + // Ya contado como strike: no dejar que el resolver lo vuelva a parsear y contar. + if (inlineUpstreamError?.code === 'upstream_waf_challenge') throw inlineUpstreamError if (attempt < maxUpstreamAttempts && isRetryableUpstreamError(inlineUpstreamError)) { logger.warn(`图片/视频请求上游返回业务错误包,准备第 ${attempt + 1} 次重试,请求ID: ${inlineUpstreamError.request_id || '未知'}`, 'CHAT') await sleep(800) @@ -1444,7 +1462,9 @@ const generateImageVideoResult = async (payload) => { break } catch (error) { logger.error('图片/视频请求失败', 'CHAT', '', buildAxiosErrorLog(error)) - const upstreamError = parseUpstreamImageError(error.response?.data) + const upstreamError = parseUpstreamBody(error.response?.data) + // Un desafio con status no-2xx ya sumo su strike aqui: el catch externo no lo re-parsea. + if (upstreamError?.code === 'upstream_waf_challenge') throw upstreamError if (attempt < maxUpstreamAttempts && isRetryableUpstreamError(upstreamError)) { logger.warn(`图片/视频请求上游返回瞬时内部错误,准备第 ${attempt + 1} 次重试`, 'CHAT') await sleep(800) @@ -1457,6 +1477,7 @@ const generateImageVideoResult = async (payload) => { if (newChatType === 't2i' || newChatType === 'image_edit') { const contentUrl = await resolveImageResultContentUrl(responseData.data, chatID, token) + noteChatAnswer() return { model, chatType: newChatType, @@ -1466,6 +1487,9 @@ const generateImageVideoResult = async (payload) => { } if (newChatType === 't2v') { + // El cuerpo t2v ya se reviso entero arriba sin desafio: Qwen acepto la tarea. No esperar + // al sondeo del video (minutos) para cerrar una media apertura. + noteChatAnswer() const contentUrl = await resolveVideoResultContentUrl(responseData.data, token, chatID) return { model, @@ -1478,6 +1502,7 @@ const generateImageVideoResult = async (payload) => { throw new Error('不支持的图片/视频类型') } catch (error) { logger.error('图片/视频主流程异常', 'CHAT', '', buildAxiosErrorLog(error)) + if (probe && error?.code !== 'upstream_waf_challenge') releaseChatProbe() if (error?.error) { throw error @@ -1528,6 +1553,13 @@ const handleImageVideoCompletion = async (req, res) => { logger.error('图片视频资源处理错误', 'CHAT', '', error) + // chat challenge 且尚未提交:返回真正的 503 + Retry-After(t2v 预设了 text/event-stream,需改回 JSON), + // 而不是把错误文本当成 200 的流式正文。 + if (downstreamStream && !res.headersSent && error?.code === 'upstream_waf_challenge') { + res.set({ 'Content-Type': 'application/json' }) + return sendUpstreamError(res, error) + } + if (downstreamStream) { return returnResponse(res, req.body.model, error?.error || error?.message || 'Service error, please try again later', true) } @@ -1789,5 +1821,7 @@ module.exports = { handleImageVideoCompletion, handleOpenAIImagesGeneration, handleOpenAIImagesEdit, - handleOpenAIVideoGeneration + handleOpenAIVideoGeneration, + // 暴露内部辅助以便测试 + parseUpstreamImageError } diff --git a/src/controllers/chat.js b/src/controllers/chat.js index ff475a87..94e4fe85 100644 --- a/src/controllers/chat.js +++ b/src/controllers/chat.js @@ -20,6 +20,8 @@ const { createUpstreamDeltaNormalizer, createClientToolNamePredicate } = require const { assertNoUpstreamFailure, describeUpstreamFailure, + isRateLimitError, + isWafChallengeError, noteRateLimitedAccount, RATE_LIMIT_OPENAI_TYPE } = require('../utils/upstream-error.js') @@ -258,13 +260,14 @@ const runWithProcessingHeartbeat = async (res, work, intervalMs = 15000) => { } } -const runWithSSEHeartbeat = async (res, work, intervalMs = 15000) => { +const runWithSSEHeartbeat = async (res, work, intervalMs = 15000, beforeBeat = null) => { const heartbeatMs = Math.max(1, Number(intervalMs) || 15000) const heartbeat = setInterval(() => { if (res.writableEnded || res.destroyed) return try { - // SSE 已经提交 200 响应后,用注释帧保活。注释不会进入 OpenAI delta, - // 但能阻止反代在长 thinking 或纠正 attempt 期间把连接判为空闲。 + // 用注释帧保活。注释不会进入 OpenAI delta,但能阻止反代在长 thinking 或纠正 + // attempt 期间把连接判为空闲。beforeBeat 先提交 SSE 首帧(延迟提交的上限)。 + if (typeof beforeBeat === 'function') beforeBeat() res.write(': qwen2api-agent-keepalive\n\n') if (typeof res.flush === 'function') res.flush() } catch (_) { @@ -391,25 +394,31 @@ const handleOpenAIAgentStream = async ( setResponseHeaders(res, true) const messageId = generateUUID() const created = Math.round(Date.now() / 1000) - let firstDelta = true - const writeDelta = (delta) => { - if (!delta || Object.keys(delta).length === 0) return - const normalizedDelta = firstDelta ? { role: 'assistant', ...delta } : delta - firstDelta = false + const writeChunk = (delta) => { res.write(`data: ${JSON.stringify({ id: `chatcmpl-${messageId}`, object: 'chat.completion.chunk', created, - choices: [{ index: 0, delta: normalizedDelta, finish_reason: null }] + choices: [{ index: 0, delta, finish_reason: null }] })}\n\n`) if (typeof res.flush === 'function') res.flush() } - - // 立即提交标准 SSE 首帧,不能等整个上游 attempt 收完后才让客户端看到响应。 - // 裸正文/工具调用仍由下方门禁缓冲;安全思考与已确认进入 final/blocked - // 包装体的正式正文会按上游节奏增量输出。 - if (typeof res.flushHeaders === 'function') res.flushHeaders() - writeDelta({ role: 'assistant' }) + // SSE 首帧(role 单独一块)在上游第一帧通过校验时提交,不等整个 attempt 收完; + // 裸正文/工具调用仍由下方门禁缓冲,安全思考与已确认的正文按上游节奏增量输出。 + // Compromiso perezoso: si el primer frame de Qwen es un chat challenge o la cuota, la + // respuesta sigue libre y sale un 503/429 real con Retry-After en vez de un frame de error. + let committed = false + const commitStream = () => { + if (committed) return + committed = true + if (typeof res.flushHeaders === 'function') res.flushHeaders() + writeChunk({ role: 'assistant' }) + } + const writeDelta = (delta) => { + if (!delta || Object.keys(delta).length === 0) return + commitStream() + writeChunk(delta) + } const liveReasoningByAttempt = new Map() const onReasoningDelta = enableThinking && !config.legacyReasoningInContent @@ -443,20 +452,26 @@ const handleOpenAIAgentStream = async ( sendChatRequest: options.sendChatRequest || sendChatRequest, on_reasoning_delta: onReasoningDelta, on_content_delta: onContentDelta, + on_upstream_frame: commitStream, isClientDisconnected: () => res.destroyed || res.writableEnded }), - options.agent_processing_heartbeat_ms + options.agent_processing_heartbeat_ms, + commitStream ) } catch (error) { logger.error('OpenAI Agent 回合处理失败', 'AGENT', '', error) if (!error.accountFailureRecorded) { noteRateLimitedAccount(error, error.failedAccountEmail ? { email: error.failedAccountEmail } : options.currentAccount) } + // Solo el chat challenge y la cuota salen como status real sin compromiso previo; el + // resto conserva el frame de error de siempre (no un 5xx que los SDK reintentan). + if (!isWafChallengeError(error) && !isRateLimitError(error)) commitStream() writeOpenAIHttpError(res, upstreamErrorShape( error, '上游 Agent 回合处理失败', 'upstream_stream_error' )) return } + commitStream() if (!runtime.ok) { writeOpenAIHttpError(res, runtime.error) return @@ -1079,34 +1094,25 @@ const handleStreamResponse = async (res, response, enable_thinking, enable_web_s res.end() } catch (error) { logger.error('聊天处理错误', 'CHAT', '', error) - // Cuota agotada -> 429 `insufficient_quota`; adjunto caido -> 503 (529 es de - // Anthropic); cualquier otro fallo conserva su etiqueta de siempre. Deteccion - // unica en utils/upstream-error.js. + // Cuota agotada -> 429 `insufficient_quota`; adjunto caido o chat challenge -> 503 + // `upstream_unavailable` (529 es de Anthropic); cualquier otro fallo conserva su + // etiqueta de siempre. Deteccion unica en utils/upstream-error.js; la forma, en + // upstreamErrorShape (la misma que el modo agente). const failure = describeUpstreamFailure(error, 502, 503) + const shape = upstreamErrorShape(error, '上游流式传输失败', 'upstream_stream_error') noteRateLimitedAccount(error, options.currentAccount) if (res.headersSent) { if (!res.writableEnded) { writeOpenAIStreamError( res, - error.publicMessage || '上游流式传输失败', - failure.rateLimited - ? RATE_LIMIT_OPENAI_TYPE - : (error.publicMessage ? error.code : 'upstream_stream_error'), - failure.rateLimited ? RATE_LIMIT_OPENAI_TYPE : 'upstream_stream_error', + shape.message, + failure.rateLimited || failure.overloaded || error.publicMessage ? shape.code : 'upstream_stream_error', + shape.type || 'upstream_stream_error', failure.retryAfter ) } } else { - if (failure.retryAfter !== null) res.set({ 'Retry-After': String(failure.retryAfter) }) - res.status(failure.status).json({ - error: { - message: error.publicMessage || '上游流式传输失败', - type: failure.rateLimited ? RATE_LIMIT_OPENAI_TYPE : 'upstream_stream_error', - code: failure.rateLimited - ? RATE_LIMIT_OPENAI_TYPE - : (error.code || 'upstream_stream_error') - } - }) + writeOpenAIHttpError(res, { type: 'upstream_stream_error', ...shape }) } } } @@ -1434,18 +1440,8 @@ const handleNonStreamResponse = async (res, response, enable_thinking, enable_we res.json(bodyTemplate) } catch (error) { logger.error('非流式聊天处理错误', 'CHAT', '', error) - const failure = describeUpstreamFailure(error, 502, 503) noteRateLimitedAccount(error, options.currentAccount) - if (!res.headersSent) { - if (failure.retryAfter !== null) res.set({ 'Retry-After': String(failure.retryAfter) }) - res.status(failure.status).json({ - error: { - message: error.publicMessage || '上游响应处理失败', - type: failure.rateLimited ? RATE_LIMIT_OPENAI_TYPE : 'upstream_error', - code: failure.rateLimited ? RATE_LIMIT_OPENAI_TYPE : (error.code || 'upstream_error') - } - }) - } + if (!res.headersSent) writeOpenAIHttpError(res, upstreamErrorShape(error, '上游响应处理失败')) } } diff --git a/src/utils/account-rotator.js b/src/utils/account-rotator.js index cb9a8f7e..c3d9d898 100644 --- a/src/utils/account-rotator.js +++ b/src/utils/account-rotator.js @@ -14,7 +14,6 @@ class AccountRotator { this.lastErrorCode = new Map() // 最近一次错误码(HTTP status 或 transport err.code) this.cooldownStartedAt = new Map() // 进入 cooldown 的起始时间戳(failureCounts 达阈值时刻) this.quotaCooldownUntil = new Map() // 日额度耗尽的账户 -> 解禁时间戳(见 recordQuotaExhausted) - this.challengeCooldownUntil = new Map() this.maxFailures = 3 // 最大失败次数 this.cooldownPeriod = 5 * 60 * 1000 // 5分钟冷却期 // 额度耗尽的默认静默期。上游给了 `data.num`(小时)时用那个,这是没给时的回退。 @@ -172,13 +171,6 @@ class AccountRotator { ) } - recordChallenge(email) { - if (!email) return - // Keep this separate from transport failures and quota cooldowns. - this.challengeCooldownUntil.set(email, Date.now() + this.cooldownPeriod) - this.recordError(email, 'upstream_waf_challenge') - } - /** * 重置账户失败计数(清除 cooldown) * 注意:不清理 lastErrorAt/lastErrorCode——它们由 endpoint 的 15 分钟窗口管理 @@ -212,10 +204,7 @@ class AccountRotator { available: this._isAccountAvailable(account), lastErrorAt: this.lastErrorAt.get(email) || null, lastErrorCode: this.lastErrorCode.get(email) || null, - cooldownEndsAt: Math.max( - cooldownStart ? cooldownStart + this.cooldownPeriod : 0, - this.challengeCooldownUntil.get(email) || 0 - ) || null, + cooldownEndsAt: cooldownStart ? cooldownStart + this.cooldownPeriod : null, quotaCooldownEndsAt: this.quotaCooldownUntil.get(email) || null } }) @@ -248,13 +237,6 @@ class AccountRotator { return false } - // A token refresh does not resolve an upstream verification challenge. - const challengeUntil = this.challengeCooldownUntil.get(account.email) - if (challengeUntil) { - if (Date.now() < challengeUntil) return false - this.challengeCooldownUntil.delete(account.email) - } - // 额度流放优先于一切:这个账户对上游来说今天已经没有配额,再选它就是白烧一轮。 const quotaUntil = this.quotaCooldownUntil.get(account.email) if (quotaUntil) { @@ -321,8 +303,7 @@ class AccountRotator { this.lastErrorAt, this.lastErrorCode, this.cooldownStartedAt, - this.quotaCooldownUntil, - this.challengeCooldownUntil + this.quotaCooldownUntil ] for (const map of maps) { for (const email of map.keys()) { @@ -344,7 +325,6 @@ class AccountRotator { this.lastErrorCode.clear() this.cooldownStartedAt.clear() this.quotaCooldownUntil.clear() - this.challengeCooldownUntil.clear() } } diff --git a/src/utils/account.js b/src/utils/account.js index 254cfd17..fcdd317a 100644 --- a/src/utils/account.js +++ b/src/utils/account.js @@ -762,10 +762,6 @@ class Account { this.accountRotator.recordQuotaExhausted(email, retryAfterSeconds) } - recordAccountChallenge(email) { - this.accountRotator.recordChallenge(email) - } - /** * 累计 daily stats(per-account) * 调用方:chat.js / anthropic.js / cli.chat.js 在成功消费完上游 usage 后 diff --git a/src/utils/agent-account-failover.js b/src/utils/agent-account-failover.js index 1352bcbc..a68b5a86 100644 --- a/src/utils/agent-account-failover.js +++ b/src/utils/agent-account-failover.js @@ -1,13 +1,11 @@ -const { isRateLimitError, isWafChallengeError, noteRateLimitedAccount } = require('./upstream-error') +const { isRateLimitError, noteRateLimitedAccount } = require('./upstream-error') function recordFailedAccount(error, account) { // Errors are logged; never attach the account object containing its token/password. error.failedAccountEmail = account?.email || null + // A chat challenge follows Qwen's load, not the account: nothing to record against it. if (isRateLimitError(error)) { error.accountFailureRecorded = noteRateLimitedAccount(error, account) - } else if (isWafChallengeError(error) && account?.email) { - require('./account').recordAccountChallenge(account.email) - error.accountFailureRecorded = true } } diff --git a/src/utils/openai-agent-runtime.js b/src/utils/openai-agent-runtime.js index c3ea6a75..62709e4e 100644 --- a/src/utils/openai-agent-runtime.js +++ b/src/utils/openai-agent-runtime.js @@ -332,6 +332,8 @@ const collectOpenAIAgentAttempt = async (upstreamResponse, options = {}) => { const decoded = isJson(frame.data) ? JSON.parse(frame.data) : null if (decoded === null) return assertNoUpstreamFailure(decoded) + // The frame passed: the controller may now commit its stream (see chat.js#commitStream). + if (typeof options.on_upstream_frame === 'function') options.on_upstream_frame() const created = normalizeCreatedMetadata(decoded) if (created) { @@ -795,7 +797,7 @@ const runOpenAIAgentTurn = async (initialResponse, options = {}) => { const replayBody = (quotaFailure || challengeFailure) && !deliveredOutput ? createAccountReplayBody(options.requestBody) : null - if (!replayBody || !currentAccount?.email || !error.accountFailureRecorded || + if (!replayBody || !currentAccount?.email || !(error.accountFailureRecorded || challengeFailure) || attemptNumber >= maxAttempts || typeof requestSender !== 'function' || options.isClientDisconnected?.() || (challengeFailure && challengeFailovers >= 1)) { throw error @@ -807,9 +809,16 @@ const runOpenAIAgentTurn = async (initialResponse, options = {}) => { if (challengeFailure) challengeFailovers += 1 logger.warn(`Agent attempt ${attemptNumber}/${maxAttempts}: ${error.code}; retrying with a different healthy account`, 'AGENT') currentAccount = replacementAccount - // Re-externalize the original complete prompt. Do not reuse another - // account's conversation IDs, shortened prompt, or cached context attachment. - currentUpstreamOptions = { ...currentUpstreamOptions, contextPrefixKey: null, allowContextCompaction: false } + // Re-externalize the original complete prompt; never reuse another account's + // conversation IDs or a shortened prompt. A quota switch also re-uploads the history. + // A chat challenge keeps the history prefix: Qwen refused before reading it, the + // prefix is not tied to an account (normal rotation reuses it across accounts too), + // and re-uploading it on every challenge is what exhausts the parse budget. + currentUpstreamOptions = { + ...currentUpstreamOptions, + contextPrefixKey: challengeFailure ? currentUpstreamOptions.contextPrefixKey : null, + allowContextCompaction: false + } const retryResponse = await sendBoundRequest(replayBody, { ...currentUpstreamOptions, chatId: null, @@ -817,6 +826,9 @@ const runOpenAIAgentTurn = async (initialResponse, options = {}) => { agentRetry: true }) if (!retryResponse?.status || !retryResponse.response) { + // A switch that could not even start keeps its cause: a challenge is still a + // retryable 503 and quota a 429, both with Retry-After — not an opaque 502. + if (challengeFailure || quotaFailure) throw error return { ok: false, error: { status: 502, message: retryResponse?.message || 'Account failover request failed', code: 'upstream_retry_failed' }, diff --git a/src/utils/request.js b/src/utils/request.js index ccbec424..9b2f926f 100644 --- a/src/utils/request.js +++ b/src/utils/request.js @@ -1,4 +1,5 @@ const axios = require('axios') +const { Readable } = require('node:stream') const accountManager = require('./account.js') const config = require('../config/index.js') const { logger } = require('./logger') @@ -7,7 +8,7 @@ const { applyProxyToAxiosConfig, getChatBaseUrl } = require('./proxy-helper'); const { generateUUID, jitter } = require('./tools.js') const { uploadAgentContextFile, buildChatFileDescriptor } = require('./upload.js') const { buildRequestHeaders } = require('./header-profile') -const { ContextExternalizationError, isTransportInterruption } = require('./upstream-error.js') +const { ContextExternalizationError, isTransportInterruption, assertChatChallengeBreakerClosed, chatChallengeFrom, isWafChallengeError, releaseChatProbe } = require('./upstream-error.js') const { contextPrefixCache, prefixMatches, canonicalHistoryHash } = require('./context-prefix-cache.js') const { TOOL_CALL_OPEN, LEDGER_HEADER, LEDGER_CAPTION, truncateToolHistoryLedger, stripRetainedThinking @@ -27,6 +28,36 @@ const isRetryableNetworkError = (error) => { const delay = (ms) => new Promise(resolve => setTimeout(resolve, ms)) +// La pagina de captcha pesa ~16 KB; el tope solo evita cargar en memoria un cuerpo inesperado. +const HTML_BODY_MAX_BYTES = 256 * 1024 + +/** + * Qwen tambien manda el chat challenge como 200 text/html (la pagina de captcha del WAF): + * sin frames `data:` nadie lo veia y el cliente recibia "respuesta vacia" tras 2-3 reintentos. + * Lanza el chat challenge; cualquier otro HTML se devuelve intacto para no cambiar su camino. + */ +const screenHtmlChallenge = async (response) => { + if (!/text\/html/i.test(String(response.headers?.['content-type'] || ''))) return response.data + const chunks = [] + let size = 0 + let readError = null + try { + for await (const chunk of response.data) { + chunks.push(Buffer.from(chunk)) + size += chunk.length + if (size >= HTML_BODY_MAX_BYTES) break + } + } catch (error) { + // Corte a media pagina: si las marcas del captcha ya llegaron, sigue siendo el desafio. + readError = error + } + const body = Buffer.concat(chunks).toString('utf8') + const challenge = chatChallengeFrom(body) + if (challenge) throw challenge + if (readError) throw readError + return Readable.from([body]) +} + const HISTORY_MARKER = '# Conversation history (JSONL)' const CURRENT_MESSAGE_MARKER = '# Current message' // Techo del ledger dentro del presupuesto inline. Ver buildBudgetedAgentPrompt. @@ -642,6 +673,19 @@ const externalizeOversizedAgentContext = async ( * @returns {Promise} 响应结果 */ const sendChatRequest = async (body, options = {}) => { + // Qwen esta rechazando la generacion: ni chat nuevo ni upload, solo 529/503 al cliente. + const probe = assertChatChallengeBreakerClosed() + try { + const result = await postChatRequest(body, options) + if (probe && !result.status) releaseChatProbe() + return result + } catch (error) { + if (probe && !isWafChallengeError(error)) releaseChatProbe() + throw error + } +} + +const postChatRequest = async (body, options = {}) => { // 获取可用的账户(包含 proxy 等完整字段) // excludeEmails:本次 HTTP 请求里已经烧掉的账户(流中途 failover)——轮换器跳过它们, // 即使它们对其他请求仍然可用。 @@ -766,7 +810,7 @@ const sendChatRequest = async (body, options = {}) => { contextPrefixReused: contextResult.reusedPrefix === true, contextSerializedBytes: contextResult.serializedBytes, status: true, - response: response.data + response: await screenHtmlChallenge(response) } } // 非 200 但是没抛——退出循环, 走下面错误分类 @@ -774,6 +818,8 @@ const sendChatRequest = async (body, options = {}) => { lastError.response = { status: response.status } break } catch (error) { + // El cliente recibe 529/503 + Retry-After; el prefijo reutilizado no tuvo la culpa. + if (isWafChallengeError(error)) throw error lastError = error if (isRetryableNetworkError(error) && attempt < totalAttempts) { logger.warn( diff --git a/src/utils/upstream-error.js b/src/utils/upstream-error.js index 2414536f..ee253b73 100644 --- a/src/utils/upstream-error.js +++ b/src/utils/upstream-error.js @@ -73,7 +73,7 @@ const findWafChallengeSignals = (payload) => { payload.data?.code, payload.data?.url, payload.error?.code - ].filter(Boolean).map(String).filter(item => WAF_SIGNAL_RE.test(item)); + ].filter(Boolean).map(String).filter(item => item.toLowerCase() === WAF_CHALLENGE_CODE || WAF_SIGNAL_RE.test(item)); }; /** * Detecta el challenge y devuelve el error canónico, o null si el payload no lo es. @@ -199,26 +199,31 @@ const rateLimitRetryAfterSeconds = (error) => { /** * MATRIZ DE ALCANZABILIDAD — que recibe el cliente de verdad, por camino y por fase. * - * El 429 solo es alcanzable mientras las cabeceras siguen libres. En streaming los dos - * controladores comprometen el 200 ANTES de leer un byte del upstream, asi que un - * paquete de cuota —que llega como PRIMER frame— nunca puede cambiar el status: + * Un status (y la cabecera Retry-After, la unica que los SDK respetan) solo es alcanzable + * mientras la respuesta sigue libre. Los caminos en streaming comprometen de forma perezosa: + * nada sale hasta que el PRIMER frame de Qwen pasa assertNoUpstreamFailure, asi que un + * chat challenge o un paquete de cuota —que llegan como primer frame— si cambian el status: * - * camino fase status senal para el cliente - * /v1/messages stream:false libre 429 body.error.type - * /v1/messages stream:true comprometida 200 evento error.type (+retry_after) - * /v1/chat/... llano stream:false libre 429 body.error.type - * /v1/chat/... llano stream:true libre 1er byte 429 body.error.type - * /v1/chat/... agente stream:true comprometida 200 frame error.type (+retry_after) + * camino se compromete en fallo en el 1er frame + * /v1/messages stream:false al final 429/529 + Retry-After + * /v1/messages stream:true 1er frame valido o 1er ping 429/529 + Retry-After + * (anthropic.js#handleAnthropicStream, ensureMessageStart) + * /v1/chat/... llano stream:false al final 429/503 + Retry-After + * /v1/chat/... llano stream:true 1er byte escrito 429/503 + Retry-After + * /v1/chat/... agente stream:true 1er frame valido o 1er latido 429/503 + Retry-After + * (chat.js#handleOpenAIAgentStream, commitStream) + * imagen stream:true al entregar el resultado 503 + Retry-After + * video (t2v) stream:true 1er keep-alive (15 s; las cabeceras 503 + Retry-After + * SSE se fijan antes pero no se envian) * - * La fila que importa es la segunda: Claude Code habla /v1/messages con stream:true, y - * las 149 negativas de cuota de los logs del usuario salen todas de ahi. Decir que este - * cambio "mapea la cuota a 429 en los dos caminos" es falso justo para el modo que el - * usuario ejecuta; lo que hace es que la negativa sea RECONOCIBLE en los dos caminos y - * en las dos fases. anthropic.js:1203/1211 y chat.js:407 son las lineas que comprometen - * la respuesta, y adelantarlas es deliberado: sin cabeceras enviadas no se pueden mandar - * `ping` dentro del protocolo (anthropic.js:1156-1160), que es como se elimino el falso - * "stream muerto" del puente. Por eso la espera viaja DENTRO del evento/frame: es el - * unico canal que queda cuando la cabecera Retry-After ya no se puede poner. + * Solo el chat challenge y la cuota se quedan sin comprometer; cualquier otro fallo antes + * del primer frame compromete y sale como antes, dentro del stream (convertirlo en 5xx + * haria que los SDK lo reintenten contra upload/parse). Tras el compromiso —p. ej. cuota a + * mitad de respuesta, las 149 negativas de los logs del usuario— la espera viaja DENTRO del + * evento/frame (`retry_after`): es el unico canal que queda. El compromiso perezoso no + * reabre el falso "stream muerto" del puente: los `ping`/latidos lo comprometen como mucho + * un intervalo despues, y el silencio largo que resolvian era el del reintento de correccion, + * que ocurre ya comprometido. */ /** @@ -226,10 +231,18 @@ const rateLimitRetryAfterSeconds = (error) => { * repetir la deteccion; el `type` de cable lo pone cada uno con su constante de arriba. * @param {unknown} error - Error capturado * @param {number} [fallbackStatus] - Status cuando NO es cuota (500 Anthropic / 502 OpenAI) - * @param {number} [overloadedStatus] - Status del adjunto de contexto caido (529 Anthropic / 503 OpenAI) + * @param {number} [overloadedStatus] - Status del adjunto de contexto caido o del chat challenge (529 Anthropic / 503 OpenAI) * @returns {{ rateLimited: boolean, overloaded: boolean, status: number, retryAfter: number|null }} */ const describeUpstreamFailure = (error, fallbackStatus = 502, overloadedStatus = 529) => { + if (isWafChallengeError(error)) { + return { + rateLimited: false, + overloaded: true, + status: overloadedStatus, + retryAfter: Number(error.retryAfter) || CHAT_CHALLENGE_RETRY_AFTER_SECONDS + }; + } if (isContextAttachmentError(error)) { return { rateLimited: false, @@ -254,7 +267,8 @@ const describeUpstreamFailure = (error, fallbackStatus = 502, overloadedStatus = * HTTP 4xx/5xx) no enfria a proposito, y ese es justo el hueco. * * El require es perezoso: account.js arranca temporizadores al cargarse y no debe - * entrar en la cadena de carga de este modulo, que es puro. + * entrar en la cadena de carga de este modulo, que no arranca ninguno (el unico estado + * que guarda es el cortacircuitos del chat challenge, mas abajo). * @param {unknown} error - Error capturado en el controlador * @param {{email?: string}|null} [account] - Cuenta que sirvio la peticion * @returns {boolean} true si se marco la cuenta @@ -272,13 +286,131 @@ const noteRateLimitedAccount = (error, account) => { } }; +/** + * Chat challenge: Qwen se niega a GENERAR ("被挤爆啦") mientras crear el chat y subir el + * historial siguen pasando. Medido en prod 2026-09-23..26: 477 de 502 envios desafiados, + * todos entre 07:00Z y 22:00Z (el pico de Pekin); la misma cuenta pasa de noche y cae de + * dia, asi que no es la cuenta ni el tamano del contexto. Cada reintento inmediato del + * cliente agentico re-sube su historial y agota el limitador de parse en segundos. + * + * Cortacircuitos (gemelo en intencion del de parse en upload.js, pero con media apertura): + * - cerrado: cada desafio suma un strike; una respuesta con `choices` los borra. + * - abierto: tras CHAT_BREAKER_STRIKES seguidos, sendChatRequest contesta 529/503 sin + * tocar Qwen durante `chatChallengeBreakerSeconds`. Una respuesta de un stream que ya + * estaba en curso borra strikes pero NO cierra: no prueba que Qwen acepte peticiones nuevas. + * - media apertura: la primera peticion tras el enfriamiento sale como UNICA sonda y + * rearma la ventana para las demas. Si la desafia, reabre; la primera respuesta de Qwen + * mientras la sonda esta en vuelo cierra, venga de la sonda o de un stream anterior. + * ponytail: distinguir la sonda de un stream viejo exigiria etiquetar cada stream; el + * peor caso es que la sonda desafiada no reabra y hagan falta 3 strikes nuevos. + * Imagen/video no pasan por assertNoUpstreamFailure: llaman a noteChatAnswer al terminar bien. + */ +const CHAT_BUSY_MESSAGE = 'Qwen 上游繁忙,触发风控验证(被挤爆啦),请稍后重试 / Qwen chat challenge: upstream busy, retry later'; +const CHAT_CAPTCHA_MESSAGE = 'Qwen 上游要求人机验证(captcha),请稍后重试 / Qwen chat challenge: captcha required, retry later'; +const CHAT_BREAKER_MESSAGE = 'Qwen 上游连续触发风控验证,已暂停发送,请稍后重试 / Qwen chat challenge: repeated upstream challenges, requests paused, retry later'; +const CHAT_BUSY_SIGNAL_RE = /RGV587|被挤爆/; +const CHAT_CHALLENGE_RETRY_AFTER_SECONDS = 30; +const CHAT_BREAKER_STRIKES = 3; +// ponytail: un breaker global, porque todas las cuentas salen por la misma egress; por egress cuando haya varias. +const chatBreaker = { strikes: 0, openUntil: 0, probing: false }; + +// Reloj inyectable, como parseClock en upload.js: los tests avanzan la ventana sin dormir. +let chatClock = () => Date.now(); +const setChatChallengeClockForTests = (fn) => { chatClock = typeof fn === 'function' ? fn : () => Date.now(); }; + +// Perezosos como el require de account.js: config valida el entorno al cargarse. +const chatBreakerSeconds = () => Math.max(0, Number(require('../config/index.js').chatChallengeBreakerSeconds) || 0); +const logger = () => require('./logger').logger; + +const resetChatChallengeBreaker = () => { + chatBreaker.strikes = 0; + chatBreaker.openUntil = 0; + chatBreaker.probing = false; +}; + +/** @returns {number} Retry-After (s) para el desafio que acaba de llegar */ +const noteChatChallenge = () => { + chatBreaker.strikes += 1; + const seconds = chatBreakerSeconds(); + if (seconds <= 0 || (!chatBreaker.probing && chatBreaker.strikes < CHAT_BREAKER_STRIKES)) { + return CHAT_CHALLENGE_RETRY_AFTER_SECONDS; + } + chatBreaker.openUntil = chatClock() + seconds * 1000; + chatBreaker.probing = false; + logger().warn(`Qwen chat challenge 连续 ${chatBreaker.strikes} 次,${seconds}s 内不再发送聊天请求`, 'UPSTREAM'); + return seconds; +}; + +const noteChatAnswer = () => { + chatBreaker.strikes = 0; + if (!chatBreaker.probing) return; + chatBreaker.openUntil = 0; + chatBreaker.probing = false; + logger().info('Qwen chat challenge 探测请求已正常返回,恢复发送聊天请求', 'UPSTREAM'); +}; + +// Los SDK de Anthropic y OpenAI solo respetan un Retry-After por debajo de 60 s; con 60 o mas +// vuelven a su backoff corto. La ventana del breaker puede ser 60: el cliente vuelve 1 s antes +// y, si aun esta abierto, recibe el segundo que falta. +const CHAT_CHALLENGE_MAX_RETRY_AFTER_SECONDS = 59; + +const chatChallengeError = (message, retryAfter, details) => { + const error = new UpstreamResponseError(message, WAF_CHALLENGE_CODE, details); + error.retryAfter = Math.min(retryAfter, CHAT_CHALLENGE_MAX_RETRY_AFTER_SECONDS); + return error; +}; + +/** + * sendChatRequest y la ruta de imagen/video lo llaman antes de crear el chat o subir nada. + * Con el enfriamiento agotado deja pasar a quien llega primero como sonda. + * @returns {boolean} true si quien llama es la sonda + */ +const assertChatChallengeBreakerClosed = () => { + if (!chatBreaker.openUntil) return false; + const now = chatClock(); + const remaining = Math.ceil((chatBreaker.openUntil - now) / 1000); + if (remaining > 0) throw chatChallengeError(CHAT_BREAKER_MESSAGE, remaining, { breakerOpen: true }); + chatBreaker.openUntil = now + chatBreakerSeconds() * 1000; + chatBreaker.probing = true; + return true; +}; + +/** + * La sonda cayo antes de que Qwen contestara o la desafiara (validacion, chat_id, upload, + * transporte): sin esto bloquearia a todos otra ventana entera. La siguiente peticion sonda. + * ponytail: una sonda que llega a Qwen y termina sin respuesta ni desafio (stream vacio) + * sigue costando una ventana; soltarla ahi exigiria seguir cada stream hasta el final. + */ +const releaseChatProbe = () => { + if (!chatBreaker.probing) return; + chatBreaker.probing = false; + chatBreaker.openUntil = chatClock(); +}; + +/** + * Si el frame es un chat challenge devuelve su error (y cuenta el strike); si no, null. + * "被挤爆啦" es saturacion; un captcha/punish sin ella es verificacion humana. + * Deteccion pura: detectWafChallenge (compartida con la ruta de imagen); aqui solo el strike. + * @param {object|string} payload - Frame ya parseado, o el cuerpo crudo (pagina HTML del captcha) + * @returns {UpstreamResponseError|null} + */ +const chatChallengeFrom = (payload) => { + const detected = detectWafChallenge(payload); + if (!detected) return null; + // Sobre todo el payload, no solo las senales del detector: "被挤爆" puede venir sin "RGV587". + const busy = CHAT_BUSY_SIGNAL_RE.test(JSON.stringify(payload)); + return chatChallengeError(busy ? CHAT_BUSY_MESSAGE : CHAT_CAPTCHA_MESSAGE, noteChatChallenge(), detected.details); +}; + /** * Qwen Web 有时以 HTTP 200 + 普通 JSON 返回 WAF/captcha 或业务失败。 * 这些帧没有 choices,若直接跳过就会被误包装成空成功或正常 stop。 + * Alimenta el cortacircuitos del chat challenge: cada desafio suma un strike y cada frame + * con `choices` cuenta como respuesta (ver noteChatAnswer). */ const assertNoUpstreamFailure = (payload) => { - const wafChallenge = detectWafChallenge(payload); - if (wafChallenge) throw wafChallenge; + const challenge = chatChallengeFrom(payload); + if (challenge) throw challenge; if (!payload || typeof payload !== 'object') return; const explicitError = payload.error; @@ -302,6 +434,8 @@ const assertNoUpstreamFailure = (payload) => { waitHours === undefined || waitHours === null ? null : { waitHours } ); } + + if (Array.isArray(payload.choices)) noteChatAnswer(); }; module.exports = { @@ -312,6 +446,12 @@ module.exports = { isRateLimitError, isWafChallengeError, isTransportInterruption, + assertChatChallengeBreakerClosed, + chatChallengeFrom, + noteChatAnswer, + releaseChatProbe, + resetChatChallengeBreaker, + setChatChallengeClockForTests, rateLimitRetryAfterSeconds, describeUpstreamFailure, noteRateLimitedAccount, diff --git a/tests/agent-account-failover.test.js b/tests/agent-account-failover.test.js index 4a6b5b7b..30e7db72 100644 --- a/tests/agent-account-failover.test.js +++ b/tests/agent-account-failover.test.js @@ -13,7 +13,7 @@ const accountManager = require('../src/utils/account') const AccountRotator = require('../src/utils/account-rotator') const { runOpenAIAgentTurn } = require('../src/utils/openai-agent-runtime') const { createAccountReplayBody } = require('../src/utils/agent-account-failover') -const { assertNoUpstreamFailure, isRateLimitError, isWafChallengeError } = require('../src/utils/upstream-error') +const { assertNoUpstreamFailure, isRateLimitError, isWafChallengeError, resetChatChallengeBreaker } = require('../src/utils/upstream-error') const { handleStreamResponse, handleNonStreamResponse } = require('../src/controllers/chat') const accounts = ['first', 'second', 'third'].map(name => ({ email: `${name}@example.invalid`, token: `${name}-test-token` })) @@ -38,6 +38,8 @@ test.beforeEach(() => { accountManager.isInitialized = true accountManager.accountRotator = new AccountRotator() accountManager.accountRotator.setAccounts(accounts) + // WAF frames here feed the process-wide chat-challenge breaker; start every test closed. + resetChatChallengeBreaker() }) test.after(() => { accountManager.destroy() }) @@ -48,15 +50,10 @@ test('quota_limit is recognized without an English message; explicit WAF is not isWafChallengeError(error) && !isRateLimitError(error)) }) -test('exhausted or challenged pools never fall back to cooled accounts', () => { +test('exhausted pools never fall back to cooled accounts', () => { const rotator = accountManager.accountRotator rotator.recordQuotaExhausted(accounts[0].email) - rotator.recordChallenge(accounts[1].email) for (let count = 0; count < rotator.maxFailures; count += 1) rotator.recordFailure(accounts[2].email, 'ECONNRESET') - assert.equal(rotator.getNextAccount(), null) - rotator.resetFailures(accounts[1].email) - assert.equal(rotator.getAccountByEmail(accounts[1].email), null, 'refreshing a token must not clear WAF cooldown') - rotator.challengeCooldownUntil.set(accounts[1].email, Date.now() - 1) assert.equal(rotator.getNextAccount().email, accounts[1].email) assert.equal(rotator.getNextAccount([accounts[1].email]), null) }) @@ -130,30 +127,41 @@ test('no account rotation after client-visible final text or reasoning', async ( } }) -test('WAF can switch once but never fans out through the entire account pool', async () => { +test('WAF can switch once but never fans out through the entire account pool, and cools no account', async () => { let retries = 0 + const prefixKeys = [] await assert.rejects(runOpenAIAgentTurn(Readable.from([failureFrame('upstream_waf_challenge')]), options({ agent_turn_max_attempts: 6, + upstreamOptions: { contextPrefixKey: 'session-prefix' }, sendChatRequest: async (body, requestOptions) => { retries += 1 + prefixKeys.push(requestOptions.contextPrefixKey) return { status: true, response: Readable.from([failureFrame('upstream_waf_challenge')]), currentAccount: requestOptions.currentAccount } } })), error => isWafChallengeError(error) && error.failedAccountEmail === accounts[1].email) assert.equal(retries, 1) - assert.equal(accountManager.accountRotator.challengeCooldownUntil.size, 2) - assert.equal(accountManager.accountRotator.getNextAccount().email, accounts[2].email) + assert.deepEqual(prefixKeys, ['session-prefix'], 'Qwen refused before reading the history: the switch reuses it') + // A chat challenge follows Qwen's load, not the account (prod 2026-09-23..26: the same + // accounts were challenged by day and answered by night). Both stay in rotation. + for (const account of accounts.slice(0, 2)) { + assert.ok(accountManager.accountRotator.getAccountByEmail(account.email), `${account.email} stays available`) + } }) test('quota failover respects the total attempt budget and preserves the final quota error', async () => { let retries = 0 + const prefixKeys = [] await assert.rejects(runOpenAIAgentTurn(Readable.from([failureFrame()]), options({ agent_turn_max_attempts: 2, + upstreamOptions: { contextPrefixKey: 'session-prefix' }, sendChatRequest: async (body, requestOptions) => { retries += 1 + prefixKeys.push(requestOptions.contextPrefixKey) return { status: true, response: Readable.from([failureFrame()]), currentAccount: requestOptions.currentAccount } } })), error => isRateLimitError(error) && error.failedAccountEmail === accounts[1].email) assert.equal(retries, 1) + assert.deepEqual(prefixKeys, [null], 'a quota switch re-uploads the full history on the new account') }) test('no healthy replacement or disconnected client stops without new upstream requests', async () => { @@ -222,6 +230,42 @@ function createResponse() { } } +test('agent stream: a first-frame challenge switches accounts before anything is committed', async () => { + const response = createResponse() + await handleStreamResponse(response, Readable.from([failureFrame('upstream_waf_challenge')]), false, false, requestBody, options({ + sendChatRequest: async (body, requestOptions) => ({ status: true, response: finishedStream(), currentAccount: requestOptions.currentAccount }) + })) + assert.equal(response.statusCode, 200) + const chunks = response.output.split('\n\n').filter(block => block.startsWith('data: {')).map(block => JSON.parse(block.slice(6))) + const roles = chunks.filter(chunk => chunk.choices?.[0]?.delta?.role) + assert.equal(roles.length, 1, 'exactly one role chunk') + assert.equal(chunks[0].choices[0].delta.role, 'assistant', 'and it comes first') + assert.doesNotMatch(response.output, /"error"/) + assert.match(response.output, /OK/) +}) + +test('agent stream: challenged on both accounts is a real 503 with Retry-After, not a 200 error frame', async () => { + const response = createResponse() + await handleStreamResponse(response, Readable.from([failureFrame('upstream_waf_challenge')]), false, false, requestBody, options({ + sendChatRequest: async (body, requestOptions) => ({ + status: true, response: Readable.from([failureFrame('upstream_waf_challenge')]), currentAccount: requestOptions.currentAccount + }) + })) + assert.equal(response.statusCode, 503) + assert.ok(Number(response.headers['Retry-After']) > 0) + assert.equal(JSON.parse(response.output).error.code, 'upstream_unavailable') +}) + +test('a switch that cannot start keeps its cause (challenge or quota) instead of an opaque 502', async () => { + await assert.rejects(runOpenAIAgentTurn(Readable.from([failureFrame('upstream_waf_challenge')]), options({ + sendChatRequest: async () => ({ status: false, message: 'chat creation failed' }) + })), isWafChallengeError) + accountManager.accountRotator.reset() + await assert.rejects(runOpenAIAgentTurn(Readable.from([failureFrame()]), options({ + sendChatRequest: async () => ({ status: false, message: 'chat creation failed' }) + })), isRateLimitError) +}) + test('OpenAI controllers attribute successful JSON and SSE usage to the replacement only', async context => { const recorded = [] context.mock.method(accountManager, 'accumulateStats', (email, kind, usage) => recorded.push({ email, kind, usage })) diff --git a/tests/agent-protocol.test.js b/tests/agent-protocol.test.js index e7079ffc..a21b6df6 100644 --- a/tests/agent-protocol.test.js +++ b/tests/agent-protocol.test.js @@ -19,7 +19,7 @@ const { externalizeOversizedAgentContext, compactAgentContextFallback } = require('../src/utils/request.js') -const { assertNoUpstreamFailure } = require('../src/utils/upstream-error.js') +const { assertNoUpstreamFailure, resetChatChallengeBreaker } = require('../src/utils/upstream-error.js') const { shouldEnableToolRuntime, ensureAgentCurrentEnvelope @@ -35,6 +35,9 @@ test.after(() => { require('../src/utils/account.js').destroy() }) +// WAF frames feed the process-wide chat-challenge breaker; every test starts closed. +test.beforeEach(() => resetChatChallengeBreaker()) + const createMockResponse = () => ({ output: '', headers: {}, @@ -1106,9 +1109,11 @@ test('Qwen HTTP-200 bare JSON WAF response reaches OpenAI clients explicitly', a { messages: [] }, {} ) - assert.equal(res.statusCode, 502) - assert.match(res.output, /upstream_waf_challenge/) - assert.match(res.output, /WAF\\u002fcaptcha|WAF\/captcha/) + // A chat challenge is Qwen saying "busy, retry later": retryable 503 with a wait, not a 502. + assert.equal(res.statusCode, 503) + assert.equal(res.headers['Retry-After'], '30') + assert.match(res.output, /upstream_unavailable/) + assert.match(res.output, /chat challenge/) }) test('Anthropic stream emits thinking signature, max_tokens and tool parse errors', async () => { diff --git a/tests/chat-challenge.test.js b/tests/chat-challenge.test.js new file mode 100644 index 00000000..a1159a52 --- /dev/null +++ b/tests/chat-challenge.test.js @@ -0,0 +1,509 @@ +const test = require('node:test') +const assert = require('node:assert/strict') +const { Readable } = require('node:stream') +const axios = require('axios') + +process.env.API_KEY = 'chat-challenge-test-key' +process.env.DATA_SAVE_MODE = 'none' +process.env.ACCOUNTS = '' +process.env.ENABLE_CLI = 'false' +process.env.ENABLE_FILE_LOG = 'false' +process.env.PROXY_URL = '' + +const accountManager = require('../src/utils/account') +const modelsMap = require('../src/models/models-map.js') +modelsMap.getLatestModels = async () => { throw new Error('offline test: no model fetch') } +const requestModule = require('../src/utils/request') +// The real one, captured before the controllers below get a stub through the require cache. +const { sendChatRequest } = requestModule +const { + assertNoUpstreamFailure, + assertChatChallengeBreakerClosed, + resetChatChallengeBreaker, + setChatChallengeClockForTests, + describeUpstreamFailure, + isWafChallengeError, + isRateLimitError +} = require('../src/utils/upstream-error') + +let upstreamFrames = [] +const invalidatedPrefixes = [] +requestModule.sendChatRequest = async () => ({ + status: true, + response: typeof upstreamFrames === 'function' ? upstreamFrames() : Readable.from(upstreamFrames), + currentAccount: null, + contextPrefixReused: true +}) +requestModule.invalidateContextPrefix = key => { invalidatedPrefixes.push(key) } +requestModule.generateChatID = async () => 'image-test-chat' +const { handleAnthropicMessages } = require('../src/controllers/anthropic') +const { handleNonStreamResponse, handleStreamResponse } = require('../src/controllers/chat') +const { parseUpstreamImageError, handleImageVideoCompletion } = require('../src/controllers/chat.image.video') + +// The only frame Qwen sent on every chat challenge in prod, 2026-09-23..26 (477 of 502 sends). +const busy = { ret: ['FAIL_SYS_USER_VALIDATE', 'RGV587_ERROR::SM::哎哟喂,被挤爆啦,请稍后重试'] } +// The slider captcha of upstream #157: no "被挤爆", a punish URL instead. +const captcha = { ret: ['FAIL_SYS_USER_VALIDATE'], data: { url: 'https://chat.qwen.ai/api/v2/chat/completions/_____tmd_____/punish?x5secdata=x&action=captchaconnect' } } +const answer = { choices: [{ delta: { phase: 'answer', content: '你好' }, finish_reason: null }] } +const frame = payload => `data: ${JSON.stringify(payload)}\n\n` +const caught = fn => { try { fn() } catch (error) { return error } return null } +const strike = (times = 1) => { for (let n = 0; n < times; n += 1) caught(() => assertNoUpstreamFailure(busy)) } +const breakerOpen = () => caught(assertChatChallengeBreakerClosed) !== null +const settle = async (rounds = 25) => { for (let n = 0; n < rounds; n += 1) await new Promise(resolve => setImmediate(resolve)) } +const events = output => String(output).split('\n\n').map(block => /(?:^|\n)event: (.+)/.exec(block)?.[1]).filter(Boolean) +const created = { 'response.created': { chat_id: 'c1', parent_id: 'p1', response_id: 'r1', response_index: '0' } } +const stop = { choices: [{ delta: {}, finish_reason: 'stop' }], usage: { prompt_tokens: 2, completion_tokens: 1 } } +/** Yields `head`, then waits until released, then yields `tail`. */ +const gatedStream = (head, tail) => { + let release + const gate = new Promise(resolve => { release = resolve }) + return { stream: Readable.from((async function* () { + for (const chunk of head) yield Buffer.from(chunk) + await gate + for (const chunk of tail) yield Buffer.from(chunk) + })()), release } +} +const streamRequest = { body: { + model: 'qwen3-max', max_tokens: 64, stream: true, + metadata: { user_id: 'chat-challenge-session' }, + messages: [{ role: 'user', content: '你好' }] +} } + +let now = 1_000_000 +const mockResponse = () => ({ + output: '', + headers: {}, + headersSent: false, + writableEnded: false, + statusCode: 200, + set(headers) { Object.assign(this.headers, headers); return this }, + setHeader(name, value) { this.headers[name] = value }, + write(chunk) { this.headersSent = true; this.output += String(chunk); return true }, + end(chunk = '') { if (chunk) this.write(chunk); this.writableEnded = true }, + status(code) { this.statusCode = code; return this }, + json(value) { this.headersSent = true; this.output += JSON.stringify(value); this.writableEnded = true; return this } +}) + +test.before(async () => { + await accountManager._initPromise + setChatChallengeClockForTests(() => now) +}) +test.beforeEach(() => { + resetChatChallengeBreaker() + upstreamFrames = [] + invalidatedPrefixes.length = 0 +}) +test.after(() => { + setChatChallengeClockForTests(null) + accountManager.destroy() +}) + +test('a chat challenge is a retryable overload, not a 502/500, and blames neither context nor account', () => { + const error = caught(() => assertNoUpstreamFailure(busy)) + assert.ok(isWafChallengeError(error) && !isRateLimitError(error)) + assert.match(error.publicMessage, /被挤爆啦.*upstream busy/) + assert.doesNotMatch(error.publicMessage, /上下文|账号/) + assert.deepEqual(describeUpstreamFailure(error, 500), + { rateLimited: false, overloaded: true, status: 529, retryAfter: 30 }) + assert.equal(describeUpstreamFailure(error, 502, 503).status, 503) +}) + +test('a slider captcha is still a chat challenge, but is not reported as a busy upstream', () => { + const error = caught(() => assertNoUpstreamFailure(captcha)) + assert.ok(isWafChallengeError(error)) + assert.match(error.publicMessage, /captcha required/) + assert.doesNotMatch(error.publicMessage, /被挤爆|busy/) +}) + +test('"被挤爆" is reported as a busy upstream even without the RGV587 code', () => { + const error = caught(() => assertNoUpstreamFailure({ ret: ['FAIL_SYS_USER_VALIDATE', '哎哟喂,被挤爆啦,请稍后重试'] })) + assert.match(error.publicMessage, /upstream busy/) +}) + +test('three chat challenges in a row stop new requests before they reach Qwen', () => { + strike(2) + assert.equal(caught(assertChatChallengeBreakerClosed), null, 'two strikes keep it closed') + assert.equal(caught(() => assertNoUpstreamFailure(busy)).retryAfter, 59, 'the third asks for the full cooldown, capped under 60') + + const blocked = caught(assertChatChallengeBreakerClosed) + assert.ok(isWafChallengeError(blocked)) + assert.match(blocked.publicMessage, /requests paused/) + assert.deepEqual(describeUpstreamFailure(blocked, 500), + { rateLimited: false, overloaded: true, status: 529, retryAfter: 59 }) +}) + +test('an answer from a stream already in flight does not close an open breaker', () => { + strike(3) + assertNoUpstreamFailure(answer) + now += 30_000 + assert.equal(caught(assertChatChallengeBreakerClosed).retryAfter, 30, 'Retry-After is the time actually left') +}) + +test('after the cooldown exactly one probe goes out, and its answer closes the breaker', () => { + strike(3) + now += 60_000 + assert.equal(caught(assertChatChallengeBreakerClosed), null, 'the first request is the probe') + assert.equal(caught(assertChatChallengeBreakerClosed).retryAfter, 59, 'the others wait for the probe') + assertNoUpstreamFailure(answer) + assert.equal(caught(assertChatChallengeBreakerClosed), null) + assert.equal(caught(assertChatChallengeBreakerClosed), null, 'closed, not probing') +}) + +test('a challenged probe reopens the breaker for a full cooldown', () => { + strike(3) + now += 60_000 + assertChatChallengeBreakerClosed() + assert.equal(caught(() => assertNoUpstreamFailure(busy)).retryAfter, 59) + assert.equal(caught(assertChatChallengeBreakerClosed).retryAfter, 59) +}) + +test('sendChatRequest refuses while the breaker is open, before creating a chat or uploading history', async () => { + const realPost = axios.post + axios.post = async () => assert.fail('an open breaker must not reach Qwen') + try { + strike(3) + await assert.rejects( + sendChatRequest({ model: 'qwen3-max', messages: [{ role: 'user', content: '你好' }] }, + { currentAccount: { email: 'breaker@example.invalid', token: 'breaker-test-token' } }), + error => isWafChallengeError(error) && describeUpstreamFailure(error, 500).status === 529 + ) + } finally { + axios.post = realPost + } +}) + +// The captcha page as text/html on /api/v2/chat/completions (upstream #179, 2026-09-19): no `data:` frame. +const captchaPage = 'Verification' + + '
' +const postReturning = (contentType, body) => async () => ({ + status: 200, + headers: { 'content-type': contentType }, + data: Readable.from([Buffer.from(body)]) +}) +const sendToQwen = () => sendChatRequest( + { model: 'qwen3-max', messages: [{ role: 'user', content: '你好' }] }, + { chatId: 'html-test-chat', currentAccount: { email: 'html@example.invalid', token: 'html-test-token' } } +) + +test('a captcha page served as text/html is a chat challenge: one strike and a 529, not an empty answer', async () => { + const realPost = axios.post + axios.post = postReturning('text/html; charset=utf-8', captchaPage) + try { + await assert.rejects(sendToQwen(), error => + isWafChallengeError(error) && describeUpstreamFailure(error, 500).status === 529 && error.retryAfter > 0) + strike(2) + assert.ok(breakerOpen(), 'the HTML challenge counted as one strike') + } finally { + axios.post = realPost + } +}) + +test('any other text/html body is handed on byte for byte, and counts no strike', async () => { + const realPost = axios.post + axios.post = postReturning('text/html', '502 Bad Gateway') + try { + const result = await sendToQwen() + assert.equal(result.status, true) + const chunks = [] + for await (const chunk of result.response) chunks.push(Buffer.from(chunk)) + assert.equal(Buffer.concat(chunks).toString(), '502 Bad Gateway') + strike(2) + assert.equal(breakerOpen(), false) + } finally { + axios.post = realPost + } +}) + +test('a captcha page cut off mid-body is still one chat challenge: no re-POST, no account penalty', async () => { + const realPost = axios.post + const realFailure = accountManager.recordAccountFailure + const failures = [] + let posts = 0 + accountManager.recordAccountFailure = (...args) => { failures.push(args) } + axios.post = async () => { + posts += 1 + return { status: 200, headers: { 'content-type': 'text/html; charset=utf-8' }, data: Readable.from((async function* () { + yield Buffer.from(captchaPage) + throw Object.assign(new Error('aborted'), { code: 'ECONNRESET' }) + })()) } + } + try { + await assert.rejects(sendToQwen(), isWafChallengeError) + assert.equal(posts, 1) + assert.deepEqual(failures, []) + } finally { + axios.post = realPost + accountManager.recordAccountFailure = realFailure + } +}) + +test('OpenAI non-stream answers a chat challenge with 503 upstream_unavailable and Retry-After', async () => { + const res = mockResponse() + await handleNonStreamResponse(res, Readable.from([frame(busy)]), false, false, 'qwen3-max', { messages: [] }, {}) + assert.equal(res.statusCode, 503) + assert.equal(res.headers['Retry-After'], '30') + const body = JSON.parse(res.output) + assert.equal(body.error.code, 'upstream_unavailable') + assert.match(body.error.message, /upstream busy/) +}) + +test('/v1/messages keeps a reused history prefix on a chat challenge, and still forgets it on other upstream errors', async () => { + const request = { body: { + model: 'qwen3-max', max_tokens: 64, stream: false, + metadata: { user_id: 'chat-challenge-session' }, + messages: [{ role: 'user', content: '你好' }] + } } + + upstreamFrames = [frame(busy)] + const challenged = mockResponse() + await handleAnthropicMessages(request, challenged) + assert.equal(challenged.statusCode, 529) + assert.deepEqual(invalidatedPrefixes, [], 'Qwen refused before reading the prefix: re-parsing it helps nobody') + + upstreamFrames = [frame({ success: false, data: { code: 'Bad_Request', details: 'bad file' } })] + await handleAnthropicMessages(request, mockResponse()) + assert.equal(invalidatedPrefixes.length, 1, 'any other upstream error may be the stale prefix') +}) + +test('/v1/messages stream: a challenge on the first frame is a real 529 with Retry-After, and one strike', async () => { + upstreamFrames = [frame(busy)] + const res = mockResponse() + await handleAnthropicMessages(streamRequest, res) + assert.equal(res.statusCode, 529) + assert.equal(res.headers['Retry-After'], '30') + assert.notEqual(res.headers['Content-Type'], 'text/event-stream', 'nothing was committed as SSE') + assert.deepEqual(events(res.output), [], 'no message_start, no in-stream error event') + assert.equal(JSON.parse(res.output).error.type, 'overloaded_error') + strike(1) + assert.equal(breakerOpen(), false, 'the request cost one strike, so two strikes keep it closed') + strike(1) + assert.equal(breakerOpen(), true) +}) + +test('/v1/messages stream commits on the first valid frame, without waiting for content', async () => { + const { stream, release } = gatedStream([frame(created)], [frame(answer), frame(stop), 'data: [DONE]\n\n']) + upstreamFrames = () => stream + const res = mockResponse() + const done = handleAnthropicMessages(streamRequest, res) + await settle() + assert.deepEqual(events(res.output), ['message_start'], 'committed on response.created') + release() + await done + assert.equal(res.statusCode, 200) + assert.ok(events(res.output).includes('message_stop')) +}) + +test('/v1/messages stream: the first ping caps the wait, and a later challenge is an in-stream event', async (t) => { + t.mock.timers.enable({ apis: ['setInterval'] }) + const { stream, release } = gatedStream([], [frame(busy)]) + upstreamFrames = () => stream + const res = mockResponse() + const done = handleAnthropicMessages(streamRequest, res) + await settle() + assert.equal(res.output, '', 'nothing before the cap') + t.mock.timers.tick(15000) + assert.deepEqual(events(res.output), ['message_start', 'ping'], 'message_start always precedes a ping') + release() + await done + assert.equal(res.statusCode, 200) + const error = String(res.output).split('\n\n').find(block => block.includes('event: error')) + const payload = JSON.parse(/data: (.+)/.exec(error)[1]) + assert.equal(payload.error.type, 'overloaded_error') + assert.equal(payload.error.retry_after, 30) +}) + +test('/v1/messages stream: any other first-frame failure still commits and goes out in the stream', async () => { + upstreamFrames = [frame({ success: false, data: { code: 'Bad_Request', details: 'roto' } })] + const res = mockResponse() + await handleAnthropicMessages(streamRequest, res) + assert.equal(res.statusCode, 200, 'not turned into a 5xx that SDKs retry against upload/parse') + assert.deepEqual(events(res.output), ['message_start', 'error']) +}) + +test('image stream with the breaker open answers 503 + Retry-After before touching Qwen', async () => { + const realPost = axios.post + axios.post = async () => assert.fail('an open breaker must not reach Qwen') + try { + strike(3) + const res = mockResponse() + await handleImageVideoCompletion({ body: { + stream: true, chat_type: 't2i', model: 'qwen3-max', messages: [{ role: 'user', content: 'a cat' }] + } }, res) + assert.equal(res.statusCode, 503) + assert.equal(res.headers['Retry-After'], '59', 'the 60 s window, capped under the SDK limit') + assert.equal(JSON.parse(res.output).code, 'upstream_waf_challenge') + } finally { + axios.post = realPost + } +}) + +test('/v1/messages stream: an attempt with no JSON frames still commits before the turn is settled', async () => { + upstreamFrames = ['data: [DONE]\n\n'] + const res = mockResponse() + await handleAnthropicMessages(streamRequest, res) + assert.equal(res.statusCode, 200, 'an empty turn is not a first-frame challenge: it stays in the stream') + const names = events(res.output) + assert.equal(names[0], 'message_start') + assert.ok(['message_stop', 'error'].includes(names.at(-1)), `closed in-protocol, got ${names.at(-1)}`) +}) + +test('video stream (SSE headers preset) with the breaker open answers 503 as JSON', async () => { + const realPost = axios.post + axios.post = async () => assert.fail('an open breaker must not reach Qwen') + try { + strike(3) + const res = mockResponse() + await handleImageVideoCompletion({ body: { + stream: true, chat_type: 't2v', model: 'qwen3-max', messages: [{ role: 'user', content: 'a cat' }] + } }, res) + assert.equal(res.statusCode, 503) + assert.equal(res.headers['Content-Type'], 'application/json', 't2v preset text/event-stream before the flow') + assert.equal(res.headers['Retry-After'], '59') + } finally { + axios.post = realPost + } +}) + +const AGENT = { has_tools: true, tool_choice: 'auto', allowed_tool_names: ['get_time'], agent_turn_max_attempts: 1 } +const agentChunks = output => String(output).split('\n\n').filter(block => block.startsWith('data: {')).map(block => JSON.parse(block.slice(6))) + +test('agent stream: any other first-frame failure still commits and goes out as an error frame', async () => { + const res = mockResponse() + await handleStreamResponse(res, Readable.from([frame({ success: false, data: { code: 'Bad_Request', details: 'roto' } })]), + false, false, { messages: [] }, AGENT) + assert.equal(res.statusCode, 200, 'not turned into a 5xx that SDKs retry against upload/parse') + const chunks = agentChunks(res.output) + assert.equal(chunks[0].choices[0].delta.role, 'assistant') + assert.ok(chunks.some(chunk => chunk.error), 'the failure is an in-stream error frame') +}) + +test('agent stream: a turn that ends without output still commits before its error frame', async () => { + const res = mockResponse() + await handleStreamResponse(res, Readable.from(['data: [DONE]\n\n']), false, false, { messages: [] }, AGENT) + assert.equal(res.statusCode, 200) + const chunks = agentChunks(res.output) + assert.equal(chunks[0].choices[0].delta.role, 'assistant') + assert.match(res.output, /data: \[DONE\]|"error"/) +}) + +test('agent stream: the SSE heartbeat commits the role chunk before its keepalive', async () => { + const { stream, release } = gatedStream([], [frame(busy)]) + const res = mockResponse() + const done = handleStreamResponse(res, stream, false, false, { messages: [] }, { ...AGENT, agent_processing_heartbeat_ms: 5 }) + await new Promise(resolve => setTimeout(resolve, 40)) + const role = res.output.indexOf('"role":"assistant"') + const keepalive = res.output.indexOf(': qwen2api-agent-keepalive') + assert.ok(role >= 0 && keepalive > role, 'role chunk first, then the keepalive comment') + release() + await done + assert.equal(res.statusCode, 200, 'committed by the heartbeat: the challenge goes out in the stream') + assert.ok(agentChunks(res.output).some(chunk => chunk.error?.code === 'upstream_unavailable')) +}) + +test('image and video generation map a chat challenge to a retryable 503', () => { + assert.deepEqual(parseUpstreamImageError(JSON.stringify(busy)), { + error: caught(() => assertNoUpstreamFailure(busy)).publicMessage, + code: 'upstream_waf_challenge', + status: 503, + retry_after: 30 + }) +}) + +// t2v asks for responseType 'json': a body that is not JSON reaches the controller as a string. +const t2v = { stream: false, chat_type: 't2v', model: 'qwen3-max', messages: [{ role: 'user', content: 'a cat' }] } +const t2i = { stream: false, chat_type: 't2i', model: 'qwen3-max', messages: [{ role: 'user', content: 'a cat' }] } +const withPost = async (post, fn) => { + const realPost = axios.post + axios.post = post + try { return await fn() } finally { axios.post = realPost } +} +const okImage = () => Readable.from([Buffer.from(frame({ choices: [{ delta: { phase: 'image_gen', content: '![image](https://cdn.example.com/a.png)' } }] }))]) + +for (const [form, body] of [ + ['the captcha page', captchaPage.replace('', '')], + ['a data:-wrapped challenge frame', frame(captcha)] +]) { + test(`t2v: ${form} is a 503 + Retry-After and one strike, never a "video"`, async () => { + const res = mockResponse() + await withPost(async () => ({ status: 200, data: body }), () => handleImageVideoCompletion({ body: t2v }, res)) + assert.equal(res.statusCode, 503) + assert.ok(Number(res.headers['Retry-After']) > 0) + assert.doesNotMatch(res.output, /alicdn|punish/) + strike(2) + assert.ok(breakerOpen(), 'exactly one strike from the video request') + }) +} + +test('t2v: a challenge sent with a non-2xx status counts one strike, not two', async () => { + const res = mockResponse() + await withPost(async () => { throw Object.assign(new Error('Request failed with status code 403'), { response: { status: 403, data: busy } }) }, + () => handleImageVideoCompletion({ body: t2v }, res)) + assert.equal(res.statusCode, 503) + strike(1) + assert.equal(breakerOpen(), false) + strike(1) + assert.ok(breakerOpen()) +}) + +test('an answered image request closes a half-open breaker and clears strikes', async () => { + strike(3) + now += 60_000 + const res = mockResponse() + await withPost(async () => ({ status: 200, data: okImage() }), () => handleImageVideoCompletion({ body: t2i }, res)) + assert.equal(res.statusCode, 200) + assert.equal(breakerOpen(), false, 'the image was the probe and Qwen answered it') + strike(1) + await withPost(async () => ({ status: 200, data: okImage() }), () => handleImageVideoCompletion({ body: t2i }, mockResponse())) + strike(2) + assert.equal(breakerOpen(), false, 'an answered image breaks the run of strikes') +}) + +test('an image request without a prompt does not take the single half-open probe', async () => { + strike(3) + now += 60_000 + const res = mockResponse() + await withPost(async () => assert.fail('an invalid request must not reach Qwen'), + () => handleImageVideoCompletion({ body: { ...t2i, messages: [{ role: 'user', content: '' }] } }, res)) + assert.equal(res.statusCode, 400) + assert.equal(breakerOpen(), false, 'the next caller still gets to be the probe') +}) + +test('a text probe that fails before asking Qwen hands the probe to the next request', async () => { + strike(3) + now += 60_000 + const result = await withPost(async () => assert.fail('no account: nothing reaches Qwen'), + () => sendChatRequest({ model: 'qwen3-max', messages: [{ role: 'user', content: '你好' }] }, {})) + assert.equal(result.status, false) + assert.equal(breakerOpen(), false, 'released: the next caller is the probe, not refused for a window') +}) + +test('an image_edit without a text part does not keep the probe', async () => { + strike(3) + now += 60_000 + const res = mockResponse() + await withPost(async () => assert.fail('an invalid request must not reach Qwen'), () => handleImageVideoCompletion({ body: { + stream: false, chat_type: 'image_edit', model: 'qwen3-max', messages: [{ role: 'user', content: [{ type: 'image', image: 'https://cdn.example.com/a.png' }] }] + } }, res)) + assert.equal(res.statusCode, 400) + assert.equal(breakerOpen(), false) +}) + +test('t2v closes a half-open breaker when Qwen accepts the task, not minutes later when the video is ready', async () => { + strike(3) + now += 60_000 + const realGet = axios.get + let closedWhilePolling = null + axios.get = async () => { + closedWhilePolling = !breakerOpen() + return { data: { task_status: 'success', content: 'https://cdn.example.com/v.mp4' } } + } + try { + const res = mockResponse() + await withPost(async () => ({ status: 200, data: { success: true, data: { task_id: 'task-123' } } }), + () => handleImageVideoCompletion({ body: t2v }, res)) + assert.equal(res.statusCode, 200) + assert.equal(closedWhilePolling, true) + } finally { + axios.get = realGet + } +}) diff --git a/tests/expected-counts.json b/tests/expected-counts.json index ea0325cb..827d82a2 100644 --- a/tests/expected-counts.json +++ b/tests/expected-counts.json @@ -1,5 +1,5 @@ { - "tests": 1187, + "tests": 1225, "suites": 137, "note": "Authoritative count. Verify with the per-file sum in AGENTS.md (\"The test gate\"). The -a on that grep is load-bearing: tool-prompt.test.js emits bytes that make grep call the stream binary, and without -a its whole summary line — 133 tests — is silently dropped from the sum.", "updated": "2026-09-27" diff --git a/tests/upstream-quota-429.test.js b/tests/upstream-quota-429.test.js index 4b8f2356..18ce61e1 100644 --- a/tests/upstream-quota-429.test.js +++ b/tests/upstream-quota-429.test.js @@ -242,9 +242,9 @@ describe('/v1/messages: la cuota agotada sale como 429 rate_limit_error', () => }); it('streaming: a media transmision sale el EVENTO de error, no un cierre pelado', async () => { - // Aqui las cabeceras ya salieron (message_start se escribe antes de consumir el - // upstream), asi que el status HTTP ya no se puede cambiar: el unico canal que le - // queda al cliente para distinguir cuota de averia es el `type` del evento. + // Aqui las cabeceras ya salieron (el primer frame valido comprometio message_start), + // asi que el status HTTP ya no se puede cambiar: el unico canal que le queda al + // cliente para distinguir cuota de averia es el `type` del evento. upstreamFactory = () => streamOf([answerFrame('Voy a mirar'), quotaFrame()]); const res = streamRes(); await handleAnthropicMessages({ @@ -360,19 +360,28 @@ describe('/v1/chat/completions: la cuota agotada sale como 429 insufficient_quot assert.equal(String(res.headers['Retry-After']), '14400', '4 h == 14400 s'); }); - it('agentico streaming: el 429 es INALCANZABLE, y por eso el frame carga la senal', async () => { - // handleOpenAIAgentStream escribe el delta de apertura ({role:'assistant'}, chat.js:407) - // ANTES de consumir el upstream, asi que cuando llega el paquete de cuota la respuesta - // ya esta comprometida con 200 y el status HTTP no se puede cambiar. Para un cliente - // agentico con tools y stream el `type` del frame es el UNICO canal que queda. De ahi - // que arreglar el frame sea la mitad que de verdad sostiene este camino, no un extra. - // El gemelo de /v1/messages tiene la misma inalcanzabilidad, y por la misma razon: - // ver 'MATRIZ DE ALCANZABILIDAD' mas abajo. + it('agentico streaming: la cuota en el 1er frame es un 429 real, sin delta de apertura', async () => { + // handleOpenAIAgentStream compromete el delta de apertura ({role:'assistant'}) con el + // primer frame valido del upstream (commitStream). La cuota como primer frame llega antes: + // la respuesta sigue libre y sale un 429 real. Gemelo de /v1/messages. const res = streamRes(); await handleStreamResponse( res, streamOf([quotaFrame()]), false, false, { messages: [{ role: 'user', content: 'hola' }] }, AGENT_OPTS ); + assert.equal(res.statusCode, 429); + assert.doesNotMatch(res.output, /"role":"assistant"/, 'nada se comprometio'); + assert.equal(JSON.parse(res.output).error?.type, 'insufficient_quota'); + }); + + it('agentico streaming: a mitad de respuesta el frame carga la senal', async () => { + // Tras el primer frame valido la respuesta ya esta comprometida con 200: para un cliente + // agentico con tools y stream el `type` del frame es el UNICO canal que queda. + const res = streamRes(); + await handleStreamResponse( + res, streamOf([answerFrame('Voy a mirar'), quotaFrame()]), false, false, + { messages: [{ role: 'user', content: 'hola' }] }, AGENT_OPTS + ); assert.equal(res.headersSent, true, 'el delta de apertura ya comprometio la respuesta'); assert.equal(res.statusCode, 200, 'no se puede reescribir un status ya enviado'); @@ -416,28 +425,23 @@ describe('/v1/chat/completions: la cuota agotada sale como 429 insufficient_quot // =================================================================================== // MATRIZ DE ALCANZABILIDAD — lo que el cliente recibe DE VERDAD, por camino y por fase. // -// El commit original se titulaba "map ... to 429 on both API paths". Es falso para el -// unico modo que el usuario ejecuta. Claude Code habla /v1/messages con `stream: true`, -// y las 149 negativas de cuota de sus logs son TODAS "mid-stream". En ese modo el 429 -// no existe: handleAnthropicStream fija las cabeceras (anthropic.js:1203) y escribe -// message_start (:1211) ANTES de leer un solo byte del upstream, asi que cuando el -// paquete de cuota llega `res.headersSent` ya es true y el catch del controlador -// (:2659) solo puede tomar la rama del evento. +// Claude Code habla /v1/messages con `stream: true`. Antes, handleAnthropicStream fijaba +// las cabeceras y escribia message_start ANTES de leer un byte del upstream, asi que una +// cuota en el primer frame salia como 200 + evento. Ahora el compromiso es perezoso (primer +// frame valido o primer ping): la cuota —o un chat challenge— en el PRIMER frame es un 429 +// real con Retry-After; a mitad de respuesta (las 149 negativas de los logs del usuario son +// todas "mid-stream") el evento sigue siendo todo el canal y carga la espera. // -// camino fase status senal -// /v1/messages stream:false cabeceras aun libres 429 body.error.type -// /v1/messages stream:true SIEMPRE comprometida 200 evento error.type -// /v1/chat/... llano stream:false cabeceras aun libres 429 body.error.type -// /v1/chat/... llano stream:true libre hasta el 1er byte 429 body.error.type -// /v1/chat/... agente stream:true SIEMPRE comprometida 200 frame error.type +// camino cuota en el 1er frame cuota a mitad de respuesta +// /v1/messages stream:false 429 + Retry-After — +// /v1/messages stream:true 429 + Retry-After 200 + evento error.type (+retry_after) +// /v1/chat/... llano stream:false 429 + Retry-After — +// /v1/chat/... llano stream:true 429 + Retry-After 200 + frame error.type (+retry_after) +// /v1/chat/... agente stream:true 429 + Retry-After 200 + frame error.type (+retry_after) // -// Estos casos clavan la fila que la frase original negaba. Que el 429 sea inalcanzable -// no es un defecto a tapar: adelantar las cabeceras es lo que permite mandar `ping` -// dentro del protocolo (anthropic.js:1156-1160), que es como se elimino el falso -// "stream muerto" del puente ccproxy. Lo que SI era un defecto es que, sin cabecera, -// la espera se perdia — eso se arregla abajo. -describe('/v1/messages en streaming: el 429 es inalcanzable y el evento es todo el canal', () => { - it('con tools y la cuota como PRIMER frame: 200 comprometido, evento rate_limit_error', async () => { +// La tabla completa, por funcion, esta en src/utils/upstream-error.js. +describe('/v1/messages en streaming: la cuota en el 1er frame es un 429 real; despues, el evento', () => { + it('con tools y la cuota como PRIMER frame: 429 rate_limit_error, nada comprometido', async () => { upstreamFactory = () => streamOf([quotaFrame()]); const res = streamRes(); await handleAnthropicMessages({ @@ -450,33 +454,40 @@ describe('/v1/messages en streaming: el 429 es inalcanzable y el evento es todo } }, res); - assert.equal(res.headersSent, true, 'message_start ya comprometio la respuesta'); - assert.equal(res.statusCode, 200, 'el 429 NO es alcanzable en el modo que usa Claude Code'); - - const events = sseEvents(res.output); - assert.equal(events[0]?.event, 'message_start', 'las cabeceras salen antes que el upstream'); - const err = events.filter(e => e.event === 'error'); - assert.equal(err.length, 1); - assert.equal(err[0].data?.error?.type, 'rate_limit_error', 'el type es el unico canal que queda'); - assert.equal(res.headers['Retry-After'], undefined, 'ya no hay cabecera que poner'); + assert.equal(res.statusCode, 429, 'el primer frame llega antes del compromiso'); + assert.equal(sseEvents(res.output).length, 0, 'ni message_start ni evento: la respuesta es JSON'); + assert.notEqual(res.headers['Content-Type'], 'text/event-stream'); + assert.equal(JSON.parse(res.output).error?.type, 'rate_limit_error'); + assert.equal(res.headers['Retry-After'], undefined, 'sin espera real, sin cabecera'); }); - it('la espera del upstream viaja DENTRO del evento, que es donde el cliente puede verla', async () => { - // Si el evento es el unico canal, tiene que cargar todo lo que la cabecera ya no - // puede llevar. `data.num` viene en HORAS (misma lectura que chat.image.video.js:88). + it('cuota en el 1er frame con la espera del upstream: la cabecera Retry-After la lleva', async () => { upstreamFactory = () => streamOf([quotaFrame({ num: 3 })]); const res = streamRes(); await handleAnthropicMessages({ body: { model: 'qwen3-max', max_tokens: 64, stream: true, messages: [{ role: 'user', content: 'hola' }] } }, res); + assert.equal(res.statusCode, 429); + assert.equal(String(res.headers['Retry-After']), '10800', '3 h == 10800 s, en la cabecera'); + }); + + it('a mitad de respuesta la espera del upstream viaja DENTRO del evento', async () => { + // Ya comprometida, el evento es el unico canal: tiene que cargar todo lo que la cabecera + // ya no puede llevar. `data.num` viene en HORAS (misma lectura que chat.image.video.js:88). + upstreamFactory = () => streamOf([answerFrame('Voy a mirar'), quotaFrame({ num: 3 })]); + const res = streamRes(); + await handleAnthropicMessages({ + body: { model: 'qwen3-max', max_tokens: 64, stream: true, messages: [{ role: 'user', content: 'hola' }] } + }, res); + const err = sseEvents(res.output).filter(e => e.event === 'error'); assert.equal(err.length, 1); assert.equal(err[0].data?.error?.retry_after, 10800, '3 h == 10800 s, dentro del evento'); }); it('sin espera real el evento no se inventa ninguna', async () => { - upstreamFactory = () => streamOf([quotaFrame()]); + upstreamFactory = () => streamOf([answerFrame('Voy a mirar'), quotaFrame()]); const res = streamRes(); await handleAnthropicMessages({ body: { model: 'qwen3-max', max_tokens: 64, stream: true, messages: [{ role: 'user', content: 'hola' }] } @@ -517,10 +528,10 @@ describe('/v1/chat/completions en streaming: el frame carga la misma espera (gem assert.equal(errs[0].error.retry_after, 7200, '2 h == 7200 s, dentro del frame'); }); - it('agentico: el frame de error lleva retry_after — aqui el 429 no existe', async () => { + it('agentico: a mitad de respuesta el frame de error lleva retry_after', async () => { const res = streamRes(); await handleStreamResponse( - res, streamOf([quotaFrame({ num: 5 })]), false, false, + res, streamOf([answerFrame('Voy a mirar'), quotaFrame({ num: 5 })]), false, false, { messages: [{ role: 'user', content: 'hola' }] }, { has_tools: true, tool_choice: 'auto', allowed_tool_names: ['get_time'], agent_turn_max_attempts: 2 } ); @@ -530,6 +541,17 @@ describe('/v1/chat/completions en streaming: el frame carga la misma espera (gem assert.equal(errs[0].error.retry_after, 18000, '5 h == 18000 s'); }); + it('agentico: en el 1er frame la espera va en la cabecera Retry-After', async () => { + const res = streamRes(); + await handleStreamResponse( + res, streamOf([quotaFrame({ num: 5 })]), false, false, + { messages: [{ role: 'user', content: 'hola' }] }, + { has_tools: true, tool_choice: 'auto', allowed_tool_names: ['get_time'], agent_turn_max_attempts: 2 } + ); + assert.equal(res.statusCode, 429); + assert.equal(String(res.headers['Retry-After']), '18000'); + }); + it('un fallo que NO es de cuota no lleva retry_after en el frame', async () => { const res = streamRes(); await handleStreamResponse( diff --git a/tools/bun-smoke.js b/tools/bun-smoke.js index ed299936..4f833490 100644 --- a/tools/bun-smoke.js +++ b/tools/bun-smoke.js @@ -77,6 +77,7 @@ async function main() { let completionRequests = 0 let replyDelayMilliseconds = 20 let failoverMode = false + let challengeMode = false const failoverTokens = [] const upstreamServer = createServer((request, response) => { const pathname = new URL(request.url, 'http://localhost').pathname @@ -101,6 +102,17 @@ async function main() { } else if (pathname === '/api/v2/chat/completions') { completionRequests += 1 response.setHeader('Content-Type', 'text/event-stream') + if (challengeMode === 'html') { + // The same challenge as Aliyun's captcha page (upstream #179): 200 text/html, no data: frame. + response.setHeader('Content-Type', 'text/html; charset=utf-8') + response.end('Verification
') + return + } + if (challengeMode) { + // The one frame Qwen sends on a chat challenge (prod, 2026-09-23..26). + response.end(`data: ${JSON.stringify({ ret: ['FAIL_SYS_USER_VALIDATE', 'RGV587_ERROR::SM::哎哟喂,被挤爆啦,请稍后重试'] })}\n\n`) + return + } if (failoverMode) { const accountToken = request.headers.cookie?.split(';').find(part => part.trim().startsWith('token='))?.trim().slice(6) failoverTokens.push(accountToken) @@ -329,6 +341,46 @@ async function main() { assert.ok(failoverTokens.every(token => smokeAccounts.some(account => account.token === token))) assert.equal(new Set(failoverTokens).size, 2, 'Failover must send a different account token upstream') console.log('PASS: Agent mid-stream quota failover with distinct account credentials') + + // A chat challenge on the first frame must reach streaming clients as a real status with a + // Retry-After header over real Bun HTTP: SDKs ignore an in-stream error's retry_after. + challengeMode = true + const challengedMessages = await request('/v1/messages', { + method: 'POST', headers: { 'Content-Type': 'application/json', 'x-api-key': API_KEY }, + body: JSON.stringify({ model: MODEL, max_tokens: 16, stream: true, messages: [{ role: 'user', content: '你好' }] }) + }) + const challengedMessagesText = await challengedMessages.text() + assert.equal(challengedMessages.status, 529, challengedMessagesText) + assert.match(challengedMessages.headers.get('content-type') || '', /application\/json/) + assert.equal(challengedMessages.headers.get('retry-after'), '30') + assert.equal(JSON.parse(challengedMessagesText).error.type, 'overloaded_error') + const challengedAgent = await request('/v1/chat/completions', { + method: 'POST', headers: { 'Content-Type': 'application/json', 'x-api-key': API_KEY }, + body: JSON.stringify({ + model: MODEL, stream: true, messages: [{ role: 'user', content: '你好' }], + tools: [{ type: 'function', function: { name: 'get_time', parameters: { type: 'object', properties: {} } } }] + }) + }) + const challengedAgentText = await challengedAgent.text() + assert.equal(challengedAgent.status, 503, challengedAgentText) + assert.match(challengedAgent.headers.get('content-type') || '', /application\/json/) + assert.ok(Number(challengedAgent.headers.get('retry-after')) > 0) + assert.equal(JSON.parse(challengedAgentText).error.code, 'upstream_unavailable') + // The quota failover above cooled the other smoke account, so the agent has no account to + // switch to: one challenged send per request. + assert.equal(completionRequests, 8, 'one challenged send each for /v1/messages and the agent') + // The captcha page as text/html used to end as an empty answer after 2-3 sends. + challengeMode = 'html' + const htmlChallenged = await request('/v1/messages', { + method: 'POST', headers: { 'Content-Type': 'application/json', 'x-api-key': API_KEY }, + body: JSON.stringify({ model: MODEL, max_tokens: 16, stream: true, messages: [{ role: 'user', content: '你好' }] }) + }) + const htmlChallengedText = await htmlChallenged.text() + assert.equal(htmlChallenged.status, 529, htmlChallengedText) + assert.ok(Number(htmlChallenged.headers.get('retry-after')) > 0) + assert.equal(completionRequests, 9, 'a captcha page is one send, not an empty-answer retry loop') + challengeMode = false + console.log('PASS: chat challenge on the first frame is a real 529/503 with Retry-After on streams') if (standalone) { const settingsResponse = await request('/api/setRetryConfig', { method: 'POST', headers: { 'Content-Type': 'application/json', 'x-api-key': API_KEY },