diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index f3f16fae58b..4b10b4080e6 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -200,6 +200,7 @@ "start-args.test.ts": "cli", "start-ownership-publication.test.ts": "cli", "responses-core-modules.test.ts": "responses", + "responses-devin-401-replay.test.ts": "responses", "responses-grok-devin-preflight.test.ts": "responses", "responses-passthrough-transient-policy.test.ts": "responses", "responses-spend-ledger-wiring.test.ts": "responses", diff --git a/src/oauth/devin.ts b/src/oauth/devin.ts index 75d9c68ce0f..33c222ac3d6 100644 --- a/src/oauth/devin.ts +++ b/src/oauth/devin.ts @@ -19,6 +19,7 @@ import { DEFAULT_REGION, type WindsurfRegion } from "./devin/types"; import { registerUser } from "./devin/register-user"; import { DEVIN_DEFAULT_API_SERVER, resolveDevinApiBaseUrl, validateDevinApiBaseUrl } from "./devin/api-base"; import { readDevinCliCredentialOutcome } from "./devin/cli-import"; +import { CloudAuthError, mintUserJwt } from "../adapters/devin/cloud-direct/auth"; import { getCredential, listAccounts } from "./store"; import { DEPRECATED_OAUTH_PROVIDER_ALIASES } from "./index"; @@ -263,14 +264,111 @@ export async function loginDevin( return loginDevinBrowser(ctrl, DEFAULT_REGION); } +/** + * A CLI-imported account follows the CLI: after `devin auth login` rewrites the + * credential file, the copy stored at import time is stale while the file holds + * a live key. The session JWT carries no expiry, so the upstream 401 is the only + * signal, and this re-read runs only on that forced refresh. + * + * Adoption fails closed on identity. The session token carries only a + * session_id, so the file key's account is established by minting a user_jwt + * with it (GetUserJwt answers with auth_uid and email). A key Cognition refuses + * (401/403) or whose token lacks auth_uid is refused; a mint that fails for any + * other reason throws a non-terminal error, so the account is not flagged and + * the next 401 retries. The minted identity must then + * not contradict the slot's recorded accountId or email, and no other stored + * account may own the key or that identity. + * + * A slot imported from the CLI records no identity (its token has none to + * give), so its first adoption rests on the last two rules alone: it is by + * definition "whatever the CLI is signed into", and the adopted credential + * records the minted identity, which makes every later adoption strict. + */ +/** An identity probe must not hold the per-account refresh lock for the mint's full 30s. */ +const DEVIN_IDENTITY_MINT_TIMEOUT_MS = 5_000; + +/** Neither a live key nor a dead one: the refresh must fail without flagging the account. */ +class DevinIdentityProbeUnavailableError extends Error { + constructor() { + super("Could not confirm the Devin CLI session identity right now; retry shortly."); + this.name = "DevinIdentityProbeUnavailableError"; + } +} + +const normalizedEmail = (value: string | undefined): string | undefined => + value?.trim().toLowerCase() || undefined; + +async function rereadDevinCliCredential( + stored: OAuthCredentials, + signal: AbortSignal | undefined, + currentAccountId: string | undefined, +): Promise { + const outcome = readDevinCliCredentialOutcome(); + // A file that exists but cannot be read or parsed may be mid-write by `devin auth login` + // or briefly locked; like a failed identity mint, that says nothing about the key. + if (outcome.kind === "unreadable" || outcome.kind === "incomplete") throw new DevinIdentityProbeUnavailableError(); + if (outcome.kind !== "ok" || outcome.file.apiKey === stored.access) return undefined; + const apiBaseUrl = validateDevinApiBaseUrl(outcome.file.apiServerUrl); + if (apiBaseUrl === undefined) return undefined; + if (findDevinCredentialOwner("devin", outcome.file.apiKey) !== undefined) return undefined; + let minted: Record | undefined; + try { + const timeout = AbortSignal.timeout(DEVIN_IDENTITY_MINT_TIMEOUT_MS); + const probeSignal = signal ? AbortSignal.any([signal, timeout]) : timeout; + minted = decodeJwtPayload((await mintUserJwt(outcome.file.apiKey, apiBaseUrl, probeSignal)).jwt); + } catch (error) { + // Only Cognition refusing the key is evidence the file holds no live session. A timeout, + // DNS failure, 5xx or 429 says nothing about the key; flagging the account on one would + // strand it, because a needsReauth slot is never refreshed again. + if (error instanceof CloudAuthError && (error.status === 401 || error.status === 403)) return undefined; + throw new DevinIdentityProbeUnavailableError(); + } + const authUid = typeof minted?.auth_uid === "string" && minted.auth_uid ? minted.auth_uid : undefined; + if (authUid === undefined) return undefined; + // identityFromApiKey records `sub ?? auth_uid`, so a stored id may be either claim. + const mintedIds = new Set([authUid, ...(typeof minted?.sub === "string" && minted.sub ? [minted.sub] : [])]); + const rawEmail = typeof minted?.email === "string" ? minted.email.trim() : ""; + const email = rawEmail || undefined; + if (stored.accountId && !mintedIds.has(stored.accountId)) return undefined; + const storedEmail = normalizedEmail(stored.email); + if (storedEmail?.includes("@") && storedEmail !== normalizedEmail(email)) return undefined; + if (devinIdentityOwnedElsewhere(stored, currentAccountId, mintedIds, normalizedEmail(email))) return undefined; + return { ...credentialsFromApiKey(outcome.file.apiKey, apiBaseUrl, "local-cli"), accountId: authUid, ...(email ? { email } : {}) }; +} + +/** + * Every stored Devin row except the one being refreshed, across the alias-linked slots + * (the active row `getCredential` reads is one of these). Rows are skipped by id; only a + * caller that cannot name the row falls back to matching its key. + */ +function devinIdentityOwnedElsewhere( + stored: OAuthCredentials, + currentAccountId: string | undefined, + mintedIds: ReadonlySet, + email: string | undefined, +): boolean { + for (const slot of ["devin", ...devinAliasCredentialSlots("devin")]) { + for (const { id, credential } of listAccounts(slot)) { + if (currentAccountId !== undefined ? id === currentAccountId : credential.access === stored.access) continue; + if (credential.accountId !== undefined && mintedIds.has(credential.accountId)) return true; + if (email && normalizedEmail(credential.email) === email) return true; + } + } + return false; +} + export async function refreshDevinToken( _refreshToken: string, - _signal?: AbortSignal, - _credential?: OAuthCredentials, + signal?: AbortSignal, + credential?: OAuthCredentials, + accountId?: string, ): Promise { + const reread = credential?.source === "local-cli" + ? await rereadDevinCliCredential(credential, signal, accountId) : undefined; + if (reread) return reread; // Cognition has no refresh endpoint. Extending the stored expiry here is what // the carried implementation did, and it makes a revoked key look valid - // forever. Throwing lets the request path mark the account needsReauth the - // first time a forced refresh happens. + // forever. Throwing on the upstream-401 forced refresh is what marks the + // account needsReauth. throw new Error("invalid_grant: Devin API keys do not refresh. Run ocx login devin again."); } diff --git a/src/oauth/index.ts b/src/oauth/index.ts index f705bf43df2..c35caf31bde 100644 --- a/src/oauth/index.ts +++ b/src/oauth/index.ts @@ -191,6 +191,8 @@ interface OAuthProviderDef { refreshToken: string, signal?: AbortSignal, credential?: OAuthCredentials, + /** Store row being refreshed; passed by the generic lock only. */ + accountId?: string, ): Promise; /** provider entry written into config.json on first login. */ providerConfig: OcxProviderConfig; @@ -626,6 +628,7 @@ const FORCE_REFRESH_PROVIDERS = new Set([ "kiro", "google-antigravity", "orcarouter-oauth", + "devin", ]); export async function forceRefreshOAuthAccessSnapshot( @@ -834,7 +837,11 @@ function authoritative(stored:OAuthCredentials,active:boolean,now:()=>number):OA function merged(fresh: OAuthCredentials, previous: OAuthCredentials): OAuthCredentials { return { ...fresh, - source: previous.source === "local-cli" ? "oauth" : fresh.source ?? previous.source ?? "oauth", + // Shared: a refresh function returns "local-cli" only when the credential it hands back + // still is the local CLI's (Devin re-reading the CLI file, Meta Muse echoing its durable + // CLI key). Relabelling that "oauth" would stop the next forced refresh from re-reading it. + source: fresh.source === "local-cli" ? "local-cli" + : previous.source === "local-cli" ? "oauth" : fresh.source ?? previous.source ?? "oauth", ...(fresh.projectId === undefined && previous.projectId ? { projectId: previous.projectId } : {}), ...(fresh.apiBaseUrl === undefined && previous.apiBaseUrl ? { apiBaseUrl: previous.apiBaseUrl } : {}), ...(fresh.email === undefined && previous.email ? { email: previous.email } : {}), @@ -1035,7 +1042,7 @@ export async function refreshGenericAccountWithLock( } const generation = credentialGeneration(stored); try { - const fresh = merged(await def.refresh(stored.refresh, deps.signal, stored), stored); + const fresh = merged(await def.refresh(stored.refresh, deps.signal, stored, accountId), stored); const outcome = await mergeAccountCredential(provider, accountId, fresh, { expectedGeneration: generation, afterPrePersistRead: deps.afterPrePersistRead, diff --git a/src/oauth/kiro-terminal-failover.ts b/src/oauth/kiro-terminal-failover.ts index 968ac1e8b4d..e2a2dcafbb6 100644 --- a/src/oauth/kiro-terminal-failover.ts +++ b/src/oauth/kiro-terminal-failover.ts @@ -2,27 +2,41 @@ import type { OAuthAccessSnapshot } from "./index"; import type { OcxConfig } from "../types"; import { getValidAccessSnapshotForAccount } from "./index"; import { credentialGeneration, getAccountCredentialWithStatus, getAccountSet } from "./store"; -import { eligibleFailoverAccounts, isGenericOAuthFailoverEnabled, +import { eligibleFailoverAccounts, isGenericFailoverProvider, GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST } from "./generic-account-failover"; -/** Only an unchanged, allowlist-classified dead credential permits this alternate. */ -export async function tryKiroAlternateAfterTerminalRefresh( - config: OcxConfig, failedAccountId: string, failedGeneration: string, +/** + * Only an unchanged, allowlist-classified dead credential permits this alternate. + * + * Consent is the stored-login count, needsReauth rows included: the refused account was + * marked needsReauth a moment ago, and counting only healthy rows would make a two-account + * setup look like one exactly when the second account is needed. + */ +export async function tryAlternateAfterTerminalRefresh( + config: OcxConfig, providerName: string, failedAccountId: string, failedGeneration: string, ): Promise { - if (!isGenericOAuthFailoverEnabled(config, "kiro")) return null; - const failed = getAccountCredentialWithStatus("kiro", failedAccountId); + const provider = config.providers?.[providerName]; + if (!provider || !isGenericFailoverProvider(providerName, provider)) return null; + const order = getAccountSet(providerName)?.accounts.map(row => row.id) ?? []; + if (order.length < 2) return null; + const failed = getAccountCredentialWithStatus(providerName, failedAccountId); if (!failed?.needsReauth || credentialGeneration(failed.credential) !== failedGeneration) return null; - const order = getAccountSet("kiro")?.accounts.map(row => row.id) ?? []; const after = order.indexOf(failedAccountId); if (after < 0) return null; const ring = [...order.slice(after + 1), ...order.slice(0, after)]; - const eligible = new Set(eligibleFailoverAccounts("kiro")); + const eligible = new Set(eligibleFailoverAccounts(providerName)); let attempted = 0; for (const id of ring) { if (id === failedAccountId || !eligible.has(id)) continue; if (++attempted >= GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST) break; - try { return await getValidAccessSnapshotForAccount("kiro", id, { requireUsableAccount: true }); } + try { return await getValidAccessSnapshotForAccount(providerName, id, { requireUsableAccount: true }); } catch { /* Keep the original login-required result if every alternate is stale. */ } } return null; } + +export function tryKiroAlternateAfterTerminalRefresh( + config: OcxConfig, failedAccountId: string, failedGeneration: string, +): Promise { + return tryAlternateAfterTerminalRefresh(config, "kiro", failedAccountId, failedGeneration); +} diff --git a/src/server/responses/adapter-dispatch.ts b/src/server/responses/adapter-dispatch.ts index 1e7bec8baf2..8d4ab62e324 100644 --- a/src/server/responses/adapter-dispatch.ts +++ b/src/server/responses/adapter-dispatch.ts @@ -42,7 +42,7 @@ import { describeUpstreamConnectFailure } from "./upstream-error"; import type { OpaqueBlobRecoveryGuard } from "./core-opaque-recovery"; import type { AttemptRecoveryKind } from "../../usage/log"; import type { OAuthAccessSnapshot } from "../../oauth"; -import { OAuthLoginRequiredError, publicOAuthAuthenticationErrorMessage } from "../../oauth"; +import { publicOAuthAuthenticationErrorMessage } from "../../oauth"; import { tryKiroAlternateAfterTerminalRefresh } from "../../oauth/kiro-terminal-failover"; import { classifyKiroRefusal } from "../../adapters/kiro-refusal"; import { normalizeFinalKiroHttpError } from "../../adapters/kiro-retry"; @@ -610,7 +610,7 @@ export async function prepareAdapterExchange( refreshed = await refreshResolvedOAuthSelection(transportState.sentOAuthSnapshot); } catch (err) { const failed = transportState.sentOAuthSnapshot; - if (route.providerName === "kiro" && err instanceof OAuthLoginRequiredError && failed + if (route.providerName === "kiro" && failed && transportState.genericFailovers < transportState.genericFailoverLimit) { const alternate = await tryKiroAlternateAfterTerminalRefresh(config, failed.accountId, failed.generation); if (alternate) { diff --git a/src/server/responses/request-transport.ts b/src/server/responses/request-transport.ts index b9a3de849dd..34d3e69cc35 100644 --- a/src/server/responses/request-transport.ts +++ b/src/server/responses/request-transport.ts @@ -117,6 +117,8 @@ export async function prepareResponsesTransport( || route.providerName === "kiro" || route.providerName === "google-antigravity" || route.providerName === "orcarouter-oauth" + // runTurn transport: the replay runs on the first-event preflight in run-turn-execution. + || route.providerName === "devin" ) && route.provider.authMode === "oauth"; let sentOAuthSnapshot: OAuthAccessSnapshot | undefined; let replayOAuthCredentialSnapshot: Pick | undefined; @@ -193,9 +195,10 @@ export async function prepareResponsesTransport( }; const refreshResolvedOAuthSelection = async (sent: OAuthAccessSnapshot): Promise => { const current = captureOAuthAccountSelection(route.providerName); - const unchanged = current?.accountId === oauthSelection?.accountId - && current?.revision === oauthSelection?.revision; - const candidate = unchanged ? await forceRefreshOAuthAccessSnapshot(sent) : sent; + // Keyed on the account, not the selection revision: a selection that moved away and back + // (A -> B -> A) before the 401 landed still serves the rejected credential, and skipping + // the refresh would replay it and spend the one recovery attempt. + const candidate = current?.accountId === sent.accountId ? await forceRefreshOAuthAccessSnapshot(sent) : sent; const admitted = await commitResolvedOAuthSelection(candidate); if (!admitted) throw new Error("OAuth selection changed during credential recovery"); if (kiroLoadEnabled && options.accountLoad?.lease?.accountId !== admitted.accountId) { diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index 736be252156..2d84f3fca16 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -27,6 +27,8 @@ import { rotateGenericOAuthAccountOn429, failoverAccountSnapshot, } from "../../oauth/generic-account-failover"; +import { publicOAuthAuthenticationErrorMessage, type OAuthAccessSnapshot } from "../../oauth/index"; +import { tryAlternateAfterTerminalRefresh } from "../../oauth/kiro-terminal-failover"; import { resolveWireProtocolOverride } from "../adapter-resolve"; import { formatErrorResponse, bridgeToResponsesSSE, buildResponseJSON } from "../../bridge"; import { redactSecretString } from "../../lib/redact"; @@ -93,6 +95,9 @@ export async function executeResponsesRunTurn( | "genericFailovers" | "genericFailoverLimit" | "applyFailoverSnapshot" + | "refreshResolvedOAuthSelection" + | "isOAuth401ReplayProvider" + | "sentOAuthSnapshot" | "resolveSelectionAdapter" | "adapter" | "noteRoutedAttemptSend" @@ -122,6 +127,7 @@ export async function executeResponsesRunTurn( adapterBindings, refreshRunTurnAdapter, applyFailoverSnapshot, + refreshResolvedOAuthSelection, resolveSelectionAdapter, } = transportState; const { @@ -348,6 +354,92 @@ export async function executeResponsesRunTurn( yield* await preflightRunTurnFailover(stream, iterParsed); })(); }; + // Rebind the turn to an admitted account. The failed attempt emitted no client-visible bytes, + // so replay is safe, but a Cursor conversation/checkpoint is credential-scoped: carrying its + // account identity into the next account would not be. The rotated adapter derives its own. + const adoptRunTurnAccount = (admittedSnapshot: Pick): boolean => { + parsed._cursorIdentityScope = undefined; + parsed._cursorConversationId = undefined; + if (parsed._providerContinuation?.cursor) { + const { cursor: _discardedCursor, ...otherProviderState } = parsed._providerContinuation; + parsed._providerContinuation = otherProviderState; + } + const rotatedProvider = resolveWireProtocolOverride( + route.providerName, + route.modelId, + route.provider, + inboundWire, + route.staticPolicy, + ); + const rotatedAdapter = resolveSelectionAdapter(rotatedProvider, config.cacheRetention); + if (!rotatedAdapter.runTurn) return false; + transportState.runTurnAdapter = rotatedAdapter; + bindRouteReasoningReplayScope({ + parsed, + providerName: route.providerName, + provider: rotatedProvider, + adapterName: rotatedAdapter.name, + oauthCredentialSnapshot: { + accountId: admittedSnapshot.accountId, + generation: admittedSnapshot.generation, + }, + codexAuthContext: admissionState.authCtx, + forwardHeaders: requestState.selectedForwardHeaders, + }); + sealRequestAttemptIdentity(logCtx.activeAttempt, logCtx.provider, rotatedAdapter.name, logCtx.accountLogLabel); + recordAttemptCredentialSource(logCtx.activeAttempt, route.providerName, route.provider, rotatedAdapter.name); + return true; + }; + // The runTurn half of the upstream-401 replay adapter-dispatch runs for HTTP transports: + // force-refresh the credential that was sent, once per request. A refresh that cannot + // succeed marks the account needsReauth inside the OAuth owner, so the account stops being + // selected; the turn then moves to a surviving stored account, or, with none, the client + // gets the login instruction instead of an opaque upstream 401. + let oauth401ReplayAttempted = false; + const recoverRunTurnAdapterOnPreflight401 = async ( + error: Extract, + ): Promise => { + const sent = transportState.sentOAuthSnapshot; + if (error.status !== 401 || !transportState.isOAuth401ReplayProvider || !sent || oauth401ReplayAttempted) return false; + oauth401ReplayAttempted = true; + const hop = reserveCredentialHop("auth-recovery", `${route.providerName}|${route.modelId}|runturn-oauth-401`); + if (!hop.allowed) return false; + try { + let admitted: OAuthAccessSnapshot | null; + try { + admitted = await applyFailoverSnapshot(await refreshResolvedOAuthSelection(sent)); + } catch (err) { + // Not only OAuthLoginRequiredError: a concurrent request that already flagged this + // account and moved the selection makes this refresh fail as "selection changed". The + // helper still requires the sent generation to be flagged needsReauth, so it is safe. + const alternate = transportState.genericFailovers < transportState.genericFailoverLimit + ? await tryAlternateAfterTerminalRefresh(config, route.providerName, sent.accountId, sent.generation) + : null; + admitted = alternate ? await applyFailoverSnapshot(alternate) : null; + if (admitted) transportState.genericFailovers += 1; + else Object.assign(error, { errorType: "authentication_error", message: publicOAuthAuthenticationErrorMessage(err) }); + } + if (!admitted || !adoptRunTurnAccount(admitted)) { + hop.permit?.release(); + return false; + } + sendBudgetState.pendingHopPermit = hop.permit; + return true; + } catch (err) { + // Applying the alternate can still fail, e.g. when a newer manual selection points back + // at the flagged account; the client then needs the login instruction, not the raw 401. + Object.assign(error, { errorType: "authentication_error", message: publicOAuthAuthenticationErrorMessage(err) }); + hop.permit?.release(); + return false; + } + }; + const recoverRunTurnAdapterOnPreflightError = async ( + error: Extract, + ): Promise => { + if (await recoverRunTurnAdapterOnPreflight401(error)) return "oauth-401"; + if (await rotateRunTurnAdapterOnPreflight429(error)) return "oauth-account-429"; + return undefined; + }; const rotateRunTurnAdapterOnPreflight429 = async ( error: Extract, ): Promise => { @@ -400,42 +492,10 @@ export async function executeResponsesRunTurn( hop.permit?.release(); return false; } - // A Cursor conversation/checkpoint is credential-scoped. The failed attempt emitted no - // client-visible bytes, so replay is safe, but carrying its account identity into the next - // account would not be. Let the rotated adapter derive a fresh identity and conversation. - parsed._cursorIdentityScope = undefined; - parsed._cursorConversationId = undefined; - if (parsed._providerContinuation?.cursor) { - const { cursor: _discardedCursor, ...otherProviderState } = parsed._providerContinuation; - parsed._providerContinuation = otherProviderState; - } - const rotatedProvider = resolveWireProtocolOverride( - route.providerName, - route.modelId, - route.provider, - inboundWire, - route.staticPolicy, - ); - const rotatedAdapter = resolveSelectionAdapter(rotatedProvider, config.cacheRetention); - if (!rotatedAdapter.runTurn) { + if (!adoptRunTurnAccount(admittedSnapshot)) { hop.permit?.release(); return false; } - transportState.runTurnAdapter = rotatedAdapter; - bindRouteReasoningReplayScope({ - parsed, - providerName: route.providerName, - provider: rotatedProvider, - adapterName: rotatedAdapter.name, - oauthCredentialSnapshot: { - accountId: admittedSnapshot.accountId, - generation: admittedSnapshot.generation, - }, - codexAuthContext: admissionState.authCtx, - forwardHeaders: requestState.selectedForwardHeaders, - }); - sealRequestAttemptIdentity(logCtx.activeAttempt, logCtx.provider, rotatedAdapter.name, logCtx.accountLogLabel); - recordAttemptCredentialSource(logCtx.activeAttempt, route.providerName, route.provider, rotatedAdapter.name); // The caller replays the turn on this rotation, and a runTurn adapter dispatches through // its own reservation ladder -- Cursor reserves once per physical send. Confirming here // would leave that ladder to charge the same replay a second time (#4709), so hand the @@ -463,13 +523,14 @@ export async function executeResponsesRunTurn( yield event; continue; } - if (!firstMeaningfulSeen && !replayUnsafe && event.type === "error" - && await rotateRunTurnAdapterOnPreflight429(event)) { + const recovery = !firstMeaningfulSeen && !replayUnsafe && event.type === "error" + ? await recoverRunTurnAdapterOnPreflightError(event) : undefined; + if (recovery) { const retryQueue = createAdapterEventQueue({ onBacklogExceeded: () => runTurnAbort.abort(), }); const pendingPermit = sendBudgetState.pendingHopPermit; - const retryAttempt = runTurnAttempt(retryQueue, "oauth-account-429", false, replayParsed); + const retryAttempt = runTurnAttempt(retryQueue, recovery, false, replayParsed); if (pendingPermit) { const releaseIfUnclaimed = () => { if (sendBudgetState.pendingHopPermit !== pendingPermit) return; @@ -522,15 +583,13 @@ export async function executeResponsesRunTurn( return streamAfterPreflight(preflight.stream, replayParsed, preflight.replayUnsafe); } if (preflight.ready) return streamAfterPreflight(preflight.stream, replayParsed, preflight.replayUnsafe); - if (preflight.replayUnsafe - || !preflight.error - || !(await rotateRunTurnAdapterOnPreflight429(preflight.error))) { - return preflight.stream; - } + const recovery = preflight.replayUnsafe || !preflight.error + ? undefined : await recoverRunTurnAdapterOnPreflightError(preflight.error); + if (!recovery) return preflight.stream; const retryQueue = createAdapterEventQueue({ onBacklogExceeded: () => runTurnAbort.abort(), }); - latestRetryAttempt = runTurnAttempt(retryQueue, "oauth-account-429", false, replayParsed); + latestRetryAttempt = runTurnAttempt(retryQueue, recovery, false, replayParsed); void latestRetryAttempt; source = retryQueue.stream(); } diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index 0d5d4a310fd..dc3df67508c 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -246,6 +246,22 @@ surface the pre-output 429 immediately, because the outer response cannot forwar heartbeats while it is choosing a target. An earlier replay-unsafe heartbeat or meaningful output keeps the failure on the current target. +## runTurn pre-output 401 replay + +`src/server/responses/run-turn-execution.ts` runs the `adapter-dispatch.ts` OAuth 401 replay on the +runTurn first-event preflight for `isOAuth401ReplayProvider` routes (Devin is the runTurn member): a +structured 401 before output or a replay-unsafe heartbeat force-refreshes the sent credential once +per request under an `auth-recovery` hop and replays the turn. A terminal refresh has already marked +the account needsReauth; the turn moves to a surviving stored account via +`tryAlternateAfterTerminalRefresh`, else the 401 carries the login instruction. Devin quota +`permission_denied` maps to 429 and plain `permission_denied` to 403, so neither refreshes. +By design a single upstream 401 on an `oauth`-source Devin account marks it needsReauth: Cognition has +no refresh endpoint and there is no confirming probe. A `local-cli` account instead re-reads the CLI +file and adopts a different key only if its host passes `validateDevinApiBaseUrl`, the key mints a +user_jwt (its auth_uid/email are the identity; the session token has none), that identity does not +contradict the slot's, and no other slot owns the key or identity; the adopted identity is recorded. An unreadable or half-written file, or a mint failing other than 401/403, fails the refresh unflagged. +Test: `tests/responses/responses-devin-401-replay.test.ts`. + ## Optional client transport hints `dropCodexSafetyBuffering` defaults to false. Canonical OpenAI forward Responses can remove only diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 01768b541df..4423318b3ea 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -46,6 +46,7 @@ "start-args.test.ts": "cli", "start-ownership-publication.test.ts": "cli", "responses-core-modules.test.ts": "responses", + "responses-devin-401-replay.test.ts": "responses", "responses-grok-devin-preflight.test.ts": "responses", "responses-passthrough-transient-policy.test.ts": "responses", "responses-spend-ledger-wiring.test.ts": "responses", diff --git a/tests/responses/responses-devin-401-replay.test.ts b/tests/responses/responses-devin-401-replay.test.ts new file mode 100644 index 00000000000..a35b47fbea0 --- /dev/null +++ b/tests/responses/responses-devin-401-replay.test.ts @@ -0,0 +1,379 @@ +import { afterAll, afterEach, beforeEach, expect, mock, test } from "bun:test"; +import { writeFileSync } from "node:fs"; +import type { ProviderAdapter } from "../../src/adapters/base"; +import type { AdapterEvent, OcxConfig, OcxProviderConfig } from "../../src/types"; +import { getAccountSet, saveCredential, setActiveAccount } from "../../src/oauth/store"; +import { clearGenericFailoverHealth } from "../../src/oauth/generic-account-failover"; +import { DEVIN_CLI_CREDENTIALS_ENV } from "../../src/oauth/devin/cli-import"; +import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; +import { createTempHome } from "../helpers/temp-home"; + +const DEAD = "devin-session-token$synthetic-dead"; +const LIVE = "devin-session-token$synthetic-live"; +const ROTATED = "devin-session-token$synthetic-rotated"; +const originalFetch = globalThis.fetch; + +// GetUserJwt stand-in: the identity a key mints, keyed by the key inside the protobuf body. +let mintedIdentity: Record = {}; +let mintCalls = 0; +let mintFailure: (() => Response) | undefined; +function fakeUserJwt(payload: object): string { + const part = (value: object) => Buffer.from(JSON.stringify(value)).toString("base64url"); + return `${part({ alg: "HS256", typ: "JWT" })}.${part({ ...payload, exp: 9_999_999_999 })}.c2lnbmF0dXJl`; +} +const mintFetch = (async (input: Parameters[0], init?: RequestInit) => { + const url = String(input instanceof Request ? input.url : input); + if (!url.endsWith("/exa.auth_pb.AuthService/GetUserJwt")) return originalFetch(input, init); + mintCalls++; + if (mintFailure) return mintFailure(); + const body = Buffer.from(init?.body as Uint8Array).toString("latin1"); + const key = Object.keys(mintedIdentity).find(candidate => body.includes(candidate)); + const identity = key ? mintedIdentity[key] : undefined; + if (!identity) return new Response("", { status: 401 }); + const jwt = Buffer.from(fakeUserJwt(identity)); + return new Response(Buffer.concat([Buffer.from([0x0a, jwt.length & 0x7f | 0x80, jwt.length >> 7]), jwt]), + { status: 200, headers: { "content-type": "application/proto" } }); +}) as typeof fetch; + +const resolver = await import("../../src/server/adapter-resolve"); +const originalResolve = resolver.resolveAdapter; +const originalResolverModule = { ...resolver }; +let sentKeys: string[] = []; +let rateLimited = false; +let holdDeadSend: (() => Promise) | undefined; +mock.module("../../src/server/adapter-resolve", () => ({ ...resolver, + resolveAdapter(provider: OcxProviderConfig, cache?: "none" | "short" | "long") { + if (provider.adapter !== "devin") return originalResolve(provider, cache); + return { + name: "devin", + buildRequest: () => ({ url: provider.baseUrl, method: "POST", headers: {}, body: "" }), + async *parseStream() { yield { type: "done" } as AdapterEvent; }, + async runTurn(_parsed, _incoming, emit) { + const key = String(provider.apiKey); + sentKeys.push(key); + if (rateLimited) { + emit({ type: "error", status: 429, errorType: "rate_limit_error", code: "resource_exhausted", + retryable: true, message: "Cognition chat failed (resource_exhausted)" }); + return; + } + if (key === DEAD) { + await holdDeadSend?.(); + emit({ type: "error", status: 401, errorType: "authentication_error", code: "unauthenticated", + retryable: false, message: "Devin cloud error unauthenticated: invalid api key" }); + return; + } + emit({ type: "text_delta", text: `served by ${key === ROTATED ? "rotated" : "live"}` }); + emit({ type: "done" }); + }, + } satisfies ProviderAdapter; + }, +})); +const { handleResponses } = await import("../../src/server/responses"); + +let home: ReturnType; +let release: (() => void) | undefined; +let previousCliPath: string | undefined; + +beforeEach(() => { + home = createTempHome("ocx-devin-401-replay-"); + release = acquireOwnedSpendHome(); + clearGenericFailoverHealth(); + sentKeys = []; + rateLimited = false; + holdDeadSend = undefined; + mintedIdentity = { [ROTATED]: { auth_uid: "uid-rotated", email: "rotated@example.com" } }; + mintCalls = 0; + mintFailure = undefined; + globalThis.fetch = mintFetch; + previousCliPath = process.env[DEVIN_CLI_CREDENTIALS_ENV]; + // Never let a test read the developer's real CLI credential. + process.env[DEVIN_CLI_CREDENTIALS_ENV] = home.path("devin-credentials.toml"); +}); +afterAll(() => { + mock.module("../../src/server/adapter-resolve", () => originalResolverModule); +}); +afterEach(() => { + try { + release?.(); + } finally { + if (previousCliPath === undefined) delete process.env[DEVIN_CLI_CREDENTIALS_ENV]; + else process.env[DEVIN_CLI_CREDENTIALS_ENV] = previousCliPath; + globalThis.fetch = originalFetch; + clearGenericFailoverHealth(); + home.remove(); + } +}); + +function writeCliFile(apiKey: string, apiServerUrl = "https://server.codeium.com"): void { + writeFileSync(home.path("devin-credentials.toml"), + `windsurf_api_key = "${apiKey}"\napi_server_url = "${apiServerUrl}"\n`); +} + +async function saveDevin(access: string, accountId: string, source: "oauth" | "local-cli" = "oauth") { + await saveCredential("devin", { + access, refresh: access, expires: Number.MAX_SAFE_INTEGER, accountId, source, + apiBaseUrl: "https://server.codeium.com", + }); +} + +// A CLI import records no identity: the session token it copies carries only a session_id. +async function saveCliImport(access: string, extra: { accountId?: string; email?: string } = {}) { + await saveCredential("devin", { + access, refresh: access, expires: Number.MAX_SAFE_INTEGER, source: "local-cli", + apiBaseUrl: "https://server.codeium.com", ...extra, + }, { preserveIdentityless: true }); +} + +function cliAccount() { + return getAccountSet("devin")?.accounts.find(row => row.credential.source === "local-cli"); +} + +function run(stream = false) { + const config = { + port: 0, defaultProvider: "devin", + providers: { devin: { adapter: "devin", authMode: "oauth", baseUrl: "https://server.codeium.com", models: ["swe-1-6"] } }, + } as OcxConfig; + return handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "devin/swe-1-6", input: "answer", stream }), + }), config, { model: "", provider: "", surface: "codex" }); +} + +function account(id: string) { + return getAccountSet("devin")?.accounts.find(row => row.credential.accountId === id); +} + +test.each([false, true])("a revoked key is marked needsReauth and the turn fails over (stream=%s)", async stream => { + await saveDevin(LIVE, "spare"); + await saveDevin(DEAD, "revoked"); + expect(getAccountSet("devin")?.activeAccountId).toBe(account("revoked")?.id); + + const response = await run(stream); + const body = await response.text(); + + expect(response.status).toBe(200); + expect(body).toContain("served by live"); + expect(sentKeys).toEqual([DEAD, LIVE]); + expect(account("revoked")?.needsReauth).toBe(true); + expect(getAccountSet("devin")?.activeAccountId).toBe(account("spare")?.id); + + // The dead account is not reselected by the next request. + sentKeys = []; + expect((await run()).status).toBe(200); + expect(sentKeys).toEqual([LIVE]); +}); + +test("a lone revoked account surfaces the login instruction", async () => { + await saveDevin(DEAD, "revoked"); + + // A buffered runTurn failure is a `status: "failed"` Response object, not an HTTP error. + const body = await (await run()).json() as { status: string; error: { type: string; message: string } }; + + expect(body.status).toBe("failed"); + expect(body.error.type).toBe("authentication_error"); + expect(body.error.message).toBe("Not logged in to devin. Run: ocx login devin"); + expect(sentKeys).toEqual([DEAD]); + expect(account("revoked")?.needsReauth).toBe(true); + + // Later turns fail fast on the flagged account instead of re-sending a dead key. + sentKeys = []; + const next = await run(); + expect(next.status).toBe(401); + expect(await next.text()).toContain("ocx login devin"); + expect(sentKeys).toEqual([]); +}); + +test("a CLI-imported account adopts the key a later `devin auth login` wrote", async () => { + await saveCliImport(DEAD); + writeCliFile(ROTATED); + + const response = await run(); + + expect(response.status).toBe(200); + expect(await response.text()).toContain("served by rotated"); + expect(sentKeys).toEqual([DEAD, ROTATED]); + const row = cliAccount(); + expect(row?.needsReauth).not.toBe(true); + expect(row?.credential.access).toBe(ROTATED); + // The minted identity is recorded, so the next adoption for this slot is strict. + expect(row?.credential.accountId).toBe("uid-rotated"); + expect(row?.credential.email).toBe("rotated@example.com"); + expect(row?.credential.source).toBe("local-cli"); + expect(row?.credential.expires).toBe(Number.MAX_SAFE_INTEGER); +}); + +test.each([ + ["unchanged", () => writeCliFile(DEAD)], + ["missing", () => {}], + ["off-allowlist host", () => writeCliFile(ROTATED, "https://attacker.example")], +])("a CLI-imported account with a %s credential file needs reauth", async (_label, arrange) => { + await saveCliImport(DEAD); + arrange(); + + const body = await (await run()).json() as { error: { message: string } }; + + expect(body.error.message).toBe("Not logged in to devin. Run: ocx login devin"); + expect(sentKeys).toEqual([DEAD]); + expect(cliAccount()?.needsReauth).toBe(true); + expect(cliAccount()?.credential.access).toBe(DEAD); + // The key is only ever sent to an allowlisted host, including the identity probe. + expect(mintCalls).toBe(0); +}); + +test("a CLI key another stored account already owns is not adopted", async () => { + await saveDevin(ROTATED, "other"); + await saveCliImport(DEAD); + writeCliFile(ROTATED); + + await (await run()).text(); + + expect(cliAccount()?.needsReauth).toBe(true); + expect(cliAccount()?.credential.access).toBe(DEAD); +}); + +test.each([ + ["a key Cognition refuses with 401", async () => { mintedIdentity = {}; await saveCliImport(DEAD); }], + ["a key Cognition refuses with 403", async () => { + mintFailure = () => new Response("", { status: 403 }); + await saveCliImport(DEAD); + }], + ["a minted token without auth_uid", async () => { + mintedIdentity = { [ROTATED]: { email: "rotated@example.com" } }; + await saveCliImport(DEAD); + }], + ["a slot whose recorded accountId differs", async () => { await saveCliImport(DEAD, { accountId: "uid-before" }); }], + ["a slot whose recorded email differs", async () => { await saveCliImport(DEAD, { email: "before@example.com" }); }], + ["an identity another stored account owns", async () => { + await saveDevin(LIVE, "uid-rotated"); + await saveCliImport(DEAD); + }], +])("the CLI key is not adopted for %s", async (_label, arrange) => { + await arrange(); + writeCliFile(ROTATED); + + await (await run()).text(); + + expect(sentKeys).not.toContain(ROTATED); + expect(cliAccount()?.needsReauth).toBe(true); + expect(cliAccount()?.credential.access).toBe(DEAD); +}); + +test.each([ + ["a transient 503", () => new Response("", { status: 503 })], + ["a 429", () => new Response("", { status: 429 })], + ["a network failure", () => { throw new TypeError("fetch failed"); }], +])("an identity probe that fails with %s does not flag the account", async (_label, failure) => { + await saveCliImport(DEAD); + writeCliFile(ROTATED); + mintFailure = failure; + + const body = await (await run()).json() as { error: { type: string; message: string } }; + + expect(body.error.type).toBe("authentication_error"); + expect(body.error.message).not.toContain("ocx login devin"); + expect(cliAccount()?.needsReauth).not.toBe(true); + expect(cliAccount()?.credential.access).toBe(DEAD); + + // Once the probe recovers, the next 401 adopts the rotated key. + mintFailure = undefined; + expect(await (await run()).text()).toContain("served by rotated"); + expect(cliAccount()?.credential.access).toBe(ROTATED); +}); + +test("a credential file caught mid-write does not flag the account", async () => { + await saveCliImport(DEAD); + // Half-written by `devin auth login`: the key line is there, the server line is not yet. + writeFileSync(home.path("devin-credentials.toml"), `windsurf_api_key = "${ROTATED}"\n`); + + const body = await (await run()).json() as { error: { type: string; message: string } }; + + expect(body.error.type).toBe("authentication_error"); + expect(body.error.message).not.toContain("ocx login devin"); + expect(cliAccount()?.needsReauth).not.toBe(true); + expect(mintCalls).toBe(0); + + // Once the write completes, the next 401 adopts the rotated key. + writeCliFile(ROTATED); + expect(await (await run()).text()).toContain("served by rotated"); + expect(cliAccount()?.credential.access).toBe(ROTATED); +}); + +test("a slot recorded from the key's `sub` claim still matches the minted identity", async () => { + mintedIdentity = { [ROTATED]: { auth_uid: "uid-rotated", sub: "sub-rotated", email: "rotated@example.com" } }; + await saveCliImport(DEAD, { accountId: "sub-rotated" }); + writeCliFile(ROTATED); + + expect(await (await run()).text()).toContain("served by rotated"); + expect(cliAccount()?.credential.accountId).toBe("uid-rotated"); +}); + +test("a slot whose recorded identity matches adopts the rotated key", async () => { + await saveCliImport(DEAD, { accountId: "uid-rotated", email: " Rotated@Example.com " }); + writeCliFile(ROTATED); + + expect(await (await run()).text()).toContain("served by rotated"); + expect(cliAccount()?.credential.access).toBe(ROTATED); + expect(mintCalls).toBe(1); +}); + +test("a turn whose 401 lands after another turn already failed the account over still fails over", async () => { + await saveDevin(LIVE, "spare"); + await saveDevin(DEAD, "revoked"); + // Both turns send on the revoked account; the second 401 only arrives once the first turn has + // flagged it and moved the selection, so the second refresh sees a changed selection. + const bothSent = Promise.withResolvers(); + const firstDone = Promise.withResolvers(); + let deadSends = 0; + holdDeadSend = async () => { + const order = ++deadSends; + if (order === 2) bothSent.resolve(); + await bothSent.promise; + if (order === 2) await firstDone.promise; + }; + + const first = run().then(async response => { + const text = await response.text(); + firstDone.resolve(); + return text; + }); + const second = run().then(response => response.text()); + const [firstBody, secondBody] = await Promise.all([first, second]); + + expect(deadSends).toBe(2); + expect(firstBody).toContain("served by live"); + expect(secondBody).toContain("served by live"); + expect(account("revoked")?.needsReauth).toBe(true); + expect(getAccountSet("devin")?.activeAccountId).toBe(account("spare")?.id); +}); + +test("a 429 is not treated as an authentication failure", async () => { + await saveDevin(LIVE, "limited"); + rateLimited = true; + + const body = await (await run()).json() as { error: { type: string } }; + + expect(body.error.type).toBe("rate_limit_error"); + expect(sentKeys).toEqual([LIVE]); + expect(account("limited")?.needsReauth).not.toBe(true); +}); + +test("a selection that leaves the revoked account and returns to it before the 401 still refreshes it", async () => { + await saveDevin(LIVE, "spare"); + await saveDevin(DEAD, "revoked"); + const revoked = account("revoked")!.id; + // Same account id at 401 time, newer selection revision: the refresh must still run, or the + // rejected key is replayed and the one allowed recovery is spent on it. + holdDeadSend = async () => { + holdDeadSend = undefined; + await setActiveAccount("devin", account("spare")!.id); + await setActiveAccount("devin", revoked); + }; + + const body = await (await run()).json() as { error: { message: string } }; + + expect(sentKeys.filter(key => key === DEAD)).toHaveLength(1); + expect(account("revoked")?.needsReauth).toBe(true); + // The operator's newer manual selection names the revoked account, so no automatic move + // overrides it; the client is told to log in rather than shown the raw upstream 401. + expect(body.error.message).toBe("Not logged in to devin. Run: ocx login devin"); +});