From 2724ad07e1f5476b75aa0d1f7cf2674ee6d40cd9 Mon Sep 17 00:00:00 2001 From: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com> Date: Thu, 10 Sep 2026 18:39:42 -0300 Subject: [PATCH] fix(sse): restore upstream error parsing + rate-limit lock in the provider pipeline MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit #12867 lifted the non-2xx branch out of chatCore.ts into providerExecutionPipeline.ts::toOutcome, but reimplemented it instead of delegating, dropping three behaviors the non-streaming failure path relies on: 1. A non-JSON upstream body fell into the inline JSON.parse catch and surfaced as the (empty) statusText — "upstream error" — discarding the text the client needs. Now parseUpstreamError() supplies the message, as chatCore did. 2. The body-derived retry-after ("Please retry after 20s") was never parsed: retryAfterMs was hard-coded null into applyStatusRestatement/createErrorResult. 3. recordRateLimitBody was plumbed through PipelineStateHooks and wired at both chatCore call sites, but never called — so updateFromResponseBody(), which drains the runtime reservoir on a body-derived 429, silently stopped running for every request routed through this pipeline. restatement rules also now match against the real upstream payload rather than the request body that transformedBody carried whenever the parse failed. Regression guard: tests/unit/chat-rate-limit-body-lock.test.ts (2 base-red failures on the release/v3.8.51 tip) — both green; provider-execution-pipeline 13/13; typecheck:core clean. --- .../chatCore/providerExecutionPipeline.ts | 100 +++++++++++++----- 1 file changed, 73 insertions(+), 27 deletions(-) diff --git a/open-sse/handlers/chatCore/providerExecutionPipeline.ts b/open-sse/handlers/chatCore/providerExecutionPipeline.ts index 077ff7cfe1..bce99f10d8 100644 --- a/open-sse/handlers/chatCore/providerExecutionPipeline.ts +++ b/open-sse/handlers/chatCore/providerExecutionPipeline.ts @@ -3,11 +3,17 @@ import type { getProviderCredentials } from "@/sse/services/auth.ts"; import type { updateFromHeaders, updateFromResponseBody } from "../../services/rateLimitManager.ts"; import type { writeTerminalStatus } from "@/shared/utils/terminalStatus.ts"; import type { updateProviderConnection } from "@/lib/db/providers.ts"; -import type { lockModel, recordCoreOwnedAntigravityQuotaState } from "../../services/accountFallback.ts"; -import { createErrorResult } from "../../utils/error.ts"; +import type { + lockModel, + recordCoreOwnedAntigravityQuotaState, +} from "../../services/accountFallback.ts"; +import { createErrorResult, parseUpstreamError } from "../../utils/error.ts"; import { applyStatusRestatement } from "../../config/upstreamStatusRestatement.ts"; import { recoverAnthropicThinkingSignature } from "./thinkingSignatureRecovery.ts"; -import { isModelUnavailableError, getNextFamilyFallback as defaultGetNextFamilyFallback } from "../../services/modelFamilyFallback.ts"; +import { + isModelUnavailableError, + getNextFamilyFallback as defaultGetNextFamilyFallback, +} from "../../services/modelFamilyFallback.ts"; import { COOLDOWN_MS } from "../../config/errorConfig.ts"; import { normalizeHeaders } from "../../utils/headers.ts"; @@ -164,7 +170,8 @@ async function toOutcome( attempt: ChatCoreExecutorResult, model: string, connectionId: string, - provider: string + provider: string, + state: PipelineStateHooks ): Promise { const status = attempt.response.status; if (status >= 200 && status < 300) { @@ -178,29 +185,40 @@ async function toOutcome( connectionId, }; } - let message = attempt.response.statusText || "upstream error"; - let body: unknown = attempt.transformedBody; - try { - // clone() is the drain. sendProviderAttempt must not cancel() a streaming - // non-2xx body before we get here (BYOP 422 / Codex 429 Retry-After). - body = JSON.parse(await attempt.response.clone().text()); - const err = (body as { error?: { message?: unknown } } | null)?.error; - if (err && typeof err.message === "string" && err.message) message = err.message; - } catch { - // keep statusText + // Delegate to the canonical upstream-error parser instead of re-implementing it. + // #12867 lifted this branch out of chatCore.ts but replaced its parseUpstreamError() + // call with an inline JSON.parse, which silently dropped two behaviors the + // non-streaming failure path depends on (tests/unit/chat-rate-limit-body-lock.test.ts): + // 1. a non-JSON body fell into the catch and surfaced as the (empty) statusText — + // "upstream error" — discarding the upstream text the client needs to see; + // 2. the body-derived retry-after ("Please retry after 20s") was never parsed, so + // retryAfterMs stayed hard-coded null and the runtime limiter was never locked. + // clone() is still the drain: sendProviderAttempt must not cancel() a streaming + // non-2xx body before we get here (BYOP 422 / Codex 429 Retry-After), and cloning + // keeps attempt.response readable for the consumers we hand it back to below. + const details = await parseUpstreamError(attempt.response.clone(), provider); + const message = details.message || attempt.response.statusText || "upstream error"; + // #12867 plumbed recordRateLimitBody through PipelineStateHooks but never called it, so the + // body-derived rate-limit lock chatCore used to apply (updateFromResponseBody on the upstream + // error body — "Please retry after 20s" → reservoir 0) silently stopped running for every + // request routed through this pipeline. Feed the parsed upstream body back the way chatCore + // did. recordRateLimitHeaders stays chatCore's job: it already learns from the real response + // on the success path, and calling it here would re-learn from the same headers twice. + if (connectionId && details.responseBody !== null && details.responseBody !== undefined) { + state.recordRateLimitBody(provider, connectionId, details.responseBody, status, model); } + // responseBody is the parsed upstream payload (or { _rawText } for a non-JSON body), + // so restatement rules now match against the real upstream text rather than the + // request body that transformedBody carried whenever the parse failed. + const body: unknown = details.responseBody ?? attempt.transformedBody; const restatement = applyStatusRestatement({ provider, status, message, body, - retryAfterMs: null, + retryAfterMs: details.retryAfterMs, }); - const result = createErrorResult( - restatement.status, - message, - restatement.retryAfterMs - ); + const result = createErrorResult(restatement.status, message, restatement.retryAfterMs); return { kind: "error", result: { @@ -273,7 +291,13 @@ export async function runProviderExecutionPipeline( const status = attempt.response.status; if (status >= 200 && status < 300) { - return toOutcome(attempt, wire.currentModel, currentConnectionId(connection), target.provider); + return toOutcome( + attempt, + wire.currentModel, + currentConnectionId(connection), + target.provider, + state + ); } const isolateProbe = await state.isolateProbeFailures(); @@ -401,18 +425,24 @@ export async function runProviderExecutionPipeline( }; }, }); - if (signatureRecovery.attempted && signatureRecovery.succeeded && signatureRecovery.execution) { + if ( + signatureRecovery.attempted && + signatureRecovery.succeeded && + signatureRecovery.execution + ) { lastAttempt = { response: signatureRecovery.execution.response, url: signatureRecovery.execution.url ?? attempt.url, - headers: (signatureRecovery.execution.headers as Record) ?? attempt.headers, + headers: + (signatureRecovery.execution.headers as Record) ?? attempt.headers, transformedBody: signatureRecovery.execution.transformedBody ?? attempt.transformedBody, }; return toOutcome( lastAttempt, wire.currentModel, currentConnectionId(connection), - target.provider + target.provider, + state ); } } @@ -430,7 +460,11 @@ export async function runProviderExecutionPipeline( // keep statusText } if (isModelUnavailableError(status, fallbackMessage, target.provider)) { - const nextModel = resolveFamilyFallback(wire.currentModel, wire.triedModels, target.provider); + const nextModel = resolveFamilyFallback( + wire.currentModel, + wire.triedModels, + target.provider + ); if (nextModel) { wire.setBodyAndModel({ ...wire.body, model: nextModel }, nextModel); modelFallbackPending = true; @@ -439,11 +473,23 @@ export async function runProviderExecutionPipeline( } } - return toOutcome(attempt, wire.currentModel, currentConnectionId(connection), target.provider); + return toOutcome( + attempt, + wire.currentModel, + currentConnectionId(connection), + target.provider, + state + ); } if (lastAttempt) { - return toOutcome(lastAttempt, wire.currentModel, currentConnectionId(connection), target.provider); + return toOutcome( + lastAttempt, + wire.currentModel, + currentConnectionId(connection), + target.provider, + state + ); } return leaseMismatch(wire.currentModel, currentConnectionId(connection)); }