From c0042468d2bb5eb2156b0003881a9afd927134f1 Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 20:45:02 +0900 Subject: [PATCH 1/2] fix(anthropic): eject an OAuth account whose access token was revoked Account recovery handled only 429 and 403, so an exact revoked-token 401 was returned to the client while session affinity kept pinning the same account. Before any output, classify that exact envelope (bounded, fatal UTF-8, no error code), mark the sending account needsReauth under a credential-generation and send-ownership fence, clear its affinities, and allow an eligible same-pool account to take over within existing send limits. Covers native Messages (including pool-off admission order), translated dispatch and continuation, and the web-search and image sidecars. Other 401s and anything after output keep their existing behavior. Closes #6749 --- .../src/content/docs/fr/guides/claude-code.md | 2 + .../src/content/docs/guides/claude-code.md | 6 + .../src/content/docs/ja/guides/claude-code.md | 2 + .../src/content/docs/ko/guides/claude-code.md | 2 + .../src/content/docs/ru/guides/claude-code.md | 2 + .../src/content/docs/tr/guides/claude-code.md | 2 + .../content/docs/zh-cn/guides/claude-code.md | 2 + .../content/docs/zh-tw/guides/claude-code.md | 2 + scripts/test-layout/layout.json | 5 + src/images/loop.ts | 2 +- src/oauth/anthropic-account-refusal.ts | 40 ++- src/oauth/anthropic-routing.ts | 2 +- src/oauth/anthropic-send-ownership.ts | 6 +- src/oauth/store.ts | 2 +- src/server/messages-native-oauth.ts | 17 +- src/server/messages-native.ts | 6 +- src/server/responses/adapter-continuation.ts | 4 +- src/server/responses/adapter-dispatch.ts | 4 +- src/server/responses/sidecar-execution.ts | 4 +- src/web-search/loop.ts | 2 + structure/data-planes/images.md | 2 + structure/data-planes/inbound-compat.md | 2 + structure/providers-and-adapters.md | 2 +- structure/providers/anthropic-account-pool.md | 22 ++ structure/transports/inventory.md | 2 +- structure/transports/responses-failover.md | 2 +- ...anthropic-revoked-token-boundaries.test.ts | 111 +++++++++ ...thropic-revoked-token-continuation.test.ts | 127 ++++++++++ .../anthropic-revoked-token-sidecars.test.ts | 139 +++++++++++ .../anthropic/anthropic-revoked-token.test.ts | 234 ++++++++++++++++++ .../messages-revoked-token.test.ts | 215 ++++++++++++++++ tests/fixtures/test-layout-expected.json | 5 + 32 files changed, 950 insertions(+), 27 deletions(-) create mode 100644 tests/adapters/anthropic/anthropic-revoked-token-boundaries.test.ts create mode 100644 tests/adapters/anthropic/anthropic-revoked-token-continuation.test.ts create mode 100644 tests/adapters/anthropic/anthropic-revoked-token-sidecars.test.ts create mode 100644 tests/adapters/anthropic/anthropic-revoked-token.test.ts create mode 100644 tests/claude-integration/messages-revoked-token.test.ts diff --git a/docs-site/src/content/docs/fr/guides/claude-code.md b/docs-site/src/content/docs/fr/guides/claude-code.md index f1fd2649697..22f57ac1954 100644 --- a/docs-site/src/content/docs/fr/guides/claude-code.md +++ b/docs-site/src/content/docs/fr/guides/claude-code.md @@ -775,3 +775,5 @@ Claude Code **2.1.257 or newer** is required for FORCE. Plugin and built-in agen The dashboard warns about old or unknown CLI versions, unavailable targets, and either variable already present in `settings.json` → `env` (which overrides launch env). Detection is read-only and server-local: it cannot inspect another launch shell, another machine, or project-local settings. An unknown result is not proof of force support. Explicit gateway selectors on a generated agent request take precedence over its legacy `ocx-route` fallback, even if the saved force setting changes after launch. For shell or settings overrides of generated roster agents, use an explicit gateway alias; bare Claude ids retain the older-client fallback behavior. Native aliases restore their bare model before the existing credential and model-map checks. Connected launches validate force targets against a fresh authenticated gateway catalog; failed discovery skips automatic force injection, and cached context windows alone never prove availability. + +Uniquement avant toute sortie : Un HTTP 401 authentication_error (sans error.code) portant exactement le message « OAuth access token has been revoked. » marque le compte OAuth ayant envoyé la requête comme nécessitant une nouvelle connexion et efface ses affinités de session. Avant toute sortie, un compte disponible du même pool peut prendre le relais dans les limites existantes. Sans remplaçant, le 401 original est renvoyé et le compte reste exclu jusqu’à une nouvelle connexion. Les autres 401 gardent leur traitement actuel. diff --git a/docs-site/src/content/docs/guides/claude-code.md b/docs-site/src/content/docs/guides/claude-code.md index ea1c6761289..ecdc703d3ce 100644 --- a/docs-site/src/content/docs/guides/claude-code.md +++ b/docs-site/src/content/docs/guides/claude-code.md @@ -81,6 +81,12 @@ Operational contract when enabled: access, policy and unrecognized errors stay terminal. Recovery respects model routes and send limits; if no replacement is eligible, the original 403 is returned. This also works with proactive pooling off. A 403 after assistant output starts never switches accounts. +- Before output, an exact structured **401** authentication_error with no error code and the message + “OAuth access token has been revoked.” marks the sending OAuth account as requiring + a new login and clears its session affinities. Before output, an eligible account in + the same pool may take over within existing send limits. With no eligible replacement, + the original 401 is returned; the refused account remains excluded until login. + Other 401 errors retain their existing behavior. - Token-refresh credential failures retain the existing `needsReauth` policy. Subscription renewal does not require reauthentication, but the account waits for its cooldown to expire. - If every eligible account is cooling, the proxy returns **429** (not 401) with `Retry-After` diff --git a/docs-site/src/content/docs/ja/guides/claude-code.md b/docs-site/src/content/docs/ja/guides/claude-code.md index faf0313c315..c6308ab24b8 100644 --- a/docs-site/src/content/docs/ja/guides/claude-code.md +++ b/docs-site/src/content/docs/ja/guides/claude-code.md @@ -630,3 +630,5 @@ Claude Code **2.1.257 or newer** is required for FORCE. Plugin and built-in agen The dashboard warns about old or unknown CLI versions, unavailable targets, and either variable already present in `settings.json` → `env` (which overrides launch env). Detection is read-only and server-local: it cannot inspect another launch shell, another machine, or project-local settings. An unknown result is not proof of force support. Explicit gateway selectors on a generated agent request take precedence over its legacy `ocx-route` fallback, even if the saved force setting changes after launch. For shell or settings overrides of generated roster agents, use an explicit gateway alias; bare Claude ids retain the older-client fallback behavior. Native aliases restore their bare model before the existing credential and model-map checks. Connected launches validate force targets against a fresh authenticated gateway catalog; failed discovery skips automatic force injection, and cached context windows alone never prove availability. + +出力前の応答に限り、HTTP 401 の authentication_error(error.code なし) が正確に “OAuth access token has been revoked.” を返した場合、送信した OAuth アカウントを再ログインが必要な状態にし、セッションの紐付けを解除します。出力前に限り、既存の送信上限内で同じプールの利用可能なアカウントに切り替えます。候補がなければ元の 401 を返し、再ログインまで選択から除外します。他の 401 の処理は変わりません。 diff --git a/docs-site/src/content/docs/ko/guides/claude-code.md b/docs-site/src/content/docs/ko/guides/claude-code.md index 4f530872b0c..1de6c56ba53 100644 --- a/docs-site/src/content/docs/ko/guides/claude-code.md +++ b/docs-site/src/content/docs/ko/guides/claude-code.md @@ -693,3 +693,5 @@ Claude Code **2.1.257 or newer** is required for FORCE. Plugin and built-in agen The dashboard warns about old or unknown CLI versions, unavailable targets, and either variable already present in `settings.json` → `env` (which overrides launch env). Detection is read-only and server-local: it cannot inspect another launch shell, another machine, or project-local settings. An unknown result is not proof of force support. Explicit gateway selectors on a generated agent request take precedence over its legacy `ocx-route` fallback, even if the saved force setting changes after launch. For shell or settings overrides of generated roster agents, use an explicit gateway alias; bare Claude ids retain the older-client fallback behavior. Native aliases restore their bare model before the existing credential and model-map checks. Connected launches validate force targets against a fresh authenticated gateway catalog; failed discovery skips automatic force injection, and cached context windows alone never prove availability. + +출력 전의 응답에서만 정확한 HTTP 401 authentication_error (error.code 없음) 메시지가 “OAuth access token has been revoked.”이면 요청을 보낸 OAuth 계정에 재로그인이 필요하다고 표시하고 세션 연결을 해제합니다. 출력 전에는 기존 전송 제한 안에서 같은 풀의 사용 가능한 계정으로 전환할 수 있습니다. 대체 계정이 없으면 원래 401을 반환하며, 해당 계정은 재로그인할 때까지 선택에서 제외됩니다. 다른 401의 처리는 바뀌지 않습니다. diff --git a/docs-site/src/content/docs/ru/guides/claude-code.md b/docs-site/src/content/docs/ru/guides/claude-code.md index 0499b4441c3..6bb458df309 100644 --- a/docs-site/src/content/docs/ru/guides/claude-code.md +++ b/docs-site/src/content/docs/ru/guides/claude-code.md @@ -665,3 +665,5 @@ Claude Code **2.1.257 or newer** is required for FORCE. Plugin and built-in agen The dashboard warns about old or unknown CLI versions, unavailable targets, and either variable already present in `settings.json` → `env` (which overrides launch env). Detection is read-only and server-local: it cannot inspect another launch shell, another machine, or project-local settings. An unknown result is not proof of force support. Explicit gateway selectors on a generated agent request take precedence over its legacy `ocx-route` fallback, even if the saved force setting changes after launch. For shell or settings overrides of generated roster agents, use an explicit gateway alias; bare Claude ids retain the older-client fallback behavior. Native aliases restore their bare model before the existing credential and model-map checks. Connected launches validate force targets against a fresh authenticated gateway catalog; failed discovery skips automatic force injection, and cached context windows alone never prove availability. + +Только до начала вывода: Если HTTP 401 содержит authentication_error (без error.code) с точным сообщением “OAuth access token has been revoked.”, отправившая запрос учётная запись OAuth помечается как требующая нового входа, а привязки сессий очищаются. До начала вывода возможен переход к доступной записи того же пула в пределах существующих лимитов отправки. Без замены возвращается исходный 401; запись исключается до нового входа. Обработка других 401 не меняется. diff --git a/docs-site/src/content/docs/tr/guides/claude-code.md b/docs-site/src/content/docs/tr/guides/claude-code.md index 87634321e60..b0015531c14 100644 --- a/docs-site/src/content/docs/tr/guides/claude-code.md +++ b/docs-site/src/content/docs/tr/guides/claude-code.md @@ -887,3 +887,5 @@ Claude Code **2.1.257 or newer** is required for FORCE. Plugin and built-in agen The dashboard warns about old or unknown CLI versions, unavailable targets, and either variable already present in `settings.json` → `env` (which overrides launch env). Detection is read-only and server-local: it cannot inspect another launch shell, another machine, or project-local settings. An unknown result is not proof of force support. Explicit gateway selectors on a generated agent request take precedence over its legacy `ocx-route` fallback, even if the saved force setting changes after launch. For shell or settings overrides of generated roster agents, use an explicit gateway alias; bare Claude ids retain the older-client fallback behavior. Native aliases restore their bare model before the existing credential and model-map checks. Connected launches validate force targets against a fresh authenticated gateway catalog; failed discovery skips automatic force injection, and cached context windows alone never prove availability. + +Yalnızca çıktı başlamadan önce: HTTP 401 authentication_error (error.code olmadan) iletisi tam olarak “OAuth access token has been revoked.” olduğunda, isteği gönderen OAuth hesabı yeniden oturum açılması gereken durumda işaretlenir ve oturum bağları temizlenir. Çıktı başlamadan önce mevcut gönderim sınırları içinde aynı havuzdaki uygun hesaba geçilebilir. Alternatif yoksa özgün 401 döndürülür; hesap yeniden giriş yapılana kadar seçim dışı kalır. Diğer 401 yanıtlarının işlenmesi değişmez. diff --git a/docs-site/src/content/docs/zh-cn/guides/claude-code.md b/docs-site/src/content/docs/zh-cn/guides/claude-code.md index 6559c1a8977..aafb1f310a7 100644 --- a/docs-site/src/content/docs/zh-cn/guides/claude-code.md +++ b/docs-site/src/content/docs/zh-cn/guides/claude-code.md @@ -596,3 +596,5 @@ Claude Code **2.1.257 or newer** is required for FORCE. Plugin and built-in agen The dashboard warns about old or unknown CLI versions, unavailable targets, and either variable already present in `settings.json` → `env` (which overrides launch env). Detection is read-only and server-local: it cannot inspect another launch shell, another machine, or project-local settings. An unknown result is not proof of force support. Explicit gateway selectors on a generated agent request take precedence over its legacy `ocx-route` fallback, even if the saved force setting changes after launch. For shell or settings overrides of generated roster agents, use an explicit gateway alias; bare Claude ids retain the older-client fallback behavior. Native aliases restore their bare model before the existing credential and model-map checks. Connected launches validate force targets against a fresh authenticated gateway catalog; failed discovery skips automatic force injection, and cached context windows alone never prove availability. + +仅在输出开始前的响应中,仅当 HTTP 401 的 authentication_error(无 error.code) 消息完全等于 “OAuth access token has been revoked.” 时,发送请求的 OAuth 账户会被标记为需要重新登录,并清除会话绑定。输出开始前,可在现有发送限制内切换到同一池的可用账户;没有替代账户时返回原始 401,该账户在重新登录前不会被选择。其他 401 的处理保持不变。 diff --git a/docs-site/src/content/docs/zh-tw/guides/claude-code.md b/docs-site/src/content/docs/zh-tw/guides/claude-code.md index 6e6936a0ea6..c6ec3d79339 100644 --- a/docs-site/src/content/docs/zh-tw/guides/claude-code.md +++ b/docs-site/src/content/docs/zh-tw/guides/claude-code.md @@ -672,3 +672,5 @@ Claude Code **2.1.257 or newer** is required for FORCE. Plugin and built-in agen The dashboard warns about old or unknown CLI versions, unavailable targets, and either variable already present in `settings.json` → `env` (which overrides launch env). Detection is read-only and server-local: it cannot inspect another launch shell, another machine, or project-local settings. An unknown result is not proof of force support. Explicit gateway selectors on a generated agent request take precedence over its legacy `ocx-route` fallback, even if the saved force setting changes after launch. For shell or settings overrides of generated roster agents, use an explicit gateway alias; bare Claude ids retain the older-client fallback behavior. Native aliases restore their bare model before the existing credential and model-map checks. Connected launches validate force targets against a fresh authenticated gateway catalog; failed discovery skips automatic force injection, and cached context windows alone never prove availability. + +僅在輸出開始前的回應中,僅當 HTTP 401 的 authentication_error(無 error.code) 訊息完全等於 “OAuth access token has been revoked.” 時,發送請求的 OAuth 帳戶會被標記為需要重新登入,並清除工作階段綁定。輸出開始前,可在既有發送限制內切換到同一池的可用帳戶;沒有替代帳戶時回傳原始 401,該帳戶在重新登入前不會被選取。其他 401 的處理保持不變。 diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index bb3d0bce73f..277de2ad679 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -4,6 +4,11 @@ "explicit": { "messages-request-id-headers.test.ts": "claude-integration", "messages-request-id-endpoint.test.ts": "claude-integration", + "anthropic-revoked-token.test.ts": "adapters/anthropic", + "messages-revoked-token.test.ts": "claude-integration", + "anthropic-revoked-token-continuation.test.ts": "adapters/anthropic", + "anthropic-revoked-token-sidecars.test.ts": "adapters/anthropic", + "anthropic-revoked-token-boundaries.test.ts": "adapters/anthropic", "azure-vendor-metadata.test.ts": "providers", "anthropic-instance-isolation.test.ts": "adapters/anthropic", "anthropic-instance-pool-parity.test.ts": "adapters/anthropic", "anthropic-instance-quota.test.ts": "adapters/anthropic", "anthropic-instance-recovery.test.ts": "adapters/anthropic", diff --git a/src/images/loop.ts b/src/images/loop.ts index 33df568ebaa..516b3cbc6a1 100644 --- a/src/images/loop.ts +++ b/src/images/loop.ts @@ -704,7 +704,7 @@ export async function runWithImageBridge(deps: ImageBridgeDeps): Promise= 3 })); } +async function isRevokedOAuthToken(response: Response, signal?: AbortSignal): Promise { + if (response.status !== 401) return false; + try { + const body = await readBoundedResponseBody(response.clone(), { signal, fatalUtf8: true }); + if (!body.displaySafe || body.truncated) return false; + const payload: unknown = JSON.parse(body.text); + if (!payload || typeof payload !== "object" || Array.isArray(payload) + || !("type" in payload) || payload.type !== "error" || !("error" in payload)) return false; + const error = payload.error; + return !!error && typeof error === "object" && !Array.isArray(error) + && "type" in error && error.type === "authentication_error" + && "message" in error && error.message === "OAuth access token has been revoked." + && (!("code" in error) || error.code == null); + } catch { return false; } +} + async function isAccountRefusal(response: Response, signal?: AbortSignal): Promise { try { const body = await readBoundedResponseBody(response.clone(), { signal }); @@ -91,6 +107,28 @@ export async function rotateAnthropicAccountOnResponseForInstance( && credentialGeneration(row.credential) === sent.generation && (!sent.checkProviderUuid || row.credential.accountId === sent.providerAccountUuid) ? row : undefined; }; + if (response.status === 401) { + if (options.allowAccountRefusal === false) return null; + let verdict = verdicts.get(response); + if (!verdict) { verdict = isRevokedOAuthToken(response, options.signal); verdicts.set(response, verdict); } + if (!await verdict || options.signal?.aborted || !ownedCurrent()) return null; + options.currentDecision?.(); + let marked: boolean; + try { + marked = await markAccountNeedsReauthIfGeneration(instance, sent.accountId, sent.generation, undefined, undefined, store => { + const row = store[instance]?.accounts.find(account => account.id === sent.accountId); + return !options.signal?.aborted && configuredAnthropicInstance(options.config, instance) === instance + && anthropicPhysicalSendOwnershipIsCurrent(sent, store) + && (!sent.checkProviderUuid || row?.credential.accountId === sent.providerAccountUuid); + }); + } catch { return null; } + if (!marked || configuredAnthropicInstance(options.config, instance) !== instance + || !anthropicPhysicalSendOwnershipIsCurrent(sent)) return null; + routing.clearAnthropicSessionAffinityForAccount(sent.accountId); + if (!options.canRetry || options.signal?.aborted) return null; + const decision = options.currentDecision ? options.currentDecision() : options.decision ?? null; + return pickAlternateAnthropicAccount(options.config, sent.accountId, Date.now(), decision, options.model, options.excludedAccountIds); + } if (response.status === 429) { const current = ownedCurrent(); if (!current) return null; diff --git a/src/oauth/anthropic-routing.ts b/src/oauth/anthropic-routing.ts index 5ee71f197d4..dfff1d80ffb 100644 --- a/src/oauth/anthropic-routing.ts +++ b/src/oauth/anthropic-routing.ts @@ -704,7 +704,7 @@ function createAnthropicRouting(instance: AnthropicInstanceId): AnthropicRouting excludedAccountIds?: ReadonlySet, ): string | null { if (!admitted(config)) return null; - const strategy = anthropicPoolStrategy(config); + const strategy = isAnthropicAccountPoolEnabled(config) ? anthropicPoolStrategy(config) : "quota"; const eligible = routeCandidates(getEligibleAnthropicAccounts(now, model), decision).filter(id => id !== excludeId && !excludedAccountIds?.has(id)); if (strategy === "round-robin") { return peekRoundRobinAccount(poolKey, eligible, stickyLimitForPool(config)); diff --git a/src/oauth/anthropic-send-ownership.ts b/src/oauth/anthropic-send-ownership.ts index 04db07545c4..3f589a68496 100644 --- a/src/oauth/anthropic-send-ownership.ts +++ b/src/oauth/anthropic-send-ownership.ts @@ -1,6 +1,6 @@ /** Ownership captured before a physical send, independent of subsequent cooldown observations. */ import type { OAuthAccessSnapshot } from "./index"; -import { credentialGeneration, getAccountSet } from "./store"; +import { credentialGeneration, getAccountSet, type AuthStore } from "./store"; import { isAnthropicInstanceId, type AnthropicInstanceId } from "../providers/anthropic-instance-id"; import { anthropicCooldownRecoveryFor } from "../providers/quota/anthropic-cooldown-recovery"; import { captureProviderAccountQuotaEpoch } from "../providers/quota/account-cache"; @@ -32,8 +32,8 @@ export function captureAnthropicPhysicalSendOwnership(snapshot: OAuthAccessSnaps } /** Pure ownership read: never adopt or reserve the replacement account's current incarnation. */ -export function anthropicPhysicalSendOwnershipIsCurrent(owner: AnthropicPhysicalSendOwnership): boolean { - const row = getAccountSet(owner.provider)?.accounts.find(account => account.id === owner.accountId); +export function anthropicPhysicalSendOwnershipIsCurrent(owner: AnthropicPhysicalSendOwnership, store?: AuthStore): boolean { + const row = (store ? store[owner.provider] : getAccountSet(owner.provider))?.accounts.find(account => account.id === owner.accountId); return !!row && row.loginId === owner.loginId && row.addedAt === owner.addedAt && row.credential.access === owner.accessToken && credentialGeneration(row.credential) === owner.generation && anthropicCooldownRecoveryFor(owner.provider).anthropicAccountIncarnation(owner.accountId) === owner.accountIncarnation diff --git a/src/oauth/store.ts b/src/oauth/store.ts index ac750fe0541..ef952bd9d40 100644 --- a/src/oauth/store.ts +++ b/src/oauth/store.ts @@ -1508,7 +1508,7 @@ export async function markAccountNeedsReauth( export async function mergeAccountCredential(provider:string,accountId:string,credential:OAuthCredentials,opts:{expectedGeneration?:string;afterPrePersistRead?:()=>void|Promise;assertOwnership?: (store: AuthStore) => void}={}):Promise<{superseded:false}|{superseded:true;stored:OAuthCredentials}>{const safe=normalizeCredential(credential);if(!safe)throw new Error("Refusing to persist invalid OAuth credential");if(isAnthropicInstanceId(provider))assertAnthropicCredentialSource(provider,safe);return await mutateStore(async store=>{await opts.afterPrePersistRead?.();const account=store[provider]?.accounts.find(x=>x.id===accountId);if(!account)throw new Error(`OAuth account disappeared before persist: ${provider}`);if(opts.expectedGeneration!==undefined&&credentialGeneration(account.credential)!==opts.expectedGeneration)return{superseded:true,stored:account.credential};opts.assertOwnership?.(store);account.credential=safe;delete account.refreshAttentionGeneration;if(account.needsReauthReason!=="verify_account"){delete account.needsReauth;delete account.needsReauthReason;}return{superseded:false};},[provider,accountId,safe,opts.expectedGeneration]);} // A late refresh failure must not change an operator-paused account's health. Check under // the mutation lock, not before awaiting it, so a concurrent pause cannot be overwritten. -export async function markAccountNeedsReauthIfGeneration(provider:string,accountId:string,generation:string,writerGeneration=captureConfigGeneration(),reason?:"verify_account"):Promise{const key=oauthAccountKey(provider,accountId);if(writerGeneration{const account=store[provider]?.accounts.find(x=>x.id===accountId);if(!account?.credential||account.paused||credentialGeneration(account.credential)!==generation)return false;if(writerGenerationboolean):Promise{const key=oauthAccountKey(provider,accountId);if(writerGeneration{const account=store[provider]?.accounts.find(x=>x.id===accountId);if(!account?.credential||account.paused||credentialGeneration(account.credential)!==generation)return false;if(writerGeneration => { + const send = async (recovery?: "rate-limit-429" | "oauth-401" | "oauth-account-403" | "key-429" | "key-401"): Promise => { const remaining = remainingTransientSends(); if (requestTransientPolicy && remaining <= 0) { throw new Error("native Messages transient send budget exhausted before recovery dispatch"); @@ -685,7 +685,7 @@ export async function handleNativeMessages(options: HandleNativeMessagesOptions) const oauthRetryKey = {}; const triedAccountIds = new Set(); let oauthFailovers = 0; - while (oauthBinding && (response.status === 429 || response.status === 403)) { + while (oauthBinding && (response.status === 429 || response.status === 403 || response.status === 401)) { const sendingBinding = oauthBinding; const expectedRecoverySelection = sendingBinding.selection; triedAccountIds.add(sendingBinding.snapshot.accountId); @@ -703,7 +703,7 @@ export async function handleNativeMessages(options: HandleNativeMessagesOptions) canRetry: transientSendAvailable() && oauthFailovers < ANTHROPIC_POOL_MAX_FAILOVERS_PER_REQUEST, }); if (!nextAccountId) break; - const recovery = response.status === 403 ? "oauth-account-403" : "rate-limit-429"; + const recovery = response.status === 401 ? "oauth-401" : response.status === 403 ? "oauth-account-403" : "rate-limit-429"; discard(response); oauthBinding = await resolveNativeOAuthBindingForInstance(nativeInstance!, config, { routeTarget, sessionKey: options.sessionKey, model: route.modelId, candidateAccountId: nextAccountId, diff --git a/src/server/responses/adapter-continuation.ts b/src/server/responses/adapter-continuation.ts index 72e8e0c43d8..096b19db88f 100644 --- a/src/server/responses/adapter-continuation.ts +++ b/src/server/responses/adapter-continuation.ts @@ -401,7 +401,7 @@ export function createAdapterContinuations( } } if ( - (response.status === 429 || response.status === 403) + (response.status === 429 || response.status === 403 || response.status === 401) && anthropicInstance && transportState.anthropicPoolAccountId && !isNonReplayableResponse(response) @@ -428,7 +428,7 @@ export function createAdapterContinuations( ); sealRequestAttemptIdentity(logCtx.activeAttempt, logCtx.provider, transportState.activeAdapter.name, logCtx.accountLogLabel); recordAttemptCredentialSource(logCtx.activeAttempt, route.providerName, route.provider, transportState.activeAdapter.name); - nextContinuationRecoveryKind = "anthropic-oauth-429"; + nextContinuationRecoveryKind = response.status === 401 ? "oauth-401" : "anthropic-oauth-429"; continue; } catch { // fall through to emit continuation error below diff --git a/src/server/responses/adapter-dispatch.ts b/src/server/responses/adapter-dispatch.ts index 4a7cf9490a1..9d7c7215c85 100644 --- a/src/server/responses/adapter-dispatch.ts +++ b/src/server/responses/adapter-dispatch.ts @@ -1097,7 +1097,7 @@ export async function prepareAdapterExchange( // Anthropic OAuth: recover a rate limit or proven account entitlement refusal // before output, within the shared request and account rotation limits. while ( - (upstreamResponse.status === 429 || upstreamResponse.status === 403) + (upstreamResponse.status === 429 || upstreamResponse.status === 403 || upstreamResponse.status === 401) && anthropicInstance && transportState.anthropicPoolAccountId ) { @@ -1122,7 +1122,7 @@ export async function prepareAdapterExchange( ); sealRequestAttemptIdentity(logCtx.activeAttempt, logCtx.provider, transportState.activeAdapter.name, logCtx.accountLogLabel); recordAttemptCredentialSource(logCtx.activeAttempt, route.providerName, route.provider, transportState.activeAdapter.name); - const result = await rebuildAndRefetch("anthropic-oauth-429"); + const result = await rebuildAndRefetch(upstreamResponse.status === 401 ? "oauth-401" : "anthropic-oauth-429"); if ("failed" in result) return result.failed; upstreamResponse = result; if (isNonReplayableResponse(upstreamResponse)) continue recovery; diff --git a/src/server/responses/sidecar-execution.ts b/src/server/responses/sidecar-execution.ts index 831b298994f..be3105c50af 100644 --- a/src/server/responses/sidecar-execution.ts +++ b/src/server/responses/sidecar-execution.ts @@ -204,7 +204,7 @@ export async function executeResponsesSidecars( ): Promise<{ adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind } | null> => { const antigravityValidationResponse = canRunWebSearch && route.providerName === "google-antigravity" && route.provider.authMode === "oauth" && originalResponse?.status === 403; - if (route.providerName !== "kiro" && !(anthropicInstance && originalResponse?.status === 403) + if (route.providerName !== "kiro" && !(anthropicInstance && (originalResponse?.status === 403 || originalResponse?.status === 401)) && !antigravityValidationResponse && originalResponse && originalResponse.status !== 429) return null; let antigravityVerification = false; if (antigravityValidationResponse) { @@ -368,7 +368,7 @@ export async function executeResponsesSidecars( hop.permit?.release(); return null; } - recoveryKind = "anthropic-oauth-429"; + recoveryKind = originalResponse?.status === 401 ? "oauth-401" : "anthropic-oauth-429"; sidecarBudget.ownCredentialHop(hop.permit); } else { // No key pool, no generic OAuth roster, no Anthropic pool could produce a replacement diff --git a/src/web-search/loop.ts b/src/web-search/loop.ts index 83164a004d8..8720e657411 100644 --- a/src/web-search/loop.ts +++ b/src/web-search/loop.ts @@ -577,6 +577,8 @@ export async function runWithWebSearch(deps: WebSearchLoopDeps): Promise { + f = await createAnthropicInstanceFixture(); + await f.seed(); + refusal = await import("../../../src/oauth/anthropic-account-refusal"); + ownership = await import("../../../src/oauth/anthropic-send-ownership"); +}); +afterEach(async () => { await f?.dispose(); }); + +async function bound(instance: AnthropicInstanceId, response = Response.json(payload, { status: 401 })) { + const { snapshot } = await f.admit(instance, "session-one"); + expect(snapshot).not.toBeNull(); + const owner = ownership.captureAnthropicPhysicalSendOwnership(snapshot!); + expect(owner).not.toBeNull(); + refusal.bindAnthropicRefusalCredentialForSend(response, owner!); + return { response, snapshot: snapshot!, owner: owner! }; +} +function rotate(instance: AnthropicInstanceId, response: Response, canRetry = true, allowAccountRefusal = true) { + return refusal.rotateAnthropicAccountOnResponseForInstance(instance, response, { + config: f.config, accountId: f.ids[0], model: f.model, sessionKey: "session-one", + requestKey: {}, canRetry, allowAccountRefusal, + }); +} +for (const instance of instances) for (const mutation of ["uuid", "cancel", "epoch", "remove-readd"] as const) { + test(instance + ": independent queued " + mutation + " cannot acquire terminal authority", async () => { + const { response, snapshot } = await bound(instance); + refusal.bindAnthropicRefusalCredentialForSend(response, ownership.captureAnthropicPhysicalSendOwnership(snapshot)!, + f.store.getAccountCredential(instance, f.ids[0])!.accountId); + const abort = new AbortController(); + if (mutation === "epoch") f.quota.clearAccountQuotaCache(); + if (mutation === "remove-readd") { + const original = structuredClone(f.store.getAccountSet(instance)!.accounts[0]!); + await f.store.removeAccount(instance,f.ids[0]); + const cache=await import("../../../src/providers/quota/account-cache"); + const recovery=await import("../../../src/providers/quota/anthropic-cooldown-recovery"); + const sweeper=await import("../../../src/lib/state-store-sweeper"); + const context={generation:sweeper.captureConfigGeneration()+1,providerNames:new Set(instances), + oauthAccountKeys:new Set(instances.flatMap(i=>f.store.getAccountSet(i)!.accounts.map(a=>cache.accountCacheKey(i,a.id)))), + comboIds:new Set(),comboTargets:new Set(),codexAccountIds:new Set(),configRoots:new Set()}; + recovery.reconcileAllAnthropicCooldownGenerations(context); + await f.store.mutateStore(s=>{s[instance]!.accounts.unshift(original);}); + } + let release!:()=>void,entered!:()=>void; + const hold=new Promise(r=>release=r), ready=new Promise(r=>entered=r); + const blocker=f.store.mutateStore(async()=>{entered();await hold;}); await ready; + const newer=mutation === "uuid" ? f.store.mutateStore(s=>{s[instance]!.accounts[0]!.credential.accountId="55555555-5555-4555-8555-555555555555";}) : Promise.resolve(); + const pending=refusal.rotateAnthropicAccountOnResponseForInstance(instance,response,{ + config:f.config,accountId:f.ids[0],model:f.model,canRetry:true,signal:abort.signal}); + try {for(let i=0;i<100;i++) await new Promise(r=>setImmediate(r)); if(mutation==="cancel")abort.abort();} + finally {release();} + await blocker;await newer; + expect(await pending).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance,f.ids[0])?.needsReauth).toBe(false); + }); +} +test("independent Pool 2 disable while terminal writer is queued",async()=>{ + const {response}=await bound("anthropic2"); + let release!:()=>void,entered!:()=>void; + const hold=new Promise(r=>release=r),ready=new Promise(r=>entered=r); + const blocker=f.store.mutateStore(async()=>{entered();await hold;});await ready; + const pending=rotate("anthropic2",response); + try {for(let i=0;i<100 && f.store.oauthMutationTailSnapshot().active<2;i++)await new Promise(r=>setImmediate(r)); + expect(f.store.oauthMutationTailSnapshot().active).toBe(2); f.config.providers.anthropic2!.disabled=true;} + finally {release();} + await blocker;expect(await pending).toBeNull(); + expect(f.store.getAccountCredentialWithStatus("anthropic2",f.ids[0])?.needsReauth).toBe(false); +}); + +for (const instance of instances) { + test(`${instance}: independent search loop does not offer revoked-401 recovery after routed output`, async () => { + const { runWithWebSearch } = await import("../../../src/web-search/loop"); + const { parseRequest } = await import("../../../src/responses/parser"); + const { createTranslatorBudget } = await import("../../../src/lib/translator-budget"); + const translatorBudget = createTranslatorBudget(); + globalThis.fetch = (async () => Response.json({ results: [{ title: "Fixture", url: "https://example.test/result", content: "Synthetic result" }] })) as typeof fetch; + let routedSends = 0; let rotations = 0; + const adapter: import("../../../src/adapters/base").ProviderAdapter = { + name: "mock-anthropic-output", + buildRequest: () => ({ url: "https://routed.example.test/messages", method: "POST", headers: {}, body: "{}" }), + fetchResponse: async () => ++routedSends === 1 ? new Response("ok") : Response.json(payload,{status:401}), + async *parseStream() { + yield { type: "text_delta", text: "I will check. " }; + yield { type: "tool_call_start", id: "search-1", name: "web_search" }; + yield { type: "tool_call_delta", arguments: '{"query":"fixture query"}' }; + yield { type: "tool_call_end" }; yield { type: "done" }; + }, + }; + try { + const response = await runWithWebSearch({ + parsed: parseRequest({ model: `${instance}/${f.model}`, input: "Answer briefly", stream: true, tools: [{ type: "web_search" }] }), + adapter, incomingMeta: { headers: new Headers(), providerName: instance, translatorBudget }, + backend: "exa", exaApiKey: "test-exa-key", hostedTool: { type: "web_search" }, selectedForwardHeaders: new Headers(), + settings: { model: "exa-fixture-model", reasoning: "low", timeoutMs: 30_000 }, maxSearches: 1, streamRoutedModelOutput: true, + on429: () => { rotations++; return null; }, + }); + expect(await response.text()).toContain("I will check."); + expect(routedSends).toBe(2); expect(rotations).toBe(0); + } finally { translatorBudget.dispose(); } + }); +} + diff --git a/tests/adapters/anthropic/anthropic-revoked-token-continuation.test.ts b/tests/adapters/anthropic/anthropic-revoked-token-continuation.test.ts new file mode 100644 index 00000000000..45eafc3465a --- /dev/null +++ b/tests/adapters/anthropic/anthropic-revoked-token-continuation.test.ts @@ -0,0 +1,127 @@ +/** Real Responses recovery with colliding account IDs and an instance-owned send ledger. */ +import { afterEach, beforeEach, expect, spyOn, test } from "bun:test"; +import type { OcxConfig, OcxProviderConfig } from "../../../src/types"; +import type { AnthropicInstanceId } from "../../../src/providers/anthropic-instance-id"; +import type { HandleResponsesOptions } from "../../../src/server/responses/core-options"; +import { createAnthropicInstanceFixture, instanceFixtureCredential, instanceFixtureUuid, type AnthropicInstanceFixture } from "../../helpers/anthropic-instance-fixture"; + +const instances = ["anthropic", "anthropic2"] as const; +let store: typeof import("../../../src/oauth/store"); +let routing: typeof import("../../../src/oauth/anthropic-routing"); +let resolver: typeof import("../../../src/server/adapter-resolve"); +let handleResponses: typeof import("../../../src/server/responses").handleResponses; +let fixture: AnthropicInstanceFixture; +let releaseSpend: (() => void) | undefined; +let ids: Record; +let config: OcxConfig; +let sends: Array<{ instance: AnthropicInstanceId; token: string | null; body: Record }>; +let reply: (instance: AnthropicInstanceId, index: number, body: Record) => Response | Promise; + +function answer(): Response { + return Response.json({ id: "msg_synthetic", type: "message", role: "assistant", model: "claude-sonnet-4-6", + content: [{ type: "text", text: "The answer is complete." }], stop_reason: "end_turn", usage: { input_tokens: 8, output_tokens: 6 } }); +} +function refusal(status: 403 | 429): Response { + return Response.json({ type: "error", error: status === 403 + ? { type: "permission_error", message: "Your account does not have access to Claude Code" } + : { type: "rate_limit_error", message: "Synthetic shared quota exhausted" } }, { + status, headers: { "retry-after": "30", ...(status === 429 ? { "anthropic-ratelimit-unified-5h-status": "rejected" } : {}) }, + }); +} +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise(done => { resolve = done; }); + return { promise, resolve }; +} +function post(instance: AnthropicInstanceId, options: HandleResponsesOptions = {}, body: Record = {}) { + return handleResponses(new Request("http://localhost/v1/responses", { method: "POST", + headers: { "content-type": "application/json", "session-id": "equal-session", authorization: "Bearer access-token-value-test-caller-excluded" }, + body: JSON.stringify({ model: `${instance}/claude-sonnet-4-6`, input: "Answer briefly", stream: false, ...body }), + }), config, { model: "", provider: "" }, options); +} + +beforeEach(async () => { + fixture = await createAnthropicInstanceFixture({ anthropic: { enabled: false }, anthropic2: { enabled: false } }); + fixture.quota.resetProviderQuotaReconcileStateForTests(); + await fixture.seed(); + ({ store, routing, config } = fixture); + resolver = await import("../../../src/server/adapter-resolve"); + ({ handleResponses } = await import("../../../src/server/responses")); + const { acquireOwnedSpendHome } = await import("../../helpers/owned-spend-home"); + releaseSpend = acquireOwnedSpendHome(); + ids = { anthropic: [...fixture.ids], anthropic2: [...fixture.ids] }; + sends = []; + reply = () => answer(); + for (const instance of instances) { + Object.assign(config.providers[instance]!, { baseUrl: "https://instance-bridge.example.test", models: [fixture.model], + fetch: (async (_input: RequestInfo | URL, init?: RequestInit) => { + const token = new Headers(init?.headers).get("authorization"); + const row = store.getAccountSet(instance)?.accounts.find(account => `Bearer ${account.credential.access}` === token); + const body = JSON.parse(String(init?.body)) as Record; + expect(row).toBeDefined(); + // Source identity is separate evidence; the translated adapter sends no account UUID. + expect(row!.credential.anthropicIdentity?.accountUuid).toBe(instanceFixtureUuid(instance, fixture.ids.indexOf(row!.id as typeof fixture.ids[number]) + 1)); + fixture.ledger.record({ instance, accountId: row!.id, token: token!.replace(/^Bearer /, ""), model: String(body.model) }); + sends.push({ instance, token, body }); + return reply(instance, sends.length, body); + }) as typeof fetch, + }); + } + fixture.publishConfig(); +}); +afterEach(async () => { + try { + fixture.ledger.assertNoCrossSend(); + releaseSpend?.(); releaseSpend = undefined; + const { clearResponseStateForTests } = await import("../../../src/responses/state"); + clearResponseStateForTests(); + } finally { + try { await fixture.dispose(); } + finally { fixture.quota.resetProviderQuotaReconcileStateForTests(); } + } +}); + +function revoked() { + return Response.json({ type: "error", error: { type: "authentication_error", + message: "OAuth access token has been revoked." } }, { status: 401 }); +} +function answerStream(text = "The answer is complete.") { + const message = { id: "msg_synthetic", type: "message", role: "assistant", model: fixture.model, + content: [], stop_reason: null, usage: { input_tokens: 8, output_tokens: 0 } }; + const frames = [ + { type: "message_start", message }, + { type: "content_block_start", index: 0, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 0, delta: { type: "text_delta", text } }, + { type: "content_block_stop", index: 0 }, + { type: "message_delta", delta: { stop_reason: "end_turn" }, usage: { output_tokens: 6 } }, + { type: "message_stop" }, + ]; + return new Response(frames.map(frame => "event: " + frame.type + "\ndata: " + JSON.stringify(frame) + "\n\n").join(""), + { headers: { "content-type": "text/event-stream" } }); +} + +for (const instance of instances) { + test(instance + ": pre-output empty continuation recovers revoked 401", async () => { + config.emptyCompletionRetry = true; + reply = (_instance, index, body) => index === 1 ? answerStream("") + : index === 2 ? revoked() : answerStream(); + const response = await post(instance, {}, { stream: true }); + expect(await response.text()).toContain("The answer is complete."); + expect(sends.map(send => send.token)).toEqual([ + "Bearer synthetic-" + instance + "-1-access", + "Bearer synthetic-" + instance + "-1-access", + "Bearer synthetic-" + instance + "-2-access", + ]); + expect(store.getAccountCredentialWithStatus(instance, fixture.ids[0])?.needsReauth).toBe(true); + }); + test(instance + ": continuation after assistant output cannot mark or replay", async () => { + reply = (_instance, index) => index === 1 ? answerStream("I will modify the file now.") : revoked(); + const response = await post(instance, {}, { stream: true, input: "Please modify the file now", + tools: [{ type: "function", name: "read_file", description: "read a file", parameters: { type: "object" } }] }); + expect(await response.text()).toContain("I will modify the file now."); + expect(sends.map(send => send.token)).toEqual([ + "Bearer synthetic-" + instance + "-1-access", "Bearer synthetic-" + instance + "-1-access", + ]); + expect(store.getAccountCredentialWithStatus(instance, fixture.ids[0])?.needsReauth).toBe(false); + }); +} diff --git a/tests/adapters/anthropic/anthropic-revoked-token-sidecars.test.ts b/tests/adapters/anthropic/anthropic-revoked-token-sidecars.test.ts new file mode 100644 index 00000000000..88238f393fa --- /dev/null +++ b/tests/adapters/anthropic/anthropic-revoked-token-sidecars.test.ts @@ -0,0 +1,139 @@ +/** Real Responses recovery with colliding account IDs and an instance-owned send ledger. */ +import { afterEach, beforeEach, expect, spyOn, test } from "bun:test"; +import type { OcxConfig, OcxProviderConfig } from "../../../src/types"; +import type { AnthropicInstanceId } from "../../../src/providers/anthropic-instance-id"; +import type { HandleResponsesOptions } from "../../../src/server/responses/core-options"; +import { createAnthropicInstanceFixture, instanceFixtureCredential, instanceFixtureUuid, type AnthropicInstanceFixture } from "../../helpers/anthropic-instance-fixture"; + +const instances = ["anthropic", "anthropic2"] as const; +let store: typeof import("../../../src/oauth/store"); +let routing: typeof import("../../../src/oauth/anthropic-routing"); +let resolver: typeof import("../../../src/server/adapter-resolve"); +let handleResponses: typeof import("../../../src/server/responses").handleResponses; +let fixture: AnthropicInstanceFixture; +let releaseSpend: (() => void) | undefined; +let ids: Record; +let config: OcxConfig; +let sends: Array<{ instance: AnthropicInstanceId; token: string | null; body: Record }>; +let reply: (instance: AnthropicInstanceId, index: number, body: Record) => Response | Promise; + +function answer(): Response { + return Response.json({ id: "msg_synthetic", type: "message", role: "assistant", model: "claude-sonnet-4-6", + content: [{ type: "text", text: "The answer is complete." }], stop_reason: "end_turn", usage: { input_tokens: 8, output_tokens: 6 } }); +} +function refusal(status: 403 | 429): Response { + return Response.json({ type: "error", error: status === 403 + ? { type: "permission_error", message: "Your account does not have access to Claude Code" } + : { type: "rate_limit_error", message: "Synthetic shared quota exhausted" } }, { + status, headers: { "retry-after": "30", ...(status === 429 ? { "anthropic-ratelimit-unified-5h-status": "rejected" } : {}) }, + }); +} +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise(done => { resolve = done; }); + return { promise, resolve }; +} +function post(instance: AnthropicInstanceId, options: HandleResponsesOptions = {}, body: Record = {}) { + return handleResponses(new Request("http://localhost/v1/responses", { method: "POST", + headers: { "content-type": "application/json", "session-id": "equal-session", authorization: "Bearer access-token-value-test-caller-excluded" }, + body: JSON.stringify({ model: `${instance}/claude-sonnet-4-6`, input: "Answer briefly", stream: false, ...body }), + }), config, { model: "", provider: "" }, options); +} + +beforeEach(async () => { + fixture = await createAnthropicInstanceFixture({ anthropic: { enabled: false }, anthropic2: { enabled: false } }); + fixture.quota.resetProviderQuotaReconcileStateForTests(); + await fixture.seed(); + ({ store, routing, config } = fixture); + resolver = await import("../../../src/server/adapter-resolve"); + ({ handleResponses } = await import("../../../src/server/responses")); + const { acquireOwnedSpendHome } = await import("../../helpers/owned-spend-home"); + releaseSpend = acquireOwnedSpendHome(); + ids = { anthropic: [...fixture.ids], anthropic2: [...fixture.ids] }; + sends = []; + reply = () => answer(); + for (const instance of instances) { + Object.assign(config.providers[instance]!, { baseUrl: "https://instance-bridge.example.test", models: [fixture.model], + fetch: (async (_input: RequestInfo | URL, init?: RequestInit) => { + const token = new Headers(init?.headers).get("authorization"); + const row = store.getAccountSet(instance)?.accounts.find(account => `Bearer ${account.credential.access}` === token); + const body = JSON.parse(String(init?.body)) as Record; + expect(row).toBeDefined(); + // Source identity is separate evidence; the translated adapter sends no account UUID. + expect(row!.credential.anthropicIdentity?.accountUuid).toBe(instanceFixtureUuid(instance, fixture.ids.indexOf(row!.id as typeof fixture.ids[number]) + 1)); + fixture.ledger.record({ instance, accountId: row!.id, token: token!.replace(/^Bearer /, ""), model: String(body.model) }); + sends.push({ instance, token, body }); + return reply(instance, sends.length, body); + }) as typeof fetch, + }); + } + fixture.publishConfig(); +}); +afterEach(async () => { + try { + fixture.ledger.assertNoCrossSend(); + releaseSpend?.(); releaseSpend = undefined; + const { clearResponseStateForTests } = await import("../../../src/responses/state"); + clearResponseStateForTests(); + } finally { + try { await fixture.dispose(); } + finally { fixture.quota.resetProviderQuotaReconcileStateForTests(); } + } +}); + +function revoked() { + return Response.json({ type: "error", error: { type: "authentication_error", + message: "OAuth access token has been revoked." } }, { status: 401 }); +} +function answerStream(text = "The answer is complete.") { + const message = { id: "msg_synthetic", type: "message", role: "assistant", model: fixture.model, + content: [], stop_reason: null, usage: { input_tokens: 8, output_tokens: 0 } }; + const frames = [ + { type: "message_start", message }, + { type: "content_block_start", index: 0, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 0, delta: { type: "text_delta", text } }, + { type: "content_block_stop", index: 0 }, + { type: "message_delta", delta: { stop_reason: "end_turn" }, usage: { output_tokens: 6 } }, + { type: "message_stop" }, + ]; + return new Response(frames.map(frame => "event: " + frame.type + "\ndata: " + JSON.stringify(frame) + "\n\n").join(""), + { headers: { "content-type": "text/event-stream" } }); +} + +for (const instance of instances) for (const kind of ["search", "image"] as const) { + test(instance + ": exact revoked 401 reaches real " + kind + " recovery", async () => { + if (kind === "search") config.webSearchSidecar = { backend: "anthropic", enabled: true }; + else { + config.images = { bridgeEnabled: true }; + config.providers.xai = { adapter: "openai-chat", baseUrl: "https://api.x.ai/v1", + authMode: "key", apiKey: "synthetic-image-key" }; + } + reply = (_instance, index) => index === 1 ? revoked() : answerStream(); + const response = await post(instance, {}, { stream: true, + tools: [{ type: kind === "search" ? "web_search" : "image_generation" }] }); + expect(response.status).toBe(200); + expect(await response.text()).toContain("The answer is complete."); + expect(sends).toHaveLength(2); + expect(fixture.ledger.sends.map(send => [send.instance, send.accountId])).toEqual([ + [instance, fixture.ids[0]], [instance, fixture.ids[1]], + ]); + expect(store.getAccountCredentialWithStatus(instance, fixture.ids[0])?.needsReauth).toBe(true); + const other = instance === "anthropic" ? "anthropic2" : "anthropic"; + expect(store.getAccountCredentialWithStatus(other, fixture.ids[0])?.needsReauth).toBe(false); + }); + test(instance + ": unknown 401 in " + kind + " cannot mark or fail over", async () => { + if (kind === "search") config.webSearchSidecar = { backend: "anthropic", enabled: true }; + else { + config.images = { bridgeEnabled: true }; + config.providers.xai = { adapter: "openai-chat", baseUrl: "https://api.x.ai/v1", + authMode: "key", apiKey: "synthetic-image-key" }; + } + reply = () => Response.json({ type: "error", error: { type: "authentication_error", + message: "OAuth access token has expired." } }, { status: 401 }); + const response = await post(instance, {}, { stream: true, + tools: [{ type: kind === "search" ? "web_search" : "image_generation" }] }); + await response.text(); + expect(sends).toHaveLength(1); + expect(store.getAccountCredentialWithStatus(instance, fixture.ids[0])?.needsReauth).toBe(false); + }); +} diff --git a/tests/adapters/anthropic/anthropic-revoked-token.test.ts b/tests/adapters/anthropic/anthropic-revoked-token.test.ts new file mode 100644 index 00000000000..247a775cb03 --- /dev/null +++ b/tests/adapters/anthropic/anthropic-revoked-token.test.ts @@ -0,0 +1,234 @@ +import { afterEach, beforeEach, expect, spyOn, test } from "bun:test"; +import type { AnthropicInstanceId } from "../../../src/providers/anthropic-instance-id"; +import { createAnthropicInstanceFixture, type AnthropicInstanceFixture } from "../../helpers/anthropic-instance-fixture"; + +let f: AnthropicInstanceFixture; +let refusal: typeof import("../../../src/oauth/anthropic-account-refusal"); +let ownership: typeof import("../../../src/oauth/anthropic-send-ownership"); +const instances = ["anthropic", "anthropic2"] as const; +const message = "OAuth access token has been revoked."; +const payload = { type: "error", error: { type: "authentication_error", message } }; + +beforeEach(async () => { + f = await createAnthropicInstanceFixture(); + await f.seed(); + refusal = await import("../../../src/oauth/anthropic-account-refusal"); + ownership = await import("../../../src/oauth/anthropic-send-ownership"); +}); +afterEach(async () => { await f?.dispose(); }); + +async function bound(instance: AnthropicInstanceId, response = Response.json(payload, { status: 401 })) { + const { snapshot } = await f.admit(instance, "session-one"); + expect(snapshot).not.toBeNull(); + const owner = ownership.captureAnthropicPhysicalSendOwnership(snapshot!); + expect(owner).not.toBeNull(); + refusal.bindAnthropicRefusalCredentialForSend(response, owner!); + return { response, snapshot: snapshot!, owner: owner! }; +} +function rotate(instance: AnthropicInstanceId, response: Response, canRetry = true, allowAccountRefusal = true) { + return refusal.rotateAnthropicAccountOnResponseForInstance(instance, response, { + config: f.config, accountId: f.ids[0], model: f.model, sessionKey: "session-one", + requestKey: {}, canRetry, allowAccountRefusal, + }); +} +for (const instance of instances) { + const other = instance === "anthropic" ? "anthropic2" : "anthropic"; + for (const canRetry of [true, false]) { + test(instance + ": mark and clear all affinity with retry allowance " + canRetry, async () => { + const { response } = await bound(instance); + await f.admit(instance, "session-two"); + await f.admit(other, "other-session"); + const otherAffinity = f.routing.anthropicRoutingFor(other).anthropicSessionAffinitySizeForTests(); + expect(f.routing.anthropicRoutingFor(instance).anthropicSessionAffinitySizeForTests()).toBe(2); + expect(await rotate(instance, response, canRetry)).toBe(canRetry ? f.ids[1] : null); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(true); + expect(f.routing.anthropicRoutingFor(instance).anthropicSessionAffinitySizeForTests()).toBe(0); + expect(f.routing.anthropicRoutingFor(instance).getEligibleAnthropicAccounts()).toEqual([f.ids[1]]); + expect(f.store.getAccountCredentialWithStatus(other, f.ids[0])?.needsReauth).toBe(false); + expect(f.routing.anthropicRoutingFor(other).anthropicSessionAffinitySizeForTests()).toBe(otherAffinity); + expect(f.routing.anthropicRoutingFor(instance).getAnthropicAccountHealthSnapshot(f.ids[0])).toBeNull(); + expect(await response.json()).toEqual(payload); + }); + } + test(instance + ": no eligible sibling still marks reauthentication", async () => { + await f.store.setAccountPaused(instance, f.ids[1], true); + const { response } = await bound(instance); + expect(await rotate(instance, response)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(true); + expect(response.status).toBe(401); + expect(await response.json()).toEqual(payload); + }); + test(instance + ": output commitment disables account-refusal recovery", async () => { + const { response } = await bound(instance); + expect(await rotate(instance, response, true, false)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + expect(await response.json()).toEqual(payload); + }); + for (const sameBytes of [false, true]) { + test(instance + ": relogin before late refusal is preserved " + sameBytes, async () => { + const { response } = await bound(instance); + const credential = f.store.getAccountCredential(instance, f.ids[0])!; + await f.store.saveAccountCredential(instance, f.ids[0], sameBytes ? credential + : { ...credential, access: "synthetic-" + instance + "-replacement-access" }, { rotateLoginId: true }); + expect(await rotate(instance, response)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + }); + } + const negatives = [ + { ...payload, error: { type: "authentication_error", message: "OAuth access token has expired." } }, + { ...payload, error: { type: "permission_error", message } }, + { ...payload, error: { type: "authentication_error", message: "Quoted: " + message } }, + { ...payload, error: { type: "authentication_error", message: message.toLowerCase() } }, + { ...payload, error: { type: "authentication_error", message: " " + message + " " } }, + { ...payload, error: { type: "authentication_error", message, code: "invalid_api_key" } }, + { error: payload.error }, [], + ]; + for (const [index, body] of negatives.entries()) { + test(instance + ": unmatched 401 body " + index + " has no account effect", async () => { + const { response } = await bound(instance, Response.json(body, { status: 401 })); + expect(await rotate(instance, response)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + expect(f.routing.anthropicRoutingFor(instance).anthropicSessionAffinitySizeForTests()).toBe(1); + expect(await response.json()).toEqual(body); + }); + } + for (const body of ["{", JSON.stringify({ ...payload, padding: "x".repeat(65_536) })]) { + test(instance + ": malformed or oversized refusal is not authority " + body.length, async () => { + const { response } = await bound(instance, new Response(body, { status: 401 })); + expect(await rotate(instance, response)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + expect(await response.text()).toBe(body); + }); + } + for (const status of [400, 403, 429]) { + test(instance + ": exact sentence on other status is not revocation " + status, async () => { + const { response } = await bound(instance, Response.json(payload, { status })); + expect(await rotate(instance, response, false)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + }); + } + test(instance + ": unbound matching response is ignored", async () => { + expect(await rotate(instance, Response.json(payload, { status: 401 }))).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + }); + test(instance + ": cancellation preserves the account", async () => { + const { response } = await bound(instance); + const abort = new AbortController(); abort.abort(); + expect(await refusal.rotateAnthropicAccountOnResponseForInstance(instance, response, { + config: f.config, accountId: f.ids[0], model: f.model, signal: abort.signal, canRetry: true, + })).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + }); + test(instance + ": explicit relogin clears terminal state", async () => { + const { response } = await bound(instance); + await rotate(instance, response, false); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(true); + const credential = f.store.getAccountCredential(instance, f.ids[0])!; + await f.store.saveAccountCredential(instance, f.ids[0], { ...credential, + access: "synthetic-" + instance + "-new-login-access" }, { rotateLoginId: true }); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + expect(f.routing.anthropicRoutingFor(instance).getEligibleAnthropicAccounts()).toContain(f.ids[0]); + }); +} + +for (const instance of instances) { + test(instance + ": audit malformed UTF-8 is not terminal authority", async () => { + const prefix=new TextEncoder().encode(JSON.stringify({...payload,padding:"X"})); + const index=prefix.indexOf(88,prefix.length-6); expect(index).toBeGreaterThan(0); prefix[index]=255; + const {response}=await bound(instance,new Response(prefix,{status:401})); + expect(await rotate(instance,response)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance,f.ids[0])?.needsReauth).toBe(false); + }); + test(instance + ": audit delayed body relogin preserves same-byte replacement", async () => { + let release!:()=>void; const ready=new Promise(r=>release=r); + let entered!:()=>void; const begin=new Promise(r=>entered=r); + const response=new Response(new ReadableStream({ async pull(controller) { entered(); await ready; controller.enqueue(new TextEncoder().encode(JSON.stringify(payload))); controller.close(); } }),{status:401}); + await bound(instance,response); + const pending=rotate(instance,response); await begin; + await f.store.saveAccountCredential(instance,f.ids[0],f.store.getAccountCredential(instance,f.ids[0])!,{rotateLoginId:true}); + release(); expect(await pending).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance,f.ids[0])?.needsReauth).toBe(false); + }); +} + +for (const instance of instances) for (const replacement of ["same-login", "pause"]) { + test(instance + ": audit locked snapshot fences queued " + replacement, async () => { + const {response}=await bound(instance); + let release!:()=>void; const hold=new Promise(r=>release=r); + let entered!:()=>void; const ready=new Promise(r=>entered=r); + const blocker=f.store.mutateStore(async()=>{entered();await hold;}); await ready; + const newer=replacement === "pause" ? f.store.setAccountPaused(instance,f.ids[0],true) + : f.store.saveAccountCredential(instance,f.ids[0],f.store.getAccountCredential(instance,f.ids[0])!,{rotateLoginId:true}); + const pending=rotate(instance,response); + try { + for (let i=0;i<100 && f.store.oauthMutationTailSnapshot().active<3;i++) await new Promise(r=>setImmediate(r)); + expect(f.store.oauthMutationTailSnapshot().active).toBe(3); + } finally {release();} + await blocker; await newer; expect(await pending).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance,f.ids[0])?.needsReauth).toBe(false); + }); +} + + +for (const instance of instances) { + test(instance + ": exact null-code evidence is accepted", async () => { + const { response } = await bound(instance, Response.json({ ...payload, error: { ...payload.error, code: null } }, { status: 401 })); + expect(await rotate(instance, response)).toBe(f.ids[1]); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(true); + }); + test(instance + ": singleton is marked without a replay proposal", async () => { + await f.store.removeAccount(instance, f.ids[1]); + const { response } = await bound(instance); + expect(await rotate(instance, response)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(true); + expect(await response.json()).toEqual(payload); + }); + test(instance + ": rejected persistence offers no alternate and keeps the response", async () => { + const { response } = await bound(instance); + const failure = spyOn(f.store, "markAccountNeedsReauthIfGeneration").mockRejectedValue(new Error("synthetic write failure")); + try { + expect(await rotate(instance, response)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + expect(await response.json()).toEqual(payload); + } finally { failure.mockRestore(); } + }); + for (const fallback of [false, true]) { + test(instance + ": revoked recovery retains model route fallback " + fallback, async () => { + const pool = { enabled: true, routes: [{ name: "strict", match: f.model, accounts: [f.ids[0]], fallback }] }; + if (instance === "anthropic") f.config.anthropicAccountPool = pool; + else f.config.providers.anthropic2!.anthropicAccountPool = pool; + const { response } = await bound(instance); + const decision = (await import("../../../src/oauth/anthropic-model-routes")).resolveAnthropicModelRouteForInstance(instance, f.config, f.model).decision; + expect(await refusal.rotateAnthropicAccountOnResponseForInstance(instance, response, { + config: f.config, accountId: f.ids[0], model: f.model, decision, canRetry: true, + })).toBe(fallback ? f.ids[1] : null); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(true); + }); + } + test(instance + ": response from the other instance is never authority", async () => { + const other = instance === "anthropic" ? "anthropic2" : "anthropic"; + const { response } = await bound(other); + expect(await rotate(instance, response)).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + expect(f.store.getAccountCredentialWithStatus(other, f.ids[0])?.needsReauth).toBe(false); + }); + test(instance + ": cancellation during body classification makes no durable change", async () => { + let release!: () => void; + const hold = new Promise(resolve => { release = resolve; }); + let entered!: () => void; + const ready = new Promise(resolve => { entered = resolve; }); + const response = new Response(new ReadableStream({ + async pull(controller) { entered(); await hold; controller.enqueue(new TextEncoder().encode(JSON.stringify(payload))); controller.close(); }, + }), { status: 401 }); + await bound(instance, response); + const abort = new AbortController(); + const pending = refusal.rotateAnthropicAccountOnResponseForInstance(instance, response, { + config: f.config, accountId: f.ids[0], model: f.model, canRetry: true, signal: abort.signal, + }); + await ready; + abort.abort(); + release(); + expect(await pending).toBeNull(); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(false); + }); +} diff --git a/tests/claude-integration/messages-revoked-token.test.ts b/tests/claude-integration/messages-revoked-token.test.ts new file mode 100644 index 00000000000..ea2a68c4736 --- /dev/null +++ b/tests/claude-integration/messages-revoked-token.test.ts @@ -0,0 +1,215 @@ +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import { existsSync, readFileSync } from "node:fs"; +import type { AnthropicInstanceId } from "../../src/providers/anthropic-instance-id"; +import type { OcxConfig } from "../../src/types"; +import type { RouteResult } from "../../src/router"; +import type { GenerationContext } from "../../src/lib/state-store-sweeper"; +import { + anthropicInstanceBarrier, createAnthropicInstanceFixture, instanceFixtureCredential, instanceFixtureUuid, + type AnthropicInstanceFixture, +} from "../helpers/anthropic-instance-fixture"; + +type Rec = Record; +type Sent = { instance: AnthropicInstanceId; url: string; headers: Headers; body: Rec }; +const INSTANCES = ["anthropic", "anthropic2"] as const; +const CALLER = "sk-ant-fixture-caller"; +let f: AnthropicInstanceFixture; +let sent: Sent[]; +let ingress: typeof import("../../src/server/claude-messages"); +let native: typeof import("../../src/server/messages-native"); +let binding: typeof import("../../src/server/messages-native-oauth"); +let planner: typeof import("../../src/protocols/plan-snapshot"); +let settings: typeof import("../../src/protocols/settings"); +let identity: typeof import("../../src/oauth/anthropic-identity"); +let pacing: typeof import("../../src/providers/request-pacing"); +let logs: typeof import("../../src/server/request-log"); +let releaseSpend: (() => void) | undefined; +const restorations: Array<() => void> = []; + +beforeEach(async () => { + // The shared fixture creates all homes and blocks network before runtime imports. + f = await createAnthropicInstanceFixture({ anthropic: { enabled: false }, anthropic2: { enabled: false } }); + [ingress, native, binding, planner, settings, identity, pacing, logs] = await Promise.all([ + import("../../src/server/claude-messages"), import("../../src/server/messages-native"), + import("../../src/server/messages-native-oauth"), import("../../src/protocols/plan-snapshot"), + import("../../src/protocols/settings"), import("../../src/oauth/anthropic-identity"), + import("../../src/providers/request-pacing"), import("../../src/server/request-log"), + ]); + releaseSpend = (await import("../helpers/owned-spend-home")).acquireOwnedSpendHome(); + (await import("../../src/responses/reasoning-replay-cache")).clearReasoningReplayCacheForTests(); + sent = []; + f.config.protocols = { rollout: { managedMessagesNative: true, managedMessagesNativeOAuth: true } }; + for (const instance of INSTANCES) { + f.config.providers[instance]!.models = [f.model]; + } + f.publishConfig(); +}); + +afterEach(async () => { + for (const restore of restorations.splice(0).reverse()) restore(); + pacing?.resetProviderRequestPacingForTest(); + releaseSpend?.(); + releaseSpend = undefined; + await f?.dispose(); +}); + +async function seed(instances: readonly AnthropicInstanceId[] = INSTANCES) { + await f.seed(instances); + // Seed clones the pure config; attach in-process transports only after that operation. + for (const instance of INSTANCES) f.config.providers[instance]!.fetch = transport(instance); +} + +function answer(model: string): Rec { + return { id: "msg_instance_fixture", type: "message", role: "assistant", model, + content: [{ type: "text", text: "fixture reply" }], stop_reason: "end_turn", stop_sequence: null, + usage: { input_tokens: 9, output_tokens: 3 } }; +} + +/** Responses forces streaming upstream even for a non-streaming Messages caller. */ +function answerForWire(send: Sent): Response { + if (send.body.stream !== true) return Response.json(answer(f.model)); + const frames: Rec[] = [ + { type: "message_start", message: { ...answer(f.model), content: [], stop_reason: null, + usage: { input_tokens: 9, output_tokens: 0 } } }, + { type: "content_block_start", index: 0, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "fixture reply" } }, + { type: "content_block_stop", index: 0 }, + { type: "message_delta", delta: { stop_reason: "end_turn", stop_sequence: null }, usage: { output_tokens: 3 } }, + { type: "message_stop" }, + ]; + return new Response(frames.map(frame => `event: ${frame.type}\ndata: ${JSON.stringify(frame)}\n\n`).join(""), + { headers: { "content-type": "text/event-stream" } }); +} + +function transport(instance: AnthropicInstanceId, response?: (send: Sent) => Response | Promise): typeof fetch { + return (async (input, init) => { + const headers = new Headers(init?.headers); + const token = headers.get("authorization")?.replace(/^Bearer /, "") ?? ""; + const slot = [1, 2].find(candidate => token === instanceFixtureCredential(instance, candidate).access); + // Independent shared ledger detects wrong-instance physical bearer even for equal IDs. + const parsedBody = JSON.parse(String(init?.body)) as Rec; + const metadata = parsedBody.metadata as { user_id?: string } | undefined; + const uuid = metadata?.user_id ? (JSON.parse(metadata.user_id) as { account_uuid?: string }).account_uuid : undefined; + f.ledger.record({ instance, accountId: slot ? f.ids[slot - 1]! : f.ids[0], token, uuid }); + expect(headers.has("x-api-key")).toBe(false); + expect(token).not.toBe(CALLER); + const entry = { instance, url: String(input), headers, body: parsedBody }; + sent.push(entry); + return response ? await response(entry) : answerForWire(entry); + }) as typeof fetch; +} + +function body(model = `anthropic2/${f.model}`, extra: Rec = {}): Rec { + return { model, max_tokens: 64, stream: false, + metadata: { user_id: JSON.stringify({ account_uuid: instanceFixtureUuid("anthropic", 1), device_id: "fixture-device", session_id: f.sessionKey }) }, + messages: [{ role: "user", content: "fixture question" }], ...extra }; +} + +async function send(model = `anthropic2/${f.model}`, extra: Rec = {}, options: { nativeCaller?: boolean; sessionKey?: string } = {}) { + const requestId = crypto.randomUUID(); + const logCtx = { model: "", provider: "" }; + const requestBody = body(model, extra); + if (options.sessionKey) { + const metadata = requestBody.metadata as { user_id: string }; + metadata.user_id = JSON.stringify({ ...JSON.parse(metadata.user_id), session_id: options.sessionKey }); + } + const response = await ingress.handleClaudeMessages(new Request("http://localhost/v1/messages", { + // Managed parity cases use ordinary admission fixtures. Caller-forward exclusion cases + // explicitly supply a classified Anthropic caller bearer while passthrough stays enabled. + method: "POST", headers: { "content-type": "application/json", + authorization: `Bearer ${options.nativeCaller ? CALLER : "fixture-admission-token"}`, + "x-api-key": options.nativeCaller ? CALLER : "fixture-caller-key", "x-session-id": options.sessionKey ?? f.sessionKey }, + body: JSON.stringify(requestBody), + }), f.config, logCtx, { requestId, start: Date.now() }); + const text = await response.text(); + const rows = logs.getRequestLogEntries().filter(row => row.requestId === requestId); + expect(JSON.stringify(rows)).not.toContain(CALLER); + if (existsSync(f.store.getAuthStorePath())) expect(readFileSync(f.store.getAuthStorePath(), "utf8")).not.toContain(CALLER); + expect(JSON.stringify(f.config)).not.toContain(CALLER); + f.ledger.assertNoCrossSend(); + return { response, text, rows }; +} + +function revoked() { + return Response.json({ type: "error", error: { type: "authentication_error", + message: "OAuth access token has been revoked." } }, { status: 401 }); +} +for (const instance of INSTANCES) for (const enabled of [true, false]) { + for (const nativeMode of [true, false]) for (const stream of [true, false]) { + test("revoked recovery pool=" + enabled + " native=" + nativeMode + " stream=" + stream + " " + instance, async () => { + await seed(); + if (instance === "anthropic") f.config.anthropicAccountPool = { enabled }; + else f.config.providers.anthropic2!.anthropicAccountPool = { enabled }; + f.config.protocols = { rollout: { managedMessagesNative: nativeMode, managedMessagesNativeOAuth: nativeMode } }; + f.config.providers[instance]!.fetch = transport(instance, row => sent.length === 1 ? revoked() : answerForWire(row)); + const { response, text, rows } = await send(instance + "/" + f.model, { stream }); + expect(response.status, text).toBe(200); + expect(text).toContain("fixture reply"); + expect(sent).toHaveLength(2); + expect(f.ledger.sends.map(row => row.accountId)).toEqual([f.ids[0], f.ids[1]]); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(true); + const other = instance === "anthropic" ? "anthropic2" : "anthropic"; + expect(f.store.getAccountCredentialWithStatus(other, f.ids[0])?.needsReauth).toBe(false); + expect(f.store.getAccountSet(other)?.activeAccountId).toBe(f.ids[0]); + expect(JSON.stringify(rows)).toContain("oauth-401"); + expect(rows[0]?.protocolTrace?.mode).toBe(nativeMode ? "native" : "legacy-bridge"); + if (nativeMode) expect(f.ledger.sends.map(row => row.uuid)).toEqual([ + instanceFixtureUuid(instance, 1), instanceFixtureUuid(instance, 2), + ]); + const subsequent = await send(instance + "/" + f.model, { stream }); + expect(subsequent.response.status).toBe(200); + expect(sent).toHaveLength(3); + expect(f.ledger.sends[2]!.accountId).toBe(f.ids[1]); + }); + } + test(instance + ": both revoked accounts retain the final upstream 401", async () => { + await seed(); + if (instance === "anthropic") f.config.anthropicAccountPool = { enabled }; + else f.config.providers.anthropic2!.anthropicAccountPool = { enabled }; + f.config.providers[instance]!.fetch = transport(instance, () => revoked()); + const { response, text } = await send(instance + "/" + f.model); + expect(response.status, text).toBe(401); + expect(text).toContain("OAuth access token has been revoked."); + expect(sent).toHaveLength(2); + for (const id of f.ids) expect(f.store.getAccountCredentialWithStatus(instance, id)?.needsReauth).toBe(true); + await send(instance + "/" + f.model); + expect(sent).toHaveLength(2); + }); + test(instance + ": unavailable sibling preserves original 401", async () => { + await seed(); + if (instance === "anthropic") f.config.anthropicAccountPool = { enabled }; + else f.config.providers.anthropic2!.anthropicAccountPool = { enabled }; + await f.store.setAccountPaused(instance, f.ids[1], true); + f.config.providers[instance]!.fetch = transport(instance, () => revoked()); + const { response, text } = await send(instance + "/" + f.model); + expect(response.status, text).toBe(401); + expect(sent).toHaveLength(1); + expect(f.store.getAccountCredentialWithStatus(instance, f.ids[0])?.needsReauth).toBe(true); + }); +} +for (const instance of INSTANCES) { + test(instance + ": stale recovery selection does not override a newer manual choice", async () => { + await seed(); + const selection = f.store.captureOAuthAccountSelection(instance); + await f.store.setActiveAccount(instance, f.ids[1]); + const result = await binding.resolveNativeOAuthBindingForInstance(instance, f.config, { + model: f.model, candidateAccountId: f.ids[0], expectedRecoverySelection: selection, + expectedRecoveryRouteDecision: null, + }); + expect(result.snapshot.accountId).toBe(f.ids[1]); + }); + test(instance + ": recovery candidate still respects the current model route", async () => { + await seed(); + const selection = f.store.captureOAuthAccountSelection(instance); + const pool = { enabled: true, routes: [{ name: "strict", match: f.model, accounts: [f.ids[0]] }] }; + if (instance === "anthropic") f.config.anthropicAccountPool = pool; + else f.config.providers.anthropic2!.anthropicAccountPool = pool; + const route = (await import("../../src/oauth/anthropic-model-routes")).resolveAnthropicModelRouteForInstance(instance, f.config, f.model); + await expect(binding.resolveNativeOAuthBindingForInstance(instance, f.config, { + model: f.model, candidateAccountId: f.ids[1], expectedRecoverySelection: selection, + expectedRecoveryRouteDecision: route.decision, + })).rejects.toThrow(); + expect(sent).toHaveLength(0); + expect(f.store.getAccountSet(instance)?.activeAccountId).toBe(f.ids[0]); + }); +} diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 281a3729757..f3236619f76 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -1,6 +1,11 @@ { "messages-request-id-headers.test.ts": "claude-integration", "messages-request-id-endpoint.test.ts": "claude-integration", + "anthropic-revoked-token.test.ts": "adapters/anthropic", + "messages-revoked-token.test.ts": "claude-integration", + "anthropic-revoked-token-continuation.test.ts": "adapters/anthropic", + "anthropic-revoked-token-sidecars.test.ts": "adapters/anthropic", + "anthropic-revoked-token-boundaries.test.ts": "adapters/anthropic", "azure-vendor-metadata.test.ts": "providers", "anthropic-instance-isolation.test.ts": "adapters/anthropic", "anthropic-instance-pool-parity.test.ts": "adapters/anthropic", "anthropic-instance-quota.test.ts": "adapters/anthropic", "anthropic-instance-recovery.test.ts": "adapters/anthropic", From f84a73fd69ea31129a6af72b78552ee3e6206035 Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 20:55:49 +0900 Subject: [PATCH 2/2] test(anthropic): drop trailing blank line --- .../anthropic/anthropic-revoked-token-boundaries.test.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/tests/adapters/anthropic/anthropic-revoked-token-boundaries.test.ts b/tests/adapters/anthropic/anthropic-revoked-token-boundaries.test.ts index 05cfdea1751..5f7f3469b89 100644 --- a/tests/adapters/anthropic/anthropic-revoked-token-boundaries.test.ts +++ b/tests/adapters/anthropic/anthropic-revoked-token-boundaries.test.ts @@ -108,4 +108,3 @@ for (const instance of instances) { } finally { translatorBudget.dispose(); } }); } -