diff --git a/.changeset/single-flight-anthropic-refresh.md b/.changeset/single-flight-anthropic-refresh.md new file mode 100644 index 0000000000..23fc38d23b --- /dev/null +++ b/.changeset/single-flight-anthropic-refresh.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Prevent concurrent tasks from falling back when an Anthropic OAuth token rotates. +category: fix +dev: Serializes Anthropic refresh-token rotation across auth storage instances and Fusion processes. diff --git a/packages/engine/src/__tests__/auth-storage-concurrency.test.ts b/packages/engine/src/__tests__/auth-storage-concurrency.test.ts index f502ca5f7d..6a9f20304a 100644 --- a/packages/engine/src/__tests__/auth-storage-concurrency.test.ts +++ b/packages/engine/src/__tests__/auth-storage-concurrency.test.ts @@ -279,4 +279,89 @@ describe("createFusionAuthStorage — concurrent cross-process coordination", () const instanceC = createFusionAuthStorage(); expect(await instanceC.getApiKey("anthropic")).not.toBe("refreshed-from-stale-old-token"); }); + + it("single-flights a rotating Anthropic refresh token across auth storage instances", async () => { + const now = Date.now(); + const instanceA = createFusionAuthStorage(); + await instanceA.set("anthropic-subscription", { + type: "oauth", + access: "expiring-access", + refresh: "single-use-refresh", + expires: now + 1_000, + }); + + // Each agent session constructs its own auth storage instance. Anthropic refresh + // tokens rotate, so concurrent refresh requests using the same token cannot both + // succeed: the second request observes an already-consumed refresh token. + const instanceB = createFusionAuthStorage(); + let refreshConsumed = false; + const fetchMock = vi.fn(async () => { + if (refreshConsumed) { + return { + ok: false, + text: async () => JSON.stringify({ error: "invalid_grant" }), + }; + } + refreshConsumed = true; + await new Promise((resolve) => setTimeout(resolve, 20)); + return { + ok: true, + text: async () => JSON.stringify({ + access_token: "rotated-access", + refresh_token: "next-single-use-refresh", + expires_in: 3600, + }), + }; + }); + globalThis.fetch = fetchMock as unknown as typeof fetch; + + const [keyA, keyB] = await Promise.all([ + instanceA.getApiKey("anthropic"), + instanceB.getApiKey("anthropic"), + ]); + + expect(keyA).toBe("rotated-access"); + expect(keyB).toBe("rotated-access"); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(readAuthFile(homeDir)["anthropic-subscription"]).toMatchObject({ + type: "oauth", + access: "rotated-access", + refresh: "next-single-use-refresh", + }); + }); + + it("single-flights one rotating token across legacy and subscription Anthropic aliases", async () => { + const instanceA = createFusionAuthStorage(); + await instanceA.set("anthropic", { + type: "oauth", + access: "legacy-expiring-access", + refresh: "shared-alias-refresh", + expires: Date.now() + 1_000, + }); + const instanceB = createFusionAuthStorage(); + + let refreshConsumed = false; + const fetchMock = vi.fn(async () => { + if (refreshConsumed) { + return { ok: false, text: async () => JSON.stringify({ error: "invalid_grant" }) }; + } + refreshConsumed = true; + await new Promise((resolve) => setTimeout(resolve, 20)); + return { + ok: true, + text: async () => JSON.stringify({ + access_token: "alias-rotated-access", + refresh_token: "next-alias-refresh", + expires_in: 3600, + }), + }; + }); + globalThis.fetch = fetchMock as unknown as typeof fetch; + + await expect(Promise.all([ + instanceA.getApiKey("anthropic"), + instanceB.getApiKey("anthropic-subscription"), + ])).resolves.toEqual(["alias-rotated-access", "alias-rotated-access"]); + expect(fetchMock).toHaveBeenCalledTimes(1); + }); }); diff --git a/packages/engine/src/auth-storage.ts b/packages/engine/src/auth-storage.ts index a24cc7c991..2b7950fc40 100644 --- a/packages/engine/src/auth-storage.ts +++ b/packages/engine/src/auth-storage.ts @@ -68,6 +68,23 @@ function enqueueAuthWrite(authPath: string, write: () => Promise): Promise return operation; } +async function withOAuthRefreshLock( + authPath: string, + providerId: string, + refresh: () => Promise, +): Promise { + const safeProviderId = providerId.replace(/[^a-zA-Z0-9._-]/g, "_"); + // Use a distinct lock target, not authPath with a custom lockfilePath. proper-lockfile + // tracks held locks by target path internally, so two lock domains sharing authPath + // would corrupt each other's renewal/release bookkeeping. + const release = await lockfile.lock(`${authPath}.${safeProviderId}.refresh`, AUTH_LOCK_OPTIONS); + try { + return await refresh(); + } finally { + await release(); + } +} + class FusionFileAuthStorage implements FusionAuthStorage { private data: Record = {}; private modelRuntime: ModelRuntime | undefined; @@ -452,7 +469,8 @@ export async function createFusionModelRegistry(authStorage: FusionAuthStorage, } export function createFusionAuthStorage(): FusionAuthStorage { - const primary = new FusionFileAuthStorage(getFusionAuthPath()); + const authPath = getFusionAuthPath(); + const primary = new FusionFileAuthStorage(authPath); let supplementalCredentials = readSupplementalCredentials(); // models.json provider API keys — final fallback after primary auth and supplemental auth.json files let modelsJsonApiKeys = readModelsJsonApiKeys(); @@ -538,9 +556,67 @@ export function createFusionAuthStorage(): FusionAuthStorage { return existing; } - const refreshPromise = refreshOAuthCredential(storageProvider, credential) + /* + FNXC:ClaudeOAuth 2026-07-18-16:50: + Anthropic refresh tokens rotate on use. The in-memory single-flight map above is + scoped to one createFusionAuthStorage() instance, while every agent session creates + its own instance and separate Fusion processes share the same auth.json. Refreshing + outside the file lock therefore let a burst of sessions submit the same refresh + token concurrently: one request rotated it and the losers fell back to another + provider with stale auth. + + Hold a provider-specific cross-process refresh lock for the complete + read/refresh/write transaction. It deliberately differs from the auth-file write + lock so a manual login can persist while the network request is in flight. A waiter + re-reads the provider row after acquiring the refresh lock; if another session + already rotated it, the fresh credential is returned without a second request. + Re-read again before persistence so a concurrent manual login always wins. + */ + // Anthropic's legacy `anthropic` row and separated `anthropic-subscription` + // row can refer to the same rotating refresh token, so both aliases must use + // one canonical refresh-lock domain. + const refreshLockProvider = getOAuthResolutionProviderId(storageProvider); + const refreshPromise = withOAuthRefreshLock(authPath, refreshLockProvider, async () => { + primary.reload(); + const selectPersistedRefreshCredential = (): StoredCredential | undefined => { + const storedCredential = primary.get(storageProvider) as StoredCredential | undefined; + if (refreshLockProvider !== ANTHROPIC_PROVIDER_ID) { + return storedCredential; + } + + const legacyCredential = primary.get(ANTHROPIC_PROVIDER_ID) as StoredCredential | undefined; + const subscriptionCredential = primary.get(ANTHROPIC_SUBSCRIPTION_PROVIDER_ID) as StoredCredential | undefined; + return choosePreferredStoredCredential( + legacyCredential?.type === "oauth" ? legacyCredential : undefined, + subscriptionCredential?.type === "oauth" ? subscriptionCredential : undefined, + ); + }; + const initialPersistedCredential = selectPersistedRefreshCredential(); + const refreshCandidate = choosePreferredStoredCredential(initialPersistedCredential, credential) ?? credential; + + if (!shouldRefreshOAuthCredential(refreshCandidate)) { + return refreshCandidate; + } + + const refreshed = await refreshOAuthCredential(storageProvider, refreshCandidate); + if (!refreshed) { + return initialPersistedCredential; + } + + primary.reload(); + const latestPersistedCredential = selectPersistedRefreshCredential(); + const changedWhileRefreshing = initialPersistedCredential + ? !isSameOAuthCredentialIdentity(latestPersistedCredential, initialPersistedCredential) + : latestPersistedCredential !== undefined; + if (changedWhileRefreshing) { + return latestPersistedCredential; + } + + await primary.set(storageProvider, refreshed); + return refreshed; + }) .then((refreshed) => { - if (refreshed) { + if (refreshed && !shouldRefreshOAuthCredential(refreshed)) { oauthRefreshCooldownUntil.delete(storageProvider); } else { oauthRefreshCooldownUntil.set(storageProvider, Date.now() + OAUTH_REFRESH_FAILURE_COOLDOWN_MS);