From 81b75adbd5c832d33ef45dda37182d7b9b1d578c Mon Sep 17 00:00:00 2001 From: PEDRO LOBATO CARCAMO Date: Sat, 26 Sep 2026 10:57:30 -0600 Subject: [PATCH 1/7] fix(upstream): treat a chat challenge as a retryable overload and stop hammering Qwen during one MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Since 2026-09-23 11:35Z Qwen answers most chat generations with a single frame `ret: ['FAIL_SYS_USER_VALIDATE', 'RGV587_ERROR::SM::哎哟喂,被挤爆啦,请稍后重试']`, even for a bare "你好" (#180). Measured on a production deployment, 2026-09-23..26: 477 of 502 chat sends challenged, every one between 07:00Z and 22:00Z; the 25 that passed were all at night (UTC), and the same accounts passed by night and were challenged by day. Creating the chat and uploading the history file kept succeeding on the same account and egress seconds before each challenge. So it is neither the context size nor the account. What clients got instead: 502 upstream_error (OpenAI) or 500 api_error (Anthropic) with no Retry-After, and the message "Agent 上下文可能过大或账号需要 验证". Agentic clients retry within a second, each retry re-uploads and re-parses the history, and the parse rate limiter answers 529 on top. - A chat challenge now maps to 529 overloaded_error / 503 with Retry-After: 30, and says what Qwen said: upstream busy, retry later. - After 3 challenges in a row, sendChatRequest answers 529/503 for 60 s before creating a chat or uploading anything; the first response with `choices` closes the breaker. One breaker for the process: every account shares the same egress in the deployments this was measured on. - No account cooldown or warning for a challenge: it cooled healthy accounts for 5 minutes. The single agent-mode failover stays. - /v1/messages no longer forgets a reused history prefix on a chat challenge (Qwen rejected before reading it), so the client's retry does not re-parse. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/controllers/anthropic.js | 7 ++- src/utils/account-rotator.js | 24 +------- src/utils/account.js | 4 -- src/utils/agent-account-failover.js | 6 +- src/utils/openai-agent-runtime.js | 2 +- src/utils/request.js | 4 +- src/utils/upstream-error.js | 61 +++++++++++++++++++- tests/agent-account-failover.test.js | 16 +++--- tests/agent-protocol.test.js | 7 ++- tests/chat-challenge.test.js | 86 ++++++++++++++++++++++++++++ 10 files changed, 169 insertions(+), 48 deletions(-) create mode 100644 tests/chat-challenge.test.js diff --git a/src/controllers/anthropic.js b/src/controllers/anthropic.js index 8ce770ac..48b88077 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 @@ -2802,8 +2803,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/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..8386b624 100644 --- a/src/utils/openai-agent-runtime.js +++ b/src/utils/openai-agent-runtime.js @@ -795,7 +795,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 diff --git a/src/utils/request.js b/src/utils/request.js index ccbec424..de9c0559 100644 --- a/src/utils/request.js +++ b/src/utils/request.js @@ -7,7 +7,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 } = require('./upstream-error.js') const { contextPrefixCache, prefixMatches, canonicalHistoryHash } = require('./context-prefix-cache.js') const { TOOL_CALL_OPEN, LEDGER_HEADER, LEDGER_CAPTION, truncateToolHistoryLedger, stripRetainedThinking @@ -642,6 +642,8 @@ 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. + assertChatChallengeBreakerClosed() // 获取可用的账户(包含 proxy 等完整字段) // excludeEmails:本次 HTTP 请求里已经烧掉的账户(流中途 failover)——轮换器跳过它们, // 即使它们对其他请求仍然可用。 diff --git a/src/utils/upstream-error.js b/src/utils/upstream-error.js index 2414536f..8f9dbc49 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. @@ -226,10 +226,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, @@ -272,13 +280,55 @@ 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. Tras + * CHAT_BREAKER_STRIKES desafios seguidos se contesta 529/503 sin tocar Qwen durante + * CHAT_BREAKER_SECONDS; pasado ese tiempo sale una sonda, y la primera respuesta con + * `choices` lo cierra. + */ +const CHAT_CHALLENGE_MESSAGE = 'Qwen 上游繁忙,触发风控验证(被挤爆啦),请稍后重试 / Qwen chat challenge: upstream busy, retry later'; +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 CHAT_BREAKER_SECONDS = 60; +const chatBreaker = { strikes: 0, openUntil: 0 }; + +const resetChatChallengeBreaker = () => { + chatBreaker.strikes = 0; + chatBreaker.openUntil = 0; +}; + +/** @returns {number} Retry-After (s) para el desafio que acaba de llegar */ +const noteChatChallenge = () => { + chatBreaker.strikes += 1; + if (chatBreaker.strikes < CHAT_BREAKER_STRIKES) return CHAT_CHALLENGE_RETRY_AFTER_SECONDS; + chatBreaker.openUntil = Date.now() + CHAT_BREAKER_SECONDS * 1000; + return CHAT_BREAKER_SECONDS; +}; + +const chatChallengeError = (retryAfter, details) => { + const error = new UpstreamResponseError(CHAT_CHALLENGE_MESSAGE, WAF_CHALLENGE_CODE, details); + error.retryAfter = retryAfter; + return error; +}; + +/** sendChatRequest lo llama antes de crear el chat o subir el historial. */ +const assertChatChallengeBreakerClosed = () => { + const remaining = Math.ceil((chatBreaker.openUntil - Date.now()) / 1000); + if (remaining > 0) throw chatChallengeError(remaining, { breakerOpen: true }); +}; + /** * Qwen Web 有时以 HTTP 200 + 普通 JSON 返回 WAF/captcha 或业务失败。 * 这些帧没有 choices,若直接跳过就会被误包装成空成功或正常 stop。 */ const assertNoUpstreamFailure = (payload) => { const wafChallenge = detectWafChallenge(payload); - if (wafChallenge) throw wafChallenge; + if (wafChallenge) throw chatChallengeError(noteChatChallenge(), wafChallenge.details); if (!payload || typeof payload !== 'object') return; const explicitError = payload.error; @@ -302,6 +352,9 @@ const assertNoUpstreamFailure = (payload) => { waitHours === undefined || waitHours === null ? null : { waitHours } ); } + + // Qwen volvio a generar: la sonda paso. + if (Array.isArray(payload.choices)) resetChatChallengeBreaker(); }; module.exports = { @@ -312,6 +365,8 @@ module.exports = { isRateLimitError, isWafChallengeError, isTransportInterruption, + assertChatChallengeBreakerClosed, + resetChatChallengeBreaker, rateLimitRetryAfterSeconds, describeUpstreamFailure, noteRateLimitedAccount, diff --git a/tests/agent-account-failover.test.js b/tests/agent-account-failover.test.js index 4a6b5b7b..f8a7f2a2 100644 --- a/tests/agent-account-failover.test.js +++ b/tests/agent-account-failover.test.js @@ -48,15 +48,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,7 +125,7 @@ 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 await assert.rejects(runOpenAIAgentTurn(Readable.from([failureFrame('upstream_waf_challenge')]), options({ agent_turn_max_attempts: 6, @@ -140,8 +135,11 @@ test('WAF can switch once but never fans out through the entire account pool', a } })), 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) + // 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 () => { diff --git a/tests/agent-protocol.test.js b/tests/agent-protocol.test.js index e7079ffc..a1a1372b 100644 --- a/tests/agent-protocol.test.js +++ b/tests/agent-protocol.test.js @@ -1094,6 +1094,7 @@ test('Qwen HTTP-200 WAF payload is surfaced as an explicit failure', () => { }) test('Qwen HTTP-200 bare JSON WAF response reaches OpenAI clients explicitly', async () => { + require('../src/utils/upstream-error.js').resetChatChallengeBreaker() const res = createMockResponse() await handleStreamResponse( res, @@ -1106,9 +1107,11 @@ test('Qwen HTTP-200 bare JSON WAF response reaches OpenAI clients explicitly', a { messages: [] }, {} ) - assert.equal(res.statusCode, 502) + // 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_waf_challenge/) - assert.match(res.output, /WAF\\u002fcaptcha|WAF\/captcha/) + 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..49f7b114 --- /dev/null +++ b/tests/chat-challenge.test.js @@ -0,0 +1,86 @@ +const test = require('node:test') +const assert = require('node:assert/strict') +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 { sendChatRequest } = require('../src/utils/request') +const { + assertNoUpstreamFailure, + assertChatChallengeBreakerClosed, + resetChatChallengeBreaker, + describeUpstreamFailure, + isWafChallengeError, + isRateLimitError +} = require('../src/utils/upstream-error') + +// The only frame Qwen sent on every chat challenge in prod, 2026-09-23..26 (477 of 502 sends). +const challenge = { ret: ['FAIL_SYS_USER_VALIDATE', 'RGV587_ERROR::SM::哎哟喂,被挤爆啦,请稍后重试'] } +const answer = { choices: [{ delta: { phase: 'answer', content: '你好' }, finish_reason: null }] } +const thrown = fn => { try { fn() } catch (error) { return error } return null } + +test.before(async () => { await accountManager._initPromise }) +test.beforeEach(() => resetChatChallengeBreaker()) +test.after(() => { accountManager.destroy() }) + +test('a chat challenge is a retryable overload, not a 502/500, and blames neither context nor account', () => { + const error = thrown(() => assertNoUpstreamFailure(challenge)) + assert.ok(isWafChallengeError(error) && !isRateLimitError(error)) + 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('three chat challenges in a row stop new requests; a request that gets an answer closes the breaker', () => { + thrown(() => assertNoUpstreamFailure(challenge)) + thrown(() => assertNoUpstreamFailure(challenge)) + assert.equal(thrown(assertChatChallengeBreakerClosed), null, 'two strikes keep it closed') + assert.equal(thrown(() => assertNoUpstreamFailure(challenge)).retryAfter, 60, 'the third asks for the full cooldown') + + const blocked = thrown(assertChatChallengeBreakerClosed) + assert.ok(isWafChallengeError(blocked)) + assert.deepEqual(describeUpstreamFailure(blocked, 500), + { rateLimited: false, overloaded: true, status: 529, retryAfter: 60 }) + + assertNoUpstreamFailure(answer) + assert.equal(thrown(assertChatChallengeBreakerClosed), null) +}) + +test('an open breaker lets requests through again after its cooldown, and the first challenged one reopens it', () => { + const realNow = Date.now + let now = realNow() + Date.now = () => now + try { + for (let strike = 0; strike < 3; strike += 1) thrown(() => assertNoUpstreamFailure(challenge)) + now += 30_000 + assert.equal(thrown(assertChatChallengeBreakerClosed).retryAfter, 30, 'Retry-After is the time actually left') + now += 30_000 + assert.equal(thrown(assertChatChallengeBreakerClosed), null, 'cooldown over: the probe goes out') + thrown(() => assertNoUpstreamFailure(challenge)) + assert.equal(thrown(assertChatChallengeBreakerClosed).retryAfter, 60) + } finally { + Date.now = realNow + } +}) + +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 { + for (let strike = 0; strike < 3; strike += 1) thrown(() => assertNoUpstreamFailure(challenge)) + 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 + } +}) From b19384bb49a23d09fef6c28860846e831a3b704f Mon Sep 17 00:00:00 2001 From: PEDRO LOBATO CARCAMO Date: Sat, 26 Sep 2026 11:22:06 -0600 Subject: [PATCH 2/7] fix(upstream): chat-challenge breaker lets one probe through; cover images and failover MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-ups on the chat-challenge breaker: - Half-open state: once the cooldown ends, the first request goes out as the only probe and re-arms the window for the rest. The probe's answer closes the breaker; a challenge reopens it. Before, every waiting client came back in the same second. - An answer from a stream that was already running when the breaker opened now only clears strikes. It no longer closes an open breaker, since it says nothing about whether Qwen accepts new requests. - CHAT_CHALLENGE_BREAKER_SECONDS (default 60, 0 disables), an injectable clock for tests, a warn log on open and an info log on close. This mirrors the parse breaker. - A slider captcha (`/punish?`, no 被挤爆) gets its own message instead of "upstream busy". The breaker-open refusal says requests are paused. - The plain /v1/chat/completions handlers go through upstreamErrorShape like agent mode, so a chat challenge is 503 `upstream_unavailable` there too. - Image and video generation post to the same chat endpoint. They now check the breaker before creating a chat, recognise the challenge frame, and answer 503 with Retry-After instead of a generic 500. - A chat-challenge account switch in agent mode keeps the history prefix instead of re-uploading and re-parsing the whole history. Qwen refused before reading it, and normal rotation already reuses the prefix across accounts. A quota switch still re-uploads. - Tests reset the process-wide breaker in beforeEach. New tests cover: the half-open probe, an in-flight answer, the captcha message, OpenAI non-stream 503, /v1/messages keeping the prefix on a challenge (and forgetting it on other errors), the image path, and the failover prefix key. Co-Authored-By: Claude Opus 5.5 (1M context) --- .env.example | 9 +- src/config/index.js | 7 ++ src/controllers/chat.image.video.js | 55 ++++++---- src/controllers/chat.js | 39 ++----- src/utils/openai-agent-runtime.js | 13 ++- src/utils/upstream-error.js | 94 ++++++++++++---- tests/agent-account-failover.test.js | 12 +- tests/agent-protocol.test.js | 8 +- tests/chat-challenge.test.js | 158 ++++++++++++++++++++++----- 9 files changed, 284 insertions(+), 111 deletions(-) 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/chat.image.video.js b/src/controllers/chat.image.video.js index 2ce6aaec..68c4a7a9 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 } = 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) @@ -598,6 +595,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 +680,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: { @@ -1296,6 +1295,12 @@ const generateImageVideoResult = async (payload) => { ] } + try { + assertChatChallengeBreakerClosed() + } catch (challenge) { + throw imageChallengeError(challenge) + } + const chatID = await generateChatID(token, model, account, chat_type) if (!chatID) { @@ -1435,6 +1440,8 @@ const generateImageVideoResult = async (payload) => { responseData = await axios.post(`${chatBaseUrl}/api/v2/chat/completions?chat_id=${chatID}`, reqBody, requestConfig) const inlineUpstreamError = parseUpstreamImageError(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) @@ -1789,5 +1796,7 @@ module.exports = { handleImageVideoCompletion, handleOpenAIImagesGeneration, handleOpenAIImagesEdit, - handleOpenAIVideoGeneration + handleOpenAIVideoGeneration, + // 暴露内部辅助以便测试 + parseUpstreamImageError } diff --git a/src/controllers/chat.js b/src/controllers/chat.js index ff475a87..4512eba8 100644 --- a/src/controllers/chat.js +++ b/src/controllers/chat.js @@ -1079,34 +1079,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 +1425,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/openai-agent-runtime.js b/src/utils/openai-agent-runtime.js index 8386b624..c616734f 100644 --- a/src/utils/openai-agent-runtime.js +++ b/src/utils/openai-agent-runtime.js @@ -807,9 +807,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, diff --git a/src/utils/upstream-error.js b/src/utils/upstream-error.js index 8f9dbc49..00455763 100644 --- a/src/utils/upstream-error.js +++ b/src/utils/upstream-error.js @@ -262,7 +262,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 @@ -285,50 +286,102 @@ const noteRateLimitedAccount = (error, account) => { * 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. Tras - * CHAT_BREAKER_STRIKES desafios seguidos se contesta 529/503 sin tocar Qwen durante - * CHAT_BREAKER_SECONDS; pasado ese tiempo sale una sonda, y la primera respuesta con - * `choices` lo cierra. + * 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 Qwen le contesta, cierra; si la desafia, reabre. */ -const CHAT_CHALLENGE_MESSAGE = 'Qwen 上游繁忙,触发风控验证(被挤爆啦),请稍后重试 / Qwen chat challenge: upstream busy, retry later'; +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 CHAT_BREAKER_SECONDS = 60; -const chatBreaker = { strikes: 0, openUntil: 0 }; +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; - if (chatBreaker.strikes < CHAT_BREAKER_STRIKES) return CHAT_CHALLENGE_RETRY_AFTER_SECONDS; - chatBreaker.openUntil = Date.now() + CHAT_BREAKER_SECONDS * 1000; - return CHAT_BREAKER_SECONDS; + 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'); }; -const chatChallengeError = (retryAfter, details) => { - const error = new UpstreamResponseError(CHAT_CHALLENGE_MESSAGE, WAF_CHALLENGE_CODE, details); +const chatChallengeError = (message, retryAfter, details) => { + const error = new UpstreamResponseError(message, WAF_CHALLENGE_CODE, details); error.retryAfter = retryAfter; return error; }; -/** sendChatRequest lo llama antes de crear el chat o subir el historial. */ +/** + * 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. + */ const assertChatChallengeBreakerClosed = () => { - const remaining = Math.ceil((chatBreaker.openUntil - Date.now()) / 1000); - if (remaining > 0) throw chatChallengeError(remaining, { breakerOpen: true }); + if (!chatBreaker.openUntil) return; + 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; +}; + +/** + * 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; + const busy = findWafChallengeSignals(payload).some(item => CHAT_BUSY_SIGNAL_RE.test(item)); + 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 chatChallengeError(noteChatChallenge(), wafChallenge.details); + const challenge = chatChallengeFrom(payload); + if (challenge) throw challenge; if (!payload || typeof payload !== 'object') return; const explicitError = payload.error; @@ -353,8 +406,7 @@ const assertNoUpstreamFailure = (payload) => { ); } - // Qwen volvio a generar: la sonda paso. - if (Array.isArray(payload.choices)) resetChatChallengeBreaker(); + if (Array.isArray(payload.choices)) noteChatAnswer(); }; module.exports = { @@ -366,7 +418,9 @@ module.exports = { isWafChallengeError, isTransportInterruption, assertChatChallengeBreakerClosed, + chatChallengeFrom, resetChatChallengeBreaker, + setChatChallengeClockForTests, rateLimitRetryAfterSeconds, describeUpstreamFailure, noteRateLimitedAccount, diff --git a/tests/agent-account-failover.test.js b/tests/agent-account-failover.test.js index f8a7f2a2..98d634ac 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() }) @@ -127,14 +129,18 @@ 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, 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.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)) { @@ -144,14 +150,18 @@ test('WAF can switch once but never fans out through the entire account pool, an 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 () => { diff --git a/tests/agent-protocol.test.js b/tests/agent-protocol.test.js index a1a1372b..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: {}, @@ -1094,7 +1097,6 @@ test('Qwen HTTP-200 WAF payload is surfaced as an explicit failure', () => { }) test('Qwen HTTP-200 bare JSON WAF response reaches OpenAI clients explicitly', async () => { - require('../src/utils/upstream-error.js').resetChatChallengeBreaker() const res = createMockResponse() await handleStreamResponse( res, @@ -1110,7 +1112,7 @@ test('Qwen HTTP-200 bare JSON WAF response reaches OpenAI clients explicitly', a // 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_waf_challenge/) + assert.match(res.output, /upstream_unavailable/) assert.match(res.output, /chat challenge/) }) diff --git a/tests/chat-challenge.test.js b/tests/chat-challenge.test.js index 49f7b114..35ef5dd8 100644 --- a/tests/chat-challenge.test.js +++ b/tests/chat-challenge.test.js @@ -1,5 +1,6 @@ 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' @@ -10,71 +11,131 @@ process.env.ENABLE_FILE_LOG = 'false' process.env.PROXY_URL = '' const accountManager = require('../src/utils/account') -const { sendChatRequest } = require('../src/utils/request') +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: Readable.from(upstreamFrames), + currentAccount: null, + contextPrefixReused: true +}) +requestModule.invalidateContextPrefix = key => { invalidatedPrefixes.push(key) } +const { handleAnthropicMessages } = require('../src/controllers/anthropic') +const { handleNonStreamResponse } = require('../src/controllers/chat') +const { parseUpstreamImageError } = 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 challenge = { ret: ['FAIL_SYS_USER_VALIDATE', 'RGV587_ERROR::SM::哎哟喂,被挤爆啦,请稍后重试'] } +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 thrown = fn => { try { fn() } catch (error) { return error } return 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)) } + +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 }) -test.beforeEach(() => resetChatChallengeBreaker()) -test.after(() => { accountManager.destroy() }) +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 = thrown(() => assertNoUpstreamFailure(challenge)) + 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('three chat challenges in a row stop new requests; a request that gets an answer closes the breaker', () => { - thrown(() => assertNoUpstreamFailure(challenge)) - thrown(() => assertNoUpstreamFailure(challenge)) - assert.equal(thrown(assertChatChallengeBreakerClosed), null, 'two strikes keep it closed') - assert.equal(thrown(() => assertNoUpstreamFailure(challenge)).retryAfter, 60, 'the third asks for the full cooldown') +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('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, 60, 'the third asks for the full cooldown') - const blocked = thrown(assertChatChallengeBreakerClosed) + 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: 60 }) +}) +test('an answer from a stream already in flight does not close an open breaker', () => { + strike(3) assertNoUpstreamFailure(answer) - assert.equal(thrown(assertChatChallengeBreakerClosed), null) + now += 30_000 + assert.equal(caught(assertChatChallengeBreakerClosed).retryAfter, 30, 'Retry-After is the time actually left') }) -test('an open breaker lets requests through again after its cooldown, and the first challenged one reopens it', () => { - const realNow = Date.now - let now = realNow() - Date.now = () => now - try { - for (let strike = 0; strike < 3; strike += 1) thrown(() => assertNoUpstreamFailure(challenge)) - now += 30_000 - assert.equal(thrown(assertChatChallengeBreakerClosed).retryAfter, 30, 'Retry-After is the time actually left') - now += 30_000 - assert.equal(thrown(assertChatChallengeBreakerClosed), null, 'cooldown over: the probe goes out') - thrown(() => assertNoUpstreamFailure(challenge)) - assert.equal(thrown(assertChatChallengeBreakerClosed).retryAfter, 60) - } finally { - Date.now = realNow - } +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, 60, '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, 60) + assert.equal(caught(assertChatChallengeBreakerClosed).retryAfter, 60) }) 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 { - for (let strike = 0; strike < 3; strike += 1) thrown(() => assertNoUpstreamFailure(challenge)) + strike(3) await assert.rejects( sendChatRequest({ model: 'qwen3-max', messages: [{ role: 'user', content: '你好' }] }, { currentAccount: { email: 'breaker@example.invalid', token: 'breaker-test-token' } }), @@ -84,3 +145,40 @@ test('sendChatRequest refuses while the breaker is open, before creating a chat axios.post = realPost } }) + +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('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 + }) +}) From 91e02d2937bc65d325bba1b7253a0dbcaff11559 Mon Sep 17 00:00:00 2001 From: PEDRO LOBATO CARCAMO Date: Sat, 26 Sep 2026 13:48:12 -0600 Subject: [PATCH 3/7] fix(stream): commit streams on Qwen's first valid frame, so a chat challenge is a real 529/503 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit /v1/messages and the OpenAI agent stream wrote HTTP 200 plus message_start (or the role chunk) before reading anything from Qwen. So the breaker's probe request, the one that actually reaches Qwen, got its chat challenge as an in-stream error event. Claude Code shows that as "API error · Retrying in 5s" and ignores the retry_after inside it. Pre-commit refusals already produced a real 529 that it honors (~60 s spacing, observed live). Both paths now commit lazily. The SSE headers and first protocol event are written on the first upstream frame that passed assertNoUpstreamFailure (usually response.created, so thinking turns still start at once), on the first ping or heartbeat (the cap), or when an attempt ends. A chat challenge or quota on the first frame therefore arrives with nothing committed and goes out as a real 529/503/429 with Retry-After. Any other failure before the first frame commits and goes out in the stream exactly as before, so it does not become a 5xx that SDKs retry against upload/parse. There is still one reader of the stream, so no bytes are buffered or replayed and a challenge costs one breaker strike. The long silence the early pings fixed was the correction retry, which still runs committed. - anthropic.js: ensureMessageStart() on the first onDelta, before every ping (runWithAnthropicPing beforePing), after each attempt, and in the attempt catch for other failures. - chat.js: commitStream() through a new on_upstream_frame runtime hook, the SSE heartbeat (beforeBeat), writeDelta, and after the runtime returns. - openai-agent-runtime.js: a challenge switch whose replay cannot start rethrows the challenge instead of an opaque 502. - chat.image.video.js: a challenge on a stream request that is still uncommitted is a 503 as JSON (t2v presets text/event-stream). - upstream-error.js: Retry-After for a chat challenge is capped at 59 s, because the Anthropic and OpenAI SDKs ignore values of 60 or more. The reachability matrix is rewritten. - Tests: first-frame quota/challenge now pins a real status. Mid-stream variants keep the in-event retry_after coverage. New tests cover the commit points (first frame, ping/heartbeat cap, empty attempt, other failures still in-stream, t2v/t2i) and the agent challenge switch. tools/bun-smoke.js checks 529/503 + JSON + Retry-After over real Bun HTTP. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/controllers/anthropic.js | 95 +++++++++------ src/controllers/chat.image.video.js | 7 ++ src/controllers/chat.js | 47 +++++--- src/utils/openai-agent-runtime.js | 5 + src/utils/upstream-error.js | 48 +++++--- tests/agent-account-failover.test.js | 32 +++++ tests/chat-challenge.test.js | 171 +++++++++++++++++++++++++-- tests/upstream-quota-429.test.js | 114 +++++++++++------- tools/bun-smoke.js | 36 ++++++ 9 files changed, 430 insertions(+), 125 deletions(-) diff --git a/src/controllers/anthropic.js b/src/controllers/anthropic.js index 48b88077..2a4100ba 100644 --- a/src/controllers/anthropic.js +++ b/src/controllers/anthropic.js @@ -1189,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 (_) { @@ -1271,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; @@ -1755,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; @@ -1774,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: @@ -1782,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; @@ -1792,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] : [] @@ -1800,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; // 本轮收尾。解析器的尾巴属于这一轮,必须在判定之前放出来。文本通道截断之后例外: @@ -1956,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) { diff --git a/src/controllers/chat.image.video.js b/src/controllers/chat.image.video.js index 68c4a7a9..f920abe0 100644 --- a/src/controllers/chat.image.video.js +++ b/src/controllers/chat.image.video.js @@ -1535,6 +1535,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) } diff --git a/src/controllers/chat.js b/src/controllers/chat.js index 4512eba8..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 diff --git a/src/utils/openai-agent-runtime.js b/src/utils/openai-agent-runtime.js index c616734f..a228f407 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) { @@ -824,6 +826,9 @@ const runOpenAIAgentTurn = async (initialResponse, options = {}) => { agentRetry: true }) if (!retryResponse?.status || !retryResponse.response) { + // A challenge switch that could not even start keeps the challenge: it is still a + // retryable 503 with Retry-After, not an opaque 502. + if (challengeFailure) throw error return { ok: false, error: { status: 502, message: retryResponse?.message || 'Account failover request failed', code: 'upstream_retry_failed' }, diff --git a/src/utils/upstream-error.js b/src/utils/upstream-error.js index 00455763..1317e663 100644 --- a/src/utils/upstream-error.js +++ b/src/utils/upstream-error.js @@ -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. */ /** @@ -340,9 +345,14 @@ const noteChatAnswer = () => { 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 = retryAfter; + error.retryAfter = Math.min(retryAfter, CHAT_CHALLENGE_MAX_RETRY_AFTER_SECONDS); return error; }; diff --git a/tests/agent-account-failover.test.js b/tests/agent-account-failover.test.js index 98d634ac..c2b3b0ff 100644 --- a/tests/agent-account-failover.test.js +++ b/tests/agent-account-failover.test.js @@ -230,6 +230,38 @@ 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 challenge switch that cannot start keeps the challenge 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) +}) + 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/chat-challenge.test.js b/tests/chat-challenge.test.js index 35ef5dd8..acae83a2 100644 --- a/tests/chat-challenge.test.js +++ b/tests/chat-challenge.test.js @@ -30,14 +30,14 @@ let upstreamFrames = [] const invalidatedPrefixes = [] requestModule.sendChatRequest = async () => ({ status: true, - response: Readable.from(upstreamFrames), + response: typeof upstreamFrames === 'function' ? upstreamFrames() : Readable.from(upstreamFrames), currentAccount: null, contextPrefixReused: true }) requestModule.invalidateContextPrefix = key => { invalidatedPrefixes.push(key) } const { handleAnthropicMessages } = require('../src/controllers/anthropic') -const { handleNonStreamResponse } = require('../src/controllers/chat') -const { parseUpstreamImageError } = require('../src/controllers/chat.image.video') +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::哎哟喂,被挤爆啦,请稍后重试'] } @@ -47,6 +47,26 @@ const answer = { choices: [{ delta: { phase: 'answer', content: '你好' }, fini 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 = () => ({ @@ -97,13 +117,13 @@ test('a slider captcha is still a chat challenge, but is not reported as a 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, 60, 'the third asks for the full cooldown') + 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: 60 }) + { rateLimited: false, overloaded: true, status: 529, retryAfter: 59 }) }) test('an answer from a stream already in flight does not close an open breaker', () => { @@ -117,7 +137,7 @@ test('after the cooldown exactly one probe goes out, and its answer closes the b strike(3) now += 60_000 assert.equal(caught(assertChatChallengeBreakerClosed), null, 'the first request is the probe') - assert.equal(caught(assertChatChallengeBreakerClosed).retryAfter, 60, 'the others wait for 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') @@ -127,8 +147,8 @@ test('a challenged probe reopens the breaker for a full cooldown', () => { strike(3) now += 60_000 assertChatChallengeBreakerClosed() - assert.equal(caught(() => assertNoUpstreamFailure(busy)).retryAfter, 60) - assert.equal(caught(assertChatChallengeBreakerClosed).retryAfter, 60) + 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 () => { @@ -174,6 +194,141 @@ test('/v1/messages keeps a reused history prefix on a chat challenge, and still 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, 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..2ddec888 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,11 @@ async function main() { } else if (pathname === '/api/v2/chat/completions') { completionRequests += 1 response.setHeader('Content-Type', 'text/event-stream') + 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 +335,36 @@ 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') + 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 }, From e3b0c1c679858e5b03f09596cff14bc2dd57541d Mon Sep 17 00:00:00 2001 From: PEDRO LOBATO CARCAMO Date: Sat, 26 Sep 2026 13:51:54 -0600 Subject: [PATCH 4/7] fix(agent): a quota account switch that cannot start keeps the quota error The challenge case already rethrew; quota fell through to an opaque 502 upstream_retry_failed. Rethrowing it lets the agent stream answer a real 429 with Retry-After when nothing was committed yet, like a first-frame quota. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/utils/openai-agent-runtime.js | 6 +++--- tests/agent-account-failover.test.js | 6 +++++- 2 files changed, 8 insertions(+), 4 deletions(-) diff --git a/src/utils/openai-agent-runtime.js b/src/utils/openai-agent-runtime.js index a228f407..62709e4e 100644 --- a/src/utils/openai-agent-runtime.js +++ b/src/utils/openai-agent-runtime.js @@ -826,9 +826,9 @@ const runOpenAIAgentTurn = async (initialResponse, options = {}) => { agentRetry: true }) if (!retryResponse?.status || !retryResponse.response) { - // A challenge switch that could not even start keeps the challenge: it is still a - // retryable 503 with Retry-After, not an opaque 502. - if (challengeFailure) throw error + // 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/tests/agent-account-failover.test.js b/tests/agent-account-failover.test.js index c2b3b0ff..30e7db72 100644 --- a/tests/agent-account-failover.test.js +++ b/tests/agent-account-failover.test.js @@ -256,10 +256,14 @@ test('agent stream: challenged on both accounts is a real 503 with Retry-After, assert.equal(JSON.parse(response.output).error.code, 'upstream_unavailable') }) -test('a challenge switch that cannot start keeps the challenge instead of an opaque 502', async () => { +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 => { From be703668118a2f24e1de0f366f69ae3de1465f1d Mon Sep 17 00:00:00 2001 From: PEDRO LOBATO CARCAMO Date: Sat, 26 Sep 2026 19:56:45 -0600 Subject: [PATCH 5/7] fix(upstream): a captcha page served as text/html is a chat challenge too Qwen also answers /api/v2/chat/completions with 200 text/html: Aliyun's captcha page (seen 2026-09-19, the HTML form #179 fixed on the image path). The text path never looked at the body: with no `data:` frame no detector ran, the breaker counted nothing, and every surface retried 2-3 times against the WAF before reporting an empty answer (502 upstream_empty_output / an in-stream api_error) with no Retry-After. sendChatRequest now reads a text/html 200 (capped at 256 KiB) and runs it through the shared detector: the captcha page is one strike and the same 529/503 + Retry-After as the JSON challenge, thrown before any stream is handed on. Any other HTML body is handed on byte for byte, so its path is unchanged. Covered in chat-challenge.test.js (captcha page -> 529 + one strike; other HTML -> replayed, no strike; fails with the screen removed) and over real Bun HTTP in bun-smoke (one send, not a retry loop). Co-Authored-By: Claude Opus 5.5 (1M context) --- src/utils/request.js | 30 ++++++++++++++++++++++++-- tests/chat-challenge.test.js | 42 ++++++++++++++++++++++++++++++++++++ tools/bun-smoke.js | 16 ++++++++++++++ 3 files changed, 86 insertions(+), 2 deletions(-) diff --git a/src/utils/request.js b/src/utils/request.js index de9c0559..f44f9e7d 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, assertChatChallengeBreakerClosed } = require('./upstream-error.js') +const { ContextExternalizationError, isTransportInterruption, assertChatChallengeBreakerClosed, chatChallengeFrom, isWafChallengeError } = 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,29 @@ 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 + for await (const chunk of response.data) { + chunks.push(Buffer.from(chunk)) + size += chunk.length + if (size >= HTML_BODY_MAX_BYTES) break + } + const body = Buffer.concat(chunks).toString('utf8') + const challenge = chatChallengeFrom(body) + if (challenge) throw challenge + 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. @@ -768,7 +792,7 @@ const sendChatRequest = async (body, options = {}) => { contextPrefixReused: contextResult.reusedPrefix === true, contextSerializedBytes: contextResult.serializedBytes, status: true, - response: response.data + response: await screenHtmlChallenge(response) } } // 非 200 但是没抛——退出循环, 走下面错误分类 @@ -776,6 +800,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/tests/chat-challenge.test.js b/tests/chat-challenge.test.js index acae83a2..a008ccba 100644 --- a/tests/chat-challenge.test.js +++ b/tests/chat-challenge.test.js @@ -166,6 +166,48 @@ test('sendChatRequest refuses while the breaker is open, before creating a chat } }) +// 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('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: [] }, {}) diff --git a/tools/bun-smoke.js b/tools/bun-smoke.js index 2ddec888..4f833490 100644 --- a/tools/bun-smoke.js +++ b/tools/bun-smoke.js @@ -102,6 +102,12 @@ 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`) @@ -363,6 +369,16 @@ async function main() { // 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) { From d594e22a00fd7c219ea57f7cbf144a6cbb1d95d7 Mon Sep 17 00:00:00 2001 From: PEDRO LOBATO CARCAMO Date: Sat, 26 Sep 2026 20:14:02 -0600 Subject: [PATCH 6/7] fix(upstream): chat challenge on images/video and a cut captcha page, per review MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review of the rebase onto #178/#179 found, besides the new HTML screen: - t2v asks for responseType 'json', so a non-JSON body (the captcha page, or a `data:`-wrapped challenge frame) reaches the controller as a string that parseUpstreamImageError dropped on its failed JSON.parse. The captcha page's first URL came back as the "video" with a 200, no strike, no Retry-After. On origin/main too; #179's raw-text detection only ran on the streaming branch. String bodies now go through the same raw-text parser (parseUpstreamBody) on both the inline check and the reader. - A t2v challenge with a non-2xx status was parsed, and counted, twice (inner catch, then outer catch). The inner catch now throws it. - Image/video never reached noteChatAnswer: an answered image request could take the half-open probe and keep text blocked for another full window, and image answers did not break a run of strikes. Both answer paths now call it; a request without a prompt is rejected before it can take the probe. - A connection cut while reading a text/html body was handled as a transport failure of the POST (re-POST, account penalty, prefix invalidation) even when the captcha markers had already arrived. The bytes read are screened first. - The busy-vs-captcha wording looked only at the detector's filtered signals, so "被挤爆" without "RGV587" read as a captcha. It looks at the whole payload again, as before the rebase. - The breaker comment claimed only the probe's own answer closes a half-open breaker; any answer does. Documented, with its ceiling. Each fix has a test in chat-challenge.test.js that fails with the fix removed. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/controllers/chat.image.video.js | 40 ++++++++----- src/utils/request.js | 15 +++-- src/utils/upstream-error.js | 10 +++- tests/chat-challenge.test.js | 88 +++++++++++++++++++++++++++++ 4 files changed, 134 insertions(+), 19 deletions(-) diff --git a/src/controllers/chat.image.video.js b/src/controllers/chat.image.video.js index f920abe0..65f71679 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 { assertChatChallengeBreakerClosed, chatChallengeFrom } = require('../utils/upstream-error.js') +const { assertChatChallengeBreakerClosed, chatChallengeFrom, noteChatAnswer } = require('../utils/upstream-error.js') const DATA_URI_REGEX = /^data:(.+);base64,(.*)$/i const HTTP_URL_REGEX = /^https?:\/\//i @@ -188,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 - 任意负载 @@ -876,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, @@ -975,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) : '' @@ -1295,6 +1304,15 @@ const generateImageVideoResult = async (payload) => { ] } + // Antes del breaker: una peticion invalida no debe quedarse con la unica sonda. + const userPrompt = messages?.[messages.length - 1]?.content + if (!userPrompt) { + throw { + status: 400, + error: '缺少有效的提示词' + } + } + try { assertChatChallengeBreakerClosed() } catch (challenge) { @@ -1309,14 +1327,6 @@ const generateImageVideoResult = async (payload) => { reqBody.chat_id = chatID - const userPrompt = messages?.[messages.length - 1]?.content - if (!userPrompt) { - throw { - status: 400, - error: '缺少有效的提示词' - } - } - const messagesHistory = messages.filter(item => item.role === 'user' || item.role === 'assistant') const selectedImageList = [] @@ -1439,7 +1449,7 @@ 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)) { @@ -1451,7 +1461,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) @@ -1464,6 +1476,7 @@ const generateImageVideoResult = async (payload) => { if (newChatType === 't2i' || newChatType === 'image_edit') { const contentUrl = await resolveImageResultContentUrl(responseData.data, chatID, token) + noteChatAnswer() return { model, chatType: newChatType, @@ -1474,6 +1487,7 @@ const generateImageVideoResult = async (payload) => { if (newChatType === 't2v') { const contentUrl = await resolveVideoResultContentUrl(responseData.data, token, chatID) + noteChatAnswer() return { model, chatType: newChatType, diff --git a/src/utils/request.js b/src/utils/request.js index f44f9e7d..21f6882a 100644 --- a/src/utils/request.js +++ b/src/utils/request.js @@ -40,14 +40,21 @@ const screenHtmlChallenge = async (response) => { if (!/text\/html/i.test(String(response.headers?.['content-type'] || ''))) return response.data const chunks = [] let size = 0 - for await (const chunk of response.data) { - chunks.push(Buffer.from(chunk)) - size += chunk.length - if (size >= HTML_BODY_MAX_BYTES) break + 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]) } diff --git a/src/utils/upstream-error.js b/src/utils/upstream-error.js index 1317e663..e13963e3 100644 --- a/src/utils/upstream-error.js +++ b/src/utils/upstream-error.js @@ -299,7 +299,11 @@ const noteRateLimitedAccount = (error, account) => { * 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 Qwen le contesta, cierra; si la desafia, reabre. + * 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'; @@ -379,7 +383,8 @@ const assertChatChallengeBreakerClosed = () => { const chatChallengeFrom = (payload) => { const detected = detectWafChallenge(payload); if (!detected) return null; - const busy = findWafChallengeSignals(payload).some(item => CHAT_BUSY_SIGNAL_RE.test(item)); + // 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); }; @@ -429,6 +434,7 @@ module.exports = { isTransportInterruption, assertChatChallengeBreakerClosed, chatChallengeFrom, + noteChatAnswer, resetChatChallengeBreaker, setChatChallengeClockForTests, rateLimitRetryAfterSeconds, diff --git a/tests/chat-challenge.test.js b/tests/chat-challenge.test.js index a008ccba..11a60f19 100644 --- a/tests/chat-challenge.test.js +++ b/tests/chat-challenge.test.js @@ -35,6 +35,7 @@ requestModule.sendChatRequest = async () => ({ 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') @@ -114,6 +115,11 @@ test('a slider captcha is still a chat challenge, but is not reported as a busy 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') @@ -208,6 +214,29 @@ test('any other text/html body is handed on byte for byte, and counts no strike' } }) +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: [] }, {}) @@ -379,3 +408,62 @@ test('image and video generation map a chat challenge to a retryable 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') +}) From 452911c36ea1bc6260b66e13eb1a608df936cfd7 Mon Sep 17 00:00:00 2001 From: PEDRO LOBATO CARCAMO Date: Sat, 26 Sep 2026 20:29:49 -0600 Subject: [PATCH 7/7] fix(upstream): a half-open probe that never asks Qwen hands itself on Second review pass. The half-open probe re-arms the window for everyone else; a probe that failed before Qwen answered or challenged it (no account, no chat_id, upload failure, an image_edit without a text part) left it probing, and every request got 529/503 for another full window. - assertChatChallengeBreakerClosed says whether the caller is the probe; sendChatRequest (now a thin guard around postChatRequest) and the image/video flow release it (releaseChatProbe) on any failure that is not a chat challenge, so the next request probes at once. A probe that reaches Qwen and ends with neither an answer nor a challenge still costs one window; documented. - t2v closes the half-open breaker as soon as Qwen accepts the task (its whole body was already screened), instead of after the video is polled to completion, minutes later. Left out, never seen from Qwen (every challenge measured arrived with HTTP 200, whole): a non-2xx challenge on the streamed image path, and a captcha page cut mid-body on it. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/controllers/chat.image.video.js | 10 +++++--- src/utils/request.js | 15 +++++++++-- src/utils/upstream-error.js | 17 +++++++++++- tests/chat-challenge.test.js | 40 +++++++++++++++++++++++++++++ tests/expected-counts.json | 2 +- 5 files changed, 77 insertions(+), 7 deletions(-) diff --git a/src/controllers/chat.image.video.js b/src/controllers/chat.image.video.js index 65f71679..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 { assertChatChallengeBreakerClosed, chatChallengeFrom, noteChatAnswer } = 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 @@ -1283,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 = { @@ -1314,7 +1315,7 @@ const generateImageVideoResult = async (payload) => { } try { - assertChatChallengeBreakerClosed() + probe = assertChatChallengeBreakerClosed() } catch (challenge) { throw imageChallengeError(challenge) } @@ -1486,8 +1487,10 @@ const generateImageVideoResult = async (payload) => { } if (newChatType === 't2v') { - const contentUrl = await resolveVideoResultContentUrl(responseData.data, token, chatID) + // 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, chatType: newChatType, @@ -1499,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 diff --git a/src/utils/request.js b/src/utils/request.js index 21f6882a..9b2f926f 100644 --- a/src/utils/request.js +++ b/src/utils/request.js @@ -8,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, assertChatChallengeBreakerClosed, chatChallengeFrom, isWafChallengeError } = 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 @@ -674,7 +674,18 @@ const externalizeOversizedAgentContext = async ( */ const sendChatRequest = async (body, options = {}) => { // Qwen esta rechazando la generacion: ni chat nuevo ni upload, solo 529/503 al cliente. - assertChatChallengeBreakerClosed() + 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)——轮换器跳过它们, // 即使它们对其他请求仍然可用。 diff --git a/src/utils/upstream-error.js b/src/utils/upstream-error.js index e13963e3..ee253b73 100644 --- a/src/utils/upstream-error.js +++ b/src/utils/upstream-error.js @@ -363,14 +363,28 @@ const chatChallengeError = (message, retryAfter, details) => { /** * 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; + 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(); }; /** @@ -435,6 +449,7 @@ module.exports = { assertChatChallengeBreakerClosed, chatChallengeFrom, noteChatAnswer, + releaseChatProbe, resetChatChallengeBreaker, setChatChallengeClockForTests, rateLimitRetryAfterSeconds, diff --git a/tests/chat-challenge.test.js b/tests/chat-challenge.test.js index 11a60f19..a1159a52 100644 --- a/tests/chat-challenge.test.js +++ b/tests/chat-challenge.test.js @@ -467,3 +467,43 @@ test('an image request without a prompt does not take the single half-open probe 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"