From 421ff27eb623b4d64b1cc008e67fcbdb7c4443c2 Mon Sep 17 00:00:00 2001 From: can1357 Date: Mon, 8 Jun 2026 14:44:13 +0200 Subject: [PATCH] fix(ai): scoped Anthropic fallback state and fixed iterator early-close cleanup - Scoped Anthropic provider session state by base URL and model so strict-tools/fast-mode fallback state no longer leaks across unrelated endpoints or models, and fast-mode re-arming clears all matching scoped entries. - Dropped stale strict fallback error messages after successful retries and only enabled adaptive thinking display for models that advertise support, avoiding unsupported-model 400 errors. - Updated idle stream iteration to close the upstream iterator on non-continuing exits (including consumer early break) and added coverage for upstream closure behavior. --- packages/ai/src/providers/anthropic.ts | 51 +++++++++++++++---- packages/ai/src/utils/idle-iterator.ts | 11 ++++ .../ai/test/stream-timeout-defaults.test.ts | 29 +++++++++++ 3 files changed, 80 insertions(+), 11 deletions(-) diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index 652a854cd..a78b49bd3 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -314,16 +314,29 @@ function createAnthropicProviderSessionState(): AnthropicProviderSessionState { return state; } +/** + * Key the sticky strict-tools / fast-mode learning per endpoint+model. A + * grammar-too-large 400 or a fast-mode rejection is specific to the model (its + * tool grammar / entitlement) and the endpoint (direct Anthropic vs a gateway / + * Foundry / Bedrock proxy), so it MUST NOT bleed onto unrelated anthropic-messages + * requests in the same session. NUL separates the two components so neither can + * forge the boundary. + */ +function anthropicProviderSessionStateKey(baseUrl: string, modelId: string): string { + return `${ANTHROPIC_PROVIDER_SESSION_STATE_KEY}:${baseUrl}\u0000${modelId}`; +} + function getAnthropicProviderSessionState( providerSessionState: Map | undefined, + baseUrl: string, + modelId: string, ): AnthropicProviderSessionState | undefined { if (!providerSessionState) return undefined; - const existing = providerSessionState.get(ANTHROPIC_PROVIDER_SESSION_STATE_KEY) as - | AnthropicProviderSessionState - | undefined; + const key = anthropicProviderSessionStateKey(baseUrl, modelId); + const existing = providerSessionState.get(key) as AnthropicProviderSessionState | undefined; if (existing) return existing; const created = createAnthropicProviderSessionState(); - providerSessionState.set(ANTHROPIC_PROVIDER_SESSION_STATE_KEY, created); + providerSessionState.set(key, created); return created; } @@ -338,10 +351,14 @@ export function clearAnthropicFastModeFallback( providerSessionState: Map | undefined, ): void { if (!providerSessionState) return; - const state = providerSessionState.get(ANTHROPIC_PROVIDER_SESSION_STATE_KEY) as - | AnthropicProviderSessionState - | undefined; - if (state) state.fastModeDisabled = false; + // Fast mode is re-armed session-wide (user toggled `/fast on`), so clear the + // sticky flag on every per-endpoint/model Anthropic entry — plus the legacy + // unscoped key — rather than a single shared object. + const prefix = `${ANTHROPIC_PROVIDER_SESSION_STATE_KEY}:`; + for (const [key, value] of providerSessionState) { + if (key !== ANTHROPIC_PROVIDER_SESSION_STATE_KEY && !key.startsWith(prefix)) continue; + (value as AnthropicProviderSessionState).fastModeDisabled = false; + } } function isAnthropicStrictGrammarTooLargeError(error: unknown): boolean { @@ -1434,7 +1451,7 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( const baseUrl = resolveAnthropicBaseUrl(model, options?.apiKey ?? getEnvApiKey(model.provider) ?? "") ?? "https://api.anthropic.com"; - const providerSessionState = getAnthropicProviderSessionState(options?.providerSessionState); + const providerSessionState = getAnthropicProviderSessionState(options?.providerSessionState, baseUrl, model.id); let disableStrictTools = (providerSessionState?.strictToolsDisabled ?? false) || (model.compat?.disableStrictTools ?? false); let strictFallbackErrorMessage: string | undefined; @@ -1899,6 +1916,15 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( } } + // A strict-tools / grammar-too-large fallback (or a transient retry that ran + // after one) stamps `strictFallbackErrorMessage` onto `output.errorMessage` for + // the *discarded* attempt. Reaching here means a later attempt succeeded + // (the in-loop guard throws on error/aborted before this break), so drop the + // stale 400 text — otherwise a successful turn ships a phantom error string to + // telemetry/UI. Only the exact stale value is cleared, never a real refusal message. + if (output.errorMessage !== undefined && output.errorMessage === strictFallbackErrorMessage) { + output.errorMessage = undefined; + } output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; if (dropFastMode && resolveServiceTier(options?.serviceTier, model.provider) === "priority") { @@ -2466,8 +2492,11 @@ function buildParams( const adaptive: { type: "adaptive"; display?: AnthropicThinkingDisplay } = { type: "adaptive" }; // Starting with Claude Opus 4.7, adaptive thinking content is omitted from the // response by default. Opt into summarized reasoning so thinking deltas keep - // streaming with human-readable content for callers that rely on it. - if (options.thinkingDisplay !== undefined || supportsAdaptiveThinkingDisplay(model.id)) { + // streaming with human-readable content for callers that rely on it. The + // `display` field is gated strictly on model support: Opus 4.6 / Sonnet 4.6+ + // reject it with a 400, so an explicit `thinkingDisplay` MUST NOT force it onto + // a model that can't accept it (a hidden-thinking toggle must never break the request). + if (supportsAdaptiveThinkingDisplay(model.id)) { adaptive.display = options.thinkingDisplay ?? "summarized"; } thinking = adaptive; diff --git a/packages/ai/src/utils/idle-iterator.ts b/packages/ai/src/utils/idle-iterator.ts index d19b5cd57..f1c5345b0 100644 --- a/packages/ai/src/utils/idle-iterator.ts +++ b/packages/ai/src/utils/idle-iterator.ts @@ -124,8 +124,11 @@ export async function* iterateWithIdleTimeout( firstItemTimeoutMs !== undefined && firstItemTimeoutMs > 0 ? Date.now() + firstItemTimeoutMs : undefined; const abortSignal = options.abortSignal; const iterator = iterable[Symbol.asyncIterator](); + let iteratorClosed = false; const closeIterator = (): void => { + if (iteratorClosed) return; + iteratorClosed = true; const returnPromise = iterator.return?.(); if (returnPromise) { void returnPromise.catch(() => {}); @@ -212,6 +215,12 @@ export async function* iterateWithIdleTimeout( racers.push(promise); } + // Tracks whether this iteration handed an item to the consumer and resumed + // normally. Any other exit — internal throw, `done` return, or the consumer + // abandoning us via `.return()`/`.throw()` at the `yield` below — must close + // the upstream iterator so the underlying SSE body / SDK stream (and its + // socket) is released instead of being left suspended. + let continuing = false; try { const outcome = await Promise.race(racers); if (outcome.kind === "abort") { @@ -247,7 +256,9 @@ export async function* iterateWithIdleTimeout( lastProgressAt = Date.now(); } yield item; + continuing = true; } finally { + if (!continuing) closeIterator(); if (timer !== undefined) clearTimeout(timer); // Resolve dangling promises so the racers don't leak (Promise.race is one-shot). resolveTimeout?.({ kind: "timeout" }); diff --git a/packages/ai/test/stream-timeout-defaults.test.ts b/packages/ai/test/stream-timeout-defaults.test.ts index 5920ef978..ee22a94e4 100644 --- a/packages/ai/test/stream-timeout-defaults.test.ts +++ b/packages/ai/test/stream-timeout-defaults.test.ts @@ -227,4 +227,33 @@ describe("iterateWithIdleTimeout", () => { await Bun.sleep(20); expect(firstItemTimedOut).toBe(false); }); + + it("closes the upstream iterator when the consumer breaks early", async () => { + let upstreamClosed = false; + async function* source(): AsyncGenerator { + try { + let n = 0; + while (true) { + await Bun.sleep(1); + yield `item-${n++}`; + } + } finally { + // Runs only if the wrapper forwards `.return()` to us. + upstreamClosed = true; + } + } + + for await (const _item of iterateWithIdleTimeout(source(), { + idleTimeoutMs: 1_000, + errorMessage: "idle timeout", + })) { + break; // abandon the wrapper after the first item + } + + // The wrapper must propagate the consumer's early termination to the source + // so the underlying SSE body / SDK stream (and its socket) is released + // instead of being left suspended. + await Bun.sleep(5); + expect(upstreamClosed).toBe(true); + }); });