import { existsSync, mkdirSync, readFileSync, renameSync, unlinkSync, writeFileSync } from "node:fs"; import { homedir } from "node:os"; import { join } from "node:path"; import lockfile from "proper-lockfile"; import * as piAiOAuth from "@earendil-works/pi-ai/oauth"; import { getModel } from "@earendil-works/pi-ai/compat"; import type { ExtensionAPI, ProviderModelConfig } from "@earendil-works/pi-coding-agent"; import { getDefaultConfig, type PiPiConfig, readScopedFlantSettings, GLOBAL_CONFIG_PATH, writeConfigValue } from "./config.js"; import { listRegisteredSpecs, updateRegistryFromAvailableModels, setTierEnabled, isSubscriptionFallbackActive } from "./model-registry.js"; import { compareModelVersion } from "./model-version.js"; import { getLogger } from "./log.js"; import { buildUserAgent, injectBillingHeader, CC_IDENTITY } from "./billing-spoof.js"; /** * The token refreshers pi-ai used to export. * * pi >= 0.84 ships this subpath as an empty module — the binary bundles it that * way and the package publishes `export {}` — so both are absent wherever pi-pi * actually runs and every call site falls through to pi's own registry. They * are still read through here rather than deleted because an older pi-ai on the * same peer range does export them, and that path rotates and persists the * credential itself. */ const oauth = piAiOAuth as unknown as { refreshAnthropicToken?: (refreshToken: string) => Promise<{ refresh: string; access: string; expires: number }>; refreshGitHubCopilotToken?: (refreshToken: string, enterpriseUrl?: string) => Promise<{ refresh: string; access: string; expires: number; enterpriseUrl?: string }>; }; export interface OpenRouterModelData { name: string; context_length: number; max_completion_tokens: number; pricing: { prompt: number; completion: number; cacheRead: number; cacheWrite: number; }; modality: string; } export interface FlantSettings { enabled: boolean; autoUpdate: boolean; cacheTTLDays: number; lastUpdated: string | null; cachedFlantModels: string[] | null; cachedOpenRouterData: Record | null; /** * Models the last metadata fetch could map but OpenRouter does not publish, * keyed by bare gateway id with the OpenRouter id that was tried. Without * this record their permanently missing entries would re-invalidate the cache * on every startup, so the TTL would never hold. */ unmappedModels: Record | null; /** * When true, additionally register the `pp-flant-anthropic-sub` provider, * which routes Claude requests through the gateway using the user's personal * Claude OAuth token (billed against their personal Claude subscription). */ subscription: boolean; /** * Minutes between out-of-band "is the subscription limit cleared yet?" probes * while the sub→non-sub rate-limit fallback is active. On each interval a * cheap probe hits the sub model; on success the user is asked to switch back. * Default 10. */ switchBackIntervalMinutes: number; /** * When true, automatically fall back to the next provider tier on a rate * limit WITHOUT a confirmation dialogue (a non-blocking notification is shown * on every switch). When false, the legacy manual permission dialogue is * shown instead. Default true. */ autoRateLimitFallback: boolean; /** * Enable the Copilot provider tier (built-in `github-copilot` provider, keyed * off COPILOT_GITHUB_TOKEN). When enabled it sits ABOVE both Flant tiers in * precedence. Default false. */ copilotEnabled: boolean; } const GEMINI_MAP: Record = { "gemini-3-flash": "google/gemini-3.0-flash", "gemini-3.1-flash-lite": "google/gemini-3.1-flash-lite-preview", "gemini-3.1-pro": "google/gemini-3.1-pro-preview", }; export function resolveAgentDir(): string { const envKey = "PI_CODING_AGENT_DIR"; const envDir = process.env[envKey]; if (envDir) { if (envDir === "~") return homedir(); if (envDir.startsWith("~/")) return homedir() + envDir.slice(1); return envDir; } return join(homedir(), ".pi", "agent"); } const SETTINGS_DIR = join(resolveAgentDir(), "extensions", "pp", "cache"); const SETTINGS_PATH = join(SETTINGS_DIR, "flant-models.json"); const DEFAULT_SETTINGS: FlantSettings = { enabled: false, autoUpdate: true, cacheTTLDays: 3, lastUpdated: null, cachedFlantModels: null, cachedOpenRouterData: null, unmappedModels: null, subscription: false, switchBackIntervalMinutes: 10, autoRateLimitFallback: true, copilotEnabled: false, }; /** Provider name for the personal-subscription Claude routing. */ export const SUB_PROVIDER = "pp-flant-anthropic-sub"; const CLAUDE_REFRESH_MARGIN_MS = 5 * 60_000; const FORCED_REFRESH_COOLDOWN_MS = 60_000; const ANTHROPIC_TOKEN_URL = "https://platform.claude.com/v1/oauth/token"; // Anthropic's public Claude Code OAuth client id, identical to the one pi-ai // uses; the refresh grant is unauthenticated beyond it. const ANTHROPIC_CLIENT_ID = "9d1c250a-e61b-44d9-88ed-5944d1962f5e"; /** Prefix the gateway expects for personal-subscription Claude models. */ export const SUB_MODEL_PREFIX = "sub/"; /** * Read the Claude OAuth access token persisted by pi's built-in `anthropic` * OAuth provider. Returns null when absent or expired. The gateway uses this * token (forwarded as `Authorization: Bearer ...`) to bill the user's personal * Claude subscription; pi-ai automatically adds the Claude Code identity * headers because the token has the `sk-ant-oat` prefix. */ export function readClaudeOAuthToken(): string | null { const token = readStoredClaudeOAuthToken(); if (!token || token.expired) return null; return token.access; } // The persisted access token WITH its expiry state, so provider registration can // fall back to an expired token. The registration must exist before pi restores // the session's model (which happens at extension load, ahead of any refresh); // without it the restore fails and the session lands on an unrelated model. function readStoredClaudeOAuthToken(): { access: string; expired: boolean } | null { const authPath = join(resolveAgentDir(), "auth.json"); if (!existsSync(authPath)) return null; try { const raw = JSON.parse(readFileSync(authPath, "utf-8")) as { anthropic?: { access?: unknown; expires?: unknown }; }; const anthropic = raw.anthropic; if (!anthropic || typeof anthropic.access !== "string" || !anthropic.access) return null; return { access: anthropic.access, expired: typeof anthropic.expires === "number" && anthropic.expires <= Date.now(), }; } catch { return null; } } interface AnthropicOAuthCreds { type?: unknown; access?: unknown; refresh?: unknown; expires?: unknown; } interface CopilotOAuthCreds extends AnthropicOAuthCreds { enterpriseUrl?: unknown; } // pi's own ModelRegistry (ctx.modelRegistry), captured at session start. Used // as the refresh path when the direct `@earendil-works/pi-ai/oauth` import is // unavailable: the standalone pi binary (>= 0.84) bundles that subpath as an // empty module, so refreshAnthropicToken/refreshGitHubCopilotToken are // undefined at runtime there even though they type-check against node_modules. // Also the source of Claude model capabilities for the sub provider — pi keeps // its catalog fresher than the pinned pi-ai's static one, which lacks the very // models this matters for. let modelRegistryRef: { getApiKeyForProvider?: (provider: string) => Promise; find?: (provider: string, modelId: string) => { compat?: unknown; thinkingLevelMap?: unknown; contextWindow?: number; maxTokens?: number } | undefined; } | null = null; export function setModelRegistry(registry: unknown): void { modelRegistryRef = registry && typeof (registry as any).getApiKeyForProvider === "function" ? registry as any : null; } function readOAuthEntry(authPath: string, provider: string): T | undefined { try { return (JSON.parse(readFileSync(authPath, "utf-8")) as Record)[provider]; } catch { return undefined; } } /** * Refresh `provider`'s OAuth credential through pi's registry, which refreshes * and persists to auth.json itself, then re-read the persisted entry. Returns * null when no registry is available or pi could not produce a fresh token. */ async function refreshViaModelRegistry(authPath: string, provider: string): Promise { const registry = modelRegistryRef; if (!registry?.getApiKeyForProvider) return null; const apiKey = await registry.getApiKeyForProvider(provider); if (!apiKey) return null; const entry = readOAuthEntry(authPath, provider); if (!entry || typeof entry.access !== "string" || !entry.access) return null; if (typeof entry.expires === "number" && entry.expires <= Date.now()) return null; return entry; } export function readCopilotOAuthToken(): string | null { const authPath = join(resolveAgentDir(), "auth.json"); if (!existsSync(authPath)) return null; try { const raw = JSON.parse(readFileSync(authPath, "utf-8")) as { "github-copilot"?: CopilotOAuthCreds }; const copilot = raw["github-copilot"]; if (!copilot || typeof copilot.access !== "string" || !copilot.access) return null; if (typeof copilot.expires === "number" && copilot.expires <= Date.now()) return null; return copilot.access; } catch { return null; } } export async function refreshCopilotOAuthToken(): Promise { const log = getLogger(); const authPath = join(resolveAgentDir(), "auth.json"); if (!existsSync(authPath)) return null; const copilot = readOAuthEntry(authPath, "github-copilot"); if (!copilot || typeof copilot.access !== "string" || !copilot.access) return null; const REFRESH_MARGIN_MS = 5 * 60_000; const expires = typeof copilot.expires === "number" ? copilot.expires : 0; if (expires > Date.now() + REFRESH_MARGIN_MS) return copilot.access; if (typeof copilot.refresh !== "string" || !copilot.refresh) return null; if (typeof oauth.refreshGitHubCopilotToken !== "function") { try { const entry = await refreshViaModelRegistry(authPath, "github-copilot"); if (entry) return entry.access as string; log.debug({ s: "flant" }, "copilot oauth token refresh unavailable"); } catch (err: any) { log.debug({ s: "flant", err: err?.message }, "copilot oauth token refresh failed"); } return null; } const enterpriseUrl = typeof copilot.enterpriseUrl === "string" ? copilot.enterpriseUrl : undefined; let refreshed: { refresh: string; access: string; expires: number; enterpriseUrl?: string }; try { refreshed = await oauth.refreshGitHubCopilotToken!(copilot.refresh, enterpriseUrl); } catch (err: any) { log.debug({ s: "flant", err: err?.message }, "copilot oauth token refresh failed"); return null; } try { const release = lockfile.lockSync(authPath, { stale: 10000 }); try { let current: Record = {}; try { current = JSON.parse(readFileSync(authPath, "utf-8")) as Record; } catch {} const existing = current["github-copilot"] && typeof current["github-copilot"] === "object" ? current["github-copilot"] as Record : {}; const existingExpires = typeof existing.expires === "number" ? existing.expires : 0; if (existingExpires > Date.now() + REFRESH_MARGIN_MS && typeof existing.access === "string" && existing.access) { return existing.access; } current["github-copilot"] = { ...existing, type: "oauth", ...refreshed }; writeFileSync(authPath, JSON.stringify(current, null, 2) + "\n", "utf-8"); } finally { release(); } } catch (err: any) { log.debug({ s: "flant", err: err?.message }, "failed to persist copilot oauth token"); } return refreshed.access; } /** * Ensure the persisted Claude OAuth access token is fresh, refreshing it via * the stored refresh token when it has expired (or is within its safety * margin). The refreshed credentials are written back to `auth.json` under the * `anthropic` provider id in pi's own `{ type: "oauth", ... }` format, so both * this extension and pi's built-in `anthropic` provider observe the new token. * * Unlike readClaudeOAuthToken, this is async (a refresh performs a network * call) and returns the valid access token, or null when no usable * credentials exist / a refresh fails. Async entry points call this before the * synchronous readClaudeOAuthToken so downstream reads see a fresh token. */ export async function refreshClaudeOAuthToken(): Promise { const log = getLogger(); const authPath = join(resolveAgentDir(), "auth.json"); if (!existsSync(authPath)) return null; const anthropic = readOAuthEntry(authPath, "anthropic"); if (!anthropic || typeof anthropic.access !== "string" || !anthropic.access) return null; // Refresh ahead of expiry so in-flight requests never race a dying token. const expires = typeof anthropic.expires === "number" ? anthropic.expires : 0; if (expires > Date.now() + CLAUDE_REFRESH_MARGIN_MS) return anthropic.access; // Expired (or expiring soon, or no expiry recorded): try to refresh. if (typeof anthropic.refresh !== "string" || !anthropic.refresh) { log.debug({ s: "flant" }, "claude oauth token expired and no refresh token available"); return null; } if (typeof oauth.refreshAnthropicToken !== "function") { try { const entry = await refreshViaModelRegistry(authPath, "anthropic"); if (entry) { log.debug({ s: "flant" }, "refreshed claude oauth token via pi registry"); return entry.access as string; } log.debug({ s: "flant" }, "claude oauth token refresh unavailable"); } catch (err: any) { log.debug({ s: "flant", err: err?.message }, "claude oauth token refresh failed"); } return null; } let refreshed: { refresh: string; access: string; expires: number }; try { refreshed = await oauth.refreshAnthropicToken!(anthropic.refresh); } catch (err: any) { log.debug({ s: "flant", err: err?.message }, "claude oauth token refresh failed"); return null; } // A credential that could not be written is not usable: the grant consumed the // stored refresh token, so disk now holds something that cannot be refreshed // again, and every consumer reads the token from disk. const persisted = await persistClaudeCredentials(refreshed, { keepFresherExisting: true, consumedRefresh: anthropic.refresh }); if (persisted !== refreshed.access) return persisted; log.debug({ s: "flant" }, "refreshed claude oauth token"); return refreshed.access; } // Write refreshed credentials back under the `anthropic` provider id, using pi's // { type: "oauth", ... } shape and the same file lock pi uses, so pi's built-in // provider observes the rotation too. Returns the access token that ended up on // disk, or null when it could not be written. // // Two things can be on disk by the time the lock is taken, because the token was // minted before it: a credential another instance rotated to (keepFresherExisting // keeps it), or a credential minted from a DIFFERENT refresh token than the one // this rotation consumed (consumedRefresh) — meaning the other instance won the // race and ours is the loser of a superseded chain. Keeping theirs in both cases // is what stops two processes from revoking each other in a loop. async function persistClaudeCredentials( refreshed: { refresh: string; access: string; expires: number }, options: { keepFresherExisting?: boolean; consumedRefresh?: string } = {}, ): Promise { const authDir = resolveAgentDir(); const authPath = join(authDir, "auth.json"); try { if (!existsSync(authDir)) mkdirSync(authDir, { recursive: true }); if (!existsSync(authPath)) writeFileSync(authPath, "{}\n", "utf-8"); // Retried rather than failed-fast, which is why this takes the async lock: // the grant already consumed the previous refresh token, so losing the write // loses the only usable credential and forces a re-login. Mirrors the retry // policy of pi's own auth storage, whose lock this contends with. const release = await lockfile.lock(authPath, { stale: 10000, retries: { retries: 10, factor: 2, minTimeout: 100, maxTimeout: 10000, randomize: true }, }); try { let current: Record = {}; try { current = JSON.parse(readFileSync(authPath, "utf-8")) as Record; } catch { current = {}; } const existing = (current.anthropic && typeof current.anthropic === "object") ? current.anthropic as Record : {}; const existingExpires = typeof existing.expires === "number" ? existing.expires : 0; const existingAccess = typeof existing.access === "string" ? existing.access : ""; if (existingAccess) { if (options.keepFresherExisting && existingExpires > Date.now() + CLAUDE_REFRESH_MARGIN_MS) { return existingAccess; } if (options.consumedRefresh && existing.refresh !== options.consumedRefresh) { getLogger().debug({ s: "flant" }, "another instance rotated the claude credential first; keeping its token"); return existingAccess; } } current.anthropic = { ...existing, type: "oauth", access: refreshed.access, refresh: refreshed.refresh, expires: refreshed.expires, }; writeFileSync(authPath, JSON.stringify(current, null, 2) + "\n", "utf-8"); } finally { await release(); } } catch (err: any) { getLogger().warn({ s: "flant", err: err?.message }, "failed to persist refreshed claude oauth token"); return null; } return refreshed.access; } let forcedRefreshAt = 0; let forcedRefreshInFlight: Promise | null = null; // The credential this process last minted through a forced rotation. A request // rejecting THAT token means rotating again is not the fix (the gateway key, the // account, or the whole refresh chain is the problem), so the second attempt is // refused instead of rotating on every turn or every probe interval — which // would revoke the shared credential out from under every other client. let lastForcedMint: string | null = null; export type ForcedRefreshResult = /** A new credential was minted, or another instance's newer one adopted. */ | { status: "rotated"; token: string } /** Another rotation happened moments ago; this one was suppressed. */ | { status: "throttled" } /** No usable credential can be produced; the subscription is unusable. */ | { status: "failed" }; /** * Rotate the Claude OAuth credential even though the persisted one has not * expired. This is the recovery path for a token the server REVOKED — another * client rotating the same credential invalidates ours while it stays * clock-valid, so every expiry-based check (including pi's own registry * refresh) keeps handing out a dead token until it finally expires. * * `rejectedToken` must be the token the failing request actually carried, not * whatever is on disk now: when the persisted one already differs, another * instance rotated in the meantime and its token is adopted rather than rotated * away — otherwise two processes revoke each other in a loop. */ export async function forceRefreshClaudeOAuthToken(rejectedToken?: string | null): Promise { if (forcedRefreshInFlight) return forcedRefreshInFlight; forcedRefreshInFlight = (async () => { const log = getLogger(); const authPath = join(resolveAgentDir(), "auth.json"); if (!existsSync(authPath)) return { status: "failed" }; const anthropic = readOAuthEntry(authPath, "anthropic"); if (!anthropic) return { status: "failed" }; const persisted = typeof anthropic.access === "string" ? anthropic.access : ""; // Checked before the cooldown: parallel workers all fail on the same token, // and every one after the first must adopt the fresh credential rather than // be told to wait for a rotation that already happened. if (rejectedToken && persisted && persisted !== rejectedToken) { log.debug({ s: "flant" }, "claude oauth token already rotated elsewhere; reusing the persisted one"); return { status: "rotated", token: persisted }; } if (rejectedToken && lastForcedMint && rejectedToken === lastForcedMint) { log.warn({ s: "flant" }, "a freshly minted claude credential was rejected too; not rotating again"); return { status: "failed" }; } if (Date.now() - forcedRefreshAt < FORCED_REFRESH_COOLDOWN_MS) return { status: "throttled" }; if (typeof anthropic.refresh !== "string" || !anthropic.refresh) { log.debug({ s: "flant" }, "claude oauth credential rejected and no refresh token available"); return { status: "failed" }; } forcedRefreshAt = Date.now(); let refreshed: { refresh: string; access: string; expires: number }; try { refreshed = await requestClaudeTokenRefresh(anthropic.refresh); } catch (err: any) { log.warn({ s: "flant", err: err?.message }, "forced claude oauth token refresh failed"); // The usual reason a grant fails is that another instance consumed this // single-use refresh token first — in which case its replacement is on // disk by now and is exactly what this caller needs. const adopted = readOAuthEntry(authPath, "anthropic")?.access; if (typeof adopted === "string" && adopted && adopted !== persisted) { log.debug({ s: "flant" }, "adopting the credential another instance minted after our grant failed"); return { status: "rotated", token: adopted }; } return { status: "failed" }; } // A credential that could not be written is not usable: the provider is // rebound by re-reading auth.json, and pi's own provider reads the same // file, so returning an unpersisted token would register the rejected one. const stored = await persistClaudeCredentials(refreshed, { consumedRefresh: anthropic.refresh }); if (!stored) return { status: "failed" }; lastForcedMint = stored; log.info({ s: "flant" }, "force-refreshed the claude oauth token after a rejected credential"); return { status: "rotated", token: stored }; })(); try { return await forcedRefreshInFlight; } finally { forcedRefreshInFlight = null; } } /** * Record that a subscription request demonstrably succeeded. The credential in * use works, so a later rejection is a NEW failure deserving its own rotation * rather than the "already tried that" refusal above. */ export function noteSubscriptionCredentialAccepted(): void { lastForcedMint = null; } // Mint a new Claude OAuth credential from a refresh token. Prefers pi-ai's own // implementation and falls back to the same public endpoint/client it uses, // because the standalone pi binary bundles `pi-ai/oauth` as an empty module — // and pi's registry cannot stand in here, since it refuses to refresh a // credential it still considers unexpired. async function requestClaudeTokenRefresh(refreshToken: string): Promise<{ refresh: string; access: string; expires: number }> { if (typeof oauth.refreshAnthropicToken === "function") return oauth.refreshAnthropicToken(refreshToken); const res = await fetch(ANTHROPIC_TOKEN_URL, { method: "POST", headers: { "content-type": "application/json", accept: "application/json" }, body: JSON.stringify({ grant_type: "refresh_token", client_id: ANTHROPIC_CLIENT_ID, refresh_token: refreshToken }), signal: AbortSignal.timeout(30000), }); const body = await res.text(); if (!res.ok) throw new Error(`Anthropic token refresh returned ${res.status}: ${body.slice(0, 200)}`); const data = JSON.parse(body) as { refresh_token: string; access_token: string; expires_in: number }; return { refresh: data.refresh_token, access: data.access_token, expires: Date.now() + data.expires_in * 1000 - CLAUDE_REFRESH_MARGIN_MS, }; } /** Resolve the gateway API key (LLM_API_KEY preferred, FLANT_API_KEY fallback). */ export function readGatewayApiKey(): string | null { return process.env.LLM_API_KEY ?? process.env.FLANT_API_KEY ?? null; } let piRef: ExtensionAPI | null = null; let generatedFlantConfig: Partial | null = null; /** * Inputs needed to (re)register the personal-subscription Claude provider. * Cached at the last registerFlantProviders call so refreshSubProvider can * rebuild the provider with a freshly-read OAuth token without re-discovering * models. Null when subscription routing is not active. */ interface SubProviderContext { anthropicModels: string[]; metadata: Record; } let subProviderContext: SubProviderContext | null = null; /** The OAuth access token the sub provider was last registered with. */ let lastSubToken: string | null = null; export function setPI(pi: ExtensionAPI): void { piRef = pi; } export function clearFlantGeneratedConfig(): void { generatedFlantConfig = null; } export function unregisterFlantProviders(pi?: ExtensionAPI): void { // Cleared before the guard: otherwise the per-turn refreshSubProvider would // re-register the sub provider the next time the OAuth token rotates. subProviderContext = null; lastSubToken = null; const api = pi ?? piRef; if (!api) return; api.unregisterProvider("pp-flant-anthropic"); api.unregisterProvider("pp-flant-openai"); api.unregisterProvider(SUB_PROVIDER); } function ensureSettingsDir(): void { if (!existsSync(SETTINGS_DIR)) { mkdirSync(SETTINGS_DIR, { recursive: true }); } } function toNumber(value: unknown, fallback = 0): number { if (typeof value === "number" && Number.isFinite(value)) return value; if (typeof value === "string") { const parsed = Number(value); if (Number.isFinite(parsed)) return parsed; } return fallback; } function normalizeSettings(raw: unknown): FlantSettings { if (!raw || typeof raw !== "object") return { ...DEFAULT_SETTINGS }; const value = raw as Record; const cacheTTLDays = Math.max(1, Math.round(toNumber(value.cacheTTLDays, DEFAULT_SETTINGS.cacheTTLDays))); const switchBackIntervalMinutes = Math.max( 1, Math.round(toNumber(value.switchBackIntervalMinutes, DEFAULT_SETTINGS.switchBackIntervalMinutes)), ); return { enabled: !!value.enabled, autoUpdate: value.autoUpdate === undefined ? true : !!value.autoUpdate, cacheTTLDays, switchBackIntervalMinutes, autoRateLimitFallback: value.autoRateLimitFallback === undefined ? true : !!value.autoRateLimitFallback, copilotEnabled: !!value.copilotEnabled, subscription: !!value.subscription, lastUpdated: typeof value.lastUpdated === "string" ? value.lastUpdated : null, cachedFlantModels: Array.isArray(value.cachedFlantModels) ? value.cachedFlantModels.filter((m): m is string => typeof m === "string") : null, cachedOpenRouterData: value.cachedOpenRouterData && typeof value.cachedOpenRouterData === "object" ? value.cachedOpenRouterData as Record : null, unmappedModels: value.unmappedModels && typeof value.unmappedModels === "object" ? value.unmappedModels as Record : null, }; } // Item 8: durable user policy now lives in scoped PiPiConfig (`flant` section); // only these regenerable fields remain in the cache file. export interface FlantCache { lastUpdated: string | null; cachedFlantModels: string[] | null; cachedOpenRouterData: Record | null; unmappedModels: Record | null; } const DURABLE_FLANT_KEYS = [ "enabled", "subscription", "switchBackIntervalMinutes", "autoRateLimitFallback", "copilotEnabled", "autoUpdate", "cacheTTLDays", ] as const; // Values that a RELEASED version shipped as the default for each durable key. // The legacy cache file was written with JSON.stringify(settings), so it // persisted untouched defaults as if the user had chosen them. Migrating such a // value would launder a stale default into permanent explicit policy that // outranks every future default change (this is exactly how a v0.8.0-era // switchBackIntervalMinutes: 30 kept overriding the current 10). A cached value // matching a historical default is therefore treated as "never chosen". // // Accepted limitation: a user who deliberately picked exactly the old default is // indistinguishable from one who never touched it, and loses that choice. const HISTORICAL_DEFAULTS: Record = { switchBackIntervalMinutes: [30, 10], // 30 = v0.8.0/v0.9.0, 10 = v0.16.0 cacheTTLDays: [7], // 7 in every released version enabled: [false], subscription: [false], copilotEnabled: [false], autoUpdate: [true], autoRateLimitFallback: [true], }; function loadFlantCache(): FlantCache { const empty: FlantCache = { lastUpdated: null, cachedFlantModels: null, cachedOpenRouterData: null, unmappedModels: null }; if (!existsSync(SETTINGS_PATH)) return empty; try { const value = JSON.parse(readFileSync(SETTINGS_PATH, "utf-8")) as Record; return { lastUpdated: typeof value.lastUpdated === "string" ? value.lastUpdated : null, cachedFlantModels: Array.isArray(value.cachedFlantModels) ? value.cachedFlantModels.filter((m): m is string => typeof m === "string") : null, cachedOpenRouterData: value.cachedOpenRouterData && typeof value.cachedOpenRouterData === "object" ? (value.cachedOpenRouterData as Record) : null, unmappedModels: value.unmappedModels && typeof value.unmappedModels === "object" ? (value.unmappedModels as Record) : null, }; } catch { return empty; } } function saveFlantCache(cache: FlantCache): void { ensureSettingsDir(); if (!existsSync(SETTINGS_PATH)) writeFileSync(SETTINGS_PATH, "{}\n", "utf-8"); const release = lockfile.lockSync(SETTINGS_PATH, { stale: 10000 }); try { const out: FlantCache = { lastUpdated: cache.lastUpdated, cachedFlantModels: cache.cachedFlantModels, cachedOpenRouterData: cache.cachedOpenRouterData, unmappedModels: cache.unmappedModels, }; // Written through a temp file and renamed: a process killed mid-write would // otherwise leave truncated JSON, and an unreadable cache registers NO flant // providers at extension load — early enough that the host cannot restore // the session's model and falls back to an unrelated one. const temp = `${SETTINGS_PATH}.${process.pid}.tmp`; try { writeFileSync(temp, JSON.stringify(out, null, 2) + "\n", "utf-8"); renameSync(temp, SETTINGS_PATH); } catch (err) { try { unlinkSync(temp); } catch { /* the temp file may never have been created */ } throw err; } } finally { release(); } } let migrationDone = false; // One-time migration (item 8): the legacy combined cache file mixed durable // user policy with cache. Copy each durable field that is an OWN-PROPERTY of // the RAW legacy JSON (so normalization-filled defaults are NOT migrated) into // the GLOBAL config scope, skipping any key GLOBAL already sets explicitly // (project scope is never consulted — writing global never overwrites an // explicit project value). Then rewrite the cache file with cache fields only. // Idempotent: after the rewrite no durable own-props remain, so re-runs no-op. export function migrateLegacyFlantSettings(globalConfigPath = GLOBAL_CONFIG_PATH): void { if (migrationDone) return; migrationDone = true; if (!existsSync(SETTINGS_PATH)) return; let raw: Record; try { raw = JSON.parse(readFileSync(SETTINGS_PATH, "utf-8")) as Record; } catch { return; } const hasDurableOwnProp = DURABLE_FLANT_KEYS.some((k) => Object.prototype.hasOwnProperty.call(raw, k)); if (!hasDurableOwnProp) return; let globalRaw: Record | null = null; if (existsSync(globalConfigPath)) { try { globalRaw = JSON.parse(readFileSync(globalConfigPath, "utf-8")) as Record; } catch { globalRaw = null; } } const globalFlant = (globalRaw?.flant && typeof globalRaw.flant === "object") ? (globalRaw.flant as Record) : {}; // Normalize once so a legacy string/other type lands as the correct typed // config value, and so historical-default comparison happens post-coercion // (a legacy "30" must match the numeric 30). const normalized = normalizeSettings(raw) as unknown as Record; for (const key of DURABLE_FLANT_KEYS) { if (!Object.prototype.hasOwnProperty.call(raw, key)) continue; if (Object.prototype.hasOwnProperty.call(globalFlant, key)) continue; if (HISTORICAL_DEFAULTS[key]?.includes(normalized[key])) { getLogger().info({ s: "flant", key }, "skipped legacy flant setting matching a historical default"); continue; } writeConfigValue(globalConfigPath, ["flant", key], normalized[key]); getLogger().info({ s: "flant", key }, "migrated legacy flant setting into global config"); } // Rewrite the cache file WITHOUT durable fields, only after the durable // values above are safely persisted to global config. saveFlantCache(loadFlantCache()); } // Compose the runtime FlantSettings bundle from scoped config (durable policy) // + the reduced cache file. `cwd` binds project-scope overrides; omit it for // the side-effect-free global-only read used at init / in subagents. export function loadFlantSettings(cwd?: string): FlantSettings { const durable = readScopedFlantSettings(cwd); const cache = loadFlantCache(); return { enabled: durable.enabled, autoUpdate: durable.autoUpdate, cacheTTLDays: durable.cacheTTLDays, switchBackIntervalMinutes: durable.switchBackIntervalMinutes, autoRateLimitFallback: durable.autoRateLimitFallback, copilotEnabled: durable.copilotEnabled, subscription: durable.subscription, lastUpdated: cache.lastUpdated, cachedFlantModels: cache.cachedFlantModels, cachedOpenRouterData: cache.cachedOpenRouterData, unmappedModels: cache.unmappedModels, }; } // Persist ONLY cache fields (item 8). Durable policy is written through the // scoped-config mechanism (applyConfigChange), never back to the cache file. export function saveFlantSettings(settings: FlantSettings): void { saveFlantCache({ lastUpdated: settings.lastUpdated, cachedFlantModels: settings.cachedFlantModels, cachedOpenRouterData: settings.cachedOpenRouterData, unmappedModels: settings.unmappedModels, }); } function toTitleCase(token: string): string { const lower = token.toLowerCase(); if (lower === "gpt") return "GPT"; if (lower === "o3") return "O3"; if (lower === "qwen") return "Qwen"; if (lower === "claude") return "Claude"; if (lower === "api") return "API"; if (lower === "ai") return "AI"; if (!token.length) return token; return token.charAt(0).toUpperCase() + token.slice(1); } export function generateDisplayName(flantId: string): string { const parts = flantId.split("-").filter(Boolean); const out: string[] = []; let i = 0; while (i < parts.length) { const part = parts[i]; if (/^\d+$/.test(part)) { const versionParts = [part]; i += 1; while (i < parts.length && /^\d+$/.test(parts[i])) { versionParts.push(parts[i]); i += 1; } out.push(versionParts.join(".")); continue; } out.push(toTitleCase(part)); i += 1; } return out.join(" "); } function mapClaudeModelId(modelId: string): string { const rest = modelId.slice("claude-".length); const parts = rest.split("-").filter(Boolean); const firstNumber = parts.findIndex((p) => /^\d+$/.test(p)); if (firstNumber === -1) return `anthropic/claude-${rest}`; const family = parts.slice(0, firstNumber).join("-"); const version = parts.slice(firstNumber).join("."); return family.length > 0 ? `anthropic/claude-${family}-${version}` : `anthropic/claude-${version}`; } function normalizeQwenRest(rest: string): string { let value = rest; if (value.startsWith("-")) value = value.slice(1); if (!value) return "default"; return value; } function mapFlantToOpenRouterId(modelId: string): string | null { if (GEMINI_MAP[modelId]) return GEMINI_MAP[modelId]; if (modelId.startsWith("claude-")) return mapClaudeModelId(modelId); if (modelId.startsWith("gpt-")) return `openai/${modelId}`; if (modelId.startsWith("deepseek-")) return `deepseek/deepseek-${modelId.slice("deepseek-".length)}`; if (modelId.startsWith("grok-")) return `x-ai/grok-${modelId.slice("grok-".length)}`; if (modelId.startsWith("minimax-")) return `minimax/minimax-${modelId.slice("minimax-".length)}`; if (modelId.startsWith("qwen")) return `qwen/qwen-${normalizeQwenRest(modelId.slice("qwen".length))}`; if (modelId.startsWith("o3-")) return `openai/o3-${modelId.slice("o3-".length)}`; if (modelId.startsWith("sonar-")) return `perplexity/sonar-${modelId.slice("sonar-".length)}`; return null; } /** * Out-of-band probe: is the personal Claude subscription limit cleared yet? * Sends a tiny, fully throwaway request (NOT part of the session, so it cannot * pollute conversation/context/cache) to the gateway's Anthropic endpoint using * the personal Claude OAuth token. Returns: * "ok" — 200: capacity is back (offer switch-back) * "rate_limited" — 429: still limited (stay on non-sub, retry later) * "error" — credentials missing / network / other status (treat as not-back) * * The OAuth token is refreshed first so a failure is a genuine 429 and not an * expired-token false negative. `max_tokens: 1` + an explicit "just respond hi" * instruction keep output at ~1 token. */ // Derive the gateway probe model id (`sub/`) from any stored // form: `pp-flant-anthropic-sub/sub/`, `sub/`, or a bare ``. Exported // for testing the derivation without a network call. export function subProbeModelId(modelId: string): string { let bare = modelId; if (bare.startsWith(`${SUB_PROVIDER}/`)) bare = bare.slice(`${SUB_PROVIDER}/`.length); if (bare.startsWith(SUB_MODEL_PREFIX)) bare = bare.slice(SUB_MODEL_PREFIX.length); return `${SUB_MODEL_PREFIX}${bare}`; } // Local matcher for the 400 "extra usage" body (kept here to avoid a circular // import with rate-limit-fallback, which imports this module). Mirrors // isExtraUsageError's phrasing. function isExtraUsageMessage(text: string): boolean { return /extra usage|draw from[\s\S]{0,40}plan limits/i.test(text); } export async function probeSubscriptionCleared( modelId: string, ): Promise<"ok" | "rate_limited" | "error"> { const log = getLogger(); await refreshClaudeOAuthToken(); const probeModel = subProbeModelId(modelId); let { outcome, token } = await sendSubscriptionProbe(probeModel); // A 401 is a REVOKED credential, not a limit: without rotating it here the // probe returns "error" forever and the switch-back never fires — the session // stays on the fallback tier until the token happens to expire on its own. // The rotation goes through the shared recovery path so the provider is // rebound too, or a cleared probe would switch back onto a registration still // carrying the rejected token as its literal apiKey. if (outcome === "unauthenticated" && await recoverRejectedSubCredential(undefined, token) === "rotated") { ({ outcome } = await sendSubscriptionProbe(probeModel)); } switch (outcome) { case "ok": log.debug({ s: "flant", model: probeModel }, "probe: capacity back"); noteSubscriptionCredentialAccepted(); return "ok"; case "rate_limited": log.debug({ s: "flant", model: probeModel }, "probe: still rate limited"); return "rate_limited"; case "unauthenticated": log.warn({ s: "flant", model: probeModel }, "probe: credential rejected even after a forced refresh"); return "error"; default: return "error"; } } type SubscriptionProbeOutcome = "ok" | "rate_limited" | "unauthenticated" | "error"; // One minimal subscription request, shaped exactly like a live one, classified // by what came back. Callers decide what to do about each outcome. The token it // sent is reported alongside, so a rejection names the credential that actually // failed rather than whatever the file holds by the time the caller looks. async function sendSubscriptionProbe(probeModel: string): Promise<{ outcome: SubscriptionProbeOutcome; token: string | null }> { const log = getLogger(); const oauthToken = readClaudeOAuthToken(); const gatewayKey = readGatewayApiKey(); if (!oauthToken || !gatewayKey) { log.debug({ s: "flant", hasOAuth: !!oauthToken, hasGatewayKey: !!gatewayKey }, "probe skipped: missing credentials"); return { outcome: "error", token: oauthToken }; } // Build a payload that matches a LIVE subscription request's billing shape: // it must carry the Claude Code identity system block (so injectBillingHeader // passes its gate) and the billing system[0] entry the transform prepends. // Without this the probe lands in the "extra usage" bucket, keeps getting the // 400, is classified rate_limited, and the switch-back timer re-arms forever. const probePayload: Record = { model: probeModel, max_tokens: 1, system: [{ type: "text", text: CC_IDENTITY }], messages: [{ role: "user", content: "just respond hi" }], // No `temperature`: the gateway rejects it as deprecated for some models // (a separate 400 source) — sending it would make the probe 400 forever. }; injectBillingHeader(probePayload); try { const res = await fetch("https://llm-api.flant.ru/v1/messages", { method: "POST", headers: { "content-type": "application/json", "anthropic-version": "2023-06-01", // Claude Code OAuth identity headers — the gateway requires these for // subscription (sub/) routing; without them it can reject the probe and // we would never detect that the limit cleared. Mirrors doctor.ts. "anthropic-beta": "claude-code-20250219,oauth-2025-04-20", // Full-form Claude Code user-agent (same as live sub requests via the // sub provider's model.headers) — a bare claude-cli/1.0.0 is bucketed // as un-spoofed "extra usage" traffic and the probe never recovers. "user-agent": buildUserAgent(), "x-app": "cli", Authorization: `Bearer ${oauthToken}`, "x-litellm-api-key": `Bearer ${gatewayKey}`, }, body: JSON.stringify(probePayload), signal: AbortSignal.timeout(30000), }); if (res.ok) return { outcome: "ok", token: oauthToken }; if (res.status === 429) return { outcome: "rate_limited", token: oauthToken }; if (res.status === 401 || res.status === 403) return { outcome: "unauthenticated", token: oauthToken }; // A 400 "extra usage" means the subscription pool is still exhausted (the // same failure the fallback was triggered for) — treat it as still-limited so // the shared switch-back probe keeps waiting rather than declaring recovery. if (res.status === 400) { let bodyText = ""; try { bodyText = await res.text(); } catch { /* ignore */ } if (isExtraUsageMessage(bodyText)) return { outcome: "rate_limited", token: oauthToken }; } log.debug({ s: "flant", model: probeModel, status: res.status }, "probe: unexpected status"); return { outcome: "error", token: oauthToken }; } catch (err: any) { log.debug({ s: "flant", err: err?.message }, "probe failed"); return { outcome: "error", token: oauthToken }; } } /** * Establish whether the subscription credential currently works, rotating it * when it does not: * * "ok" — accepted as-is, or refused over quota rather than identity * "rotated" — was rejected, a fresh credential is now registered * "throttled" — was rejected, but a rotation just happened; retry later * "failed" — rejected and unrenewable; the subscription is unusable * "inconclusive" — nothing could be established (network, unexpected status) * * Every decision to rotate is made from THIS probe, whose token identity is * exact. Callers that only learn "something returned 401" — a finished turn, a * worker that failed minutes ago — cannot say which credential was rejected, * and guessing is what makes two instances rotate each other's token away. * * A 401 the gateway raised over its OWN key, or over a disabled account, is * indistinguishable here from a revoked token, so it does trigger one rotation. * What bounds it is that the rotation refuses to run twice on a credential this * process just minted, so the shared token cannot be spun once per turn. */ export async function reviveSubscriptionCredential( modelId: string, pi?: ExtensionAPI, ): Promise<"ok" | "inconclusive" | ForcedRefreshResult["status"]> { const { outcome, token } = await sendSubscriptionProbe(subProbeModelId(modelId)); if (outcome === "ok" || outcome === "rate_limited") { // The probe reads the token from disk, but the provider holds it as a // literal: when another instance rotated it, a probe that succeeds proves // only that the DISK credential works, and resuming without rebinding sends // the old one again — a 401 per turn, forever, each one "recovering" here. if (token && token !== lastSubToken) await refreshSubProvider(pi); // Quota, not identity: on a 429 the credential itself is fine, but it has // not been shown to be ACCEPTED, so the one-rotation guard stays armed. if (outcome === "ok") noteSubscriptionCredentialAccepted(); return "ok"; } if (outcome === "error") return "inconclusive"; getLogger().warn({ s: "flant" }, "the subscription credential was rejected; rotating it"); return recoverRejectedSubCredential(pi, token); } export async function discoverFlantModels(apiKey: string): Promise { const res = await fetch("https://llm-api.flant.ru/v1/models", { headers: { Authorization: `Bearer ${apiKey}` }, signal: AbortSignal.timeout(30000), }); if (!res.ok) { throw new Error(`Flant API returned ${res.status}`); } const payload = await res.json() as { data?: Array<{ id?: unknown }> }; const models = (payload.data ?? []) .map((m) => (typeof m.id === "string" ? m.id : "")) .filter((id) => id.length > 0 && !id.startsWith("or/")); return [...new Set(models)]; } export async function fetchOpenRouterMetadata(modelIds: string[]): Promise> { const mapping = new Map(); for (const modelId of modelIds) { // The gateway lists Claude only under `sub/`, while every consumer // (registerSubProvider, buildProviderModelConfig) looks metadata up by the // BARE id — so key the result bare, not as listed. const key = modelId.startsWith(SUB_MODEL_PREFIX) ? modelId.slice(SUB_MODEL_PREFIX.length) : modelId; const mapped = mapFlantToOpenRouterId(key); if (mapped) mapping.set(key, mapped); } if (mapping.size === 0) return {}; const res = await fetch("https://openrouter.ai/api/v1/models", { signal: AbortSignal.timeout(30000), }); if (!res.ok) { throw new Error(`OpenRouter API returned ${res.status}`); } const payload = await res.json() as { data?: Array> }; const modelMap = new Map>(); for (const item of payload.data ?? []) { const id = typeof item.id === "string" ? item.id : ""; if (id) modelMap.set(id, item); } const out: Record = {}; for (const [flantModelId, openRouterId] of mapping.entries()) { const model = modelMap.get(openRouterId); if (!model) continue; const pricing = model.pricing && typeof model.pricing === "object" ? model.pricing as Record : {}; const architecture = model.architecture && typeof model.architecture === "object" ? model.architecture as Record : {}; const topProvider = model.top_provider && typeof model.top_provider === "object" ? model.top_provider as Record : {}; out[flantModelId] = { name: typeof model.name === "string" ? model.name : generateDisplayName(flantModelId), context_length: toNumber(model.context_length, 200000), max_completion_tokens: toNumber(topProvider.max_completion_tokens, 32000), pricing: { prompt: toNumber(pricing.prompt, 0), completion: toNumber(pricing.completion, 0), cacheRead: toNumber(pricing.input_cache_read, 0), cacheWrite: toNumber(pricing.input_cache_write, 0), }, modality: typeof architecture.modality === "string" ? architecture.modality : "text", }; } return out; } /** * Returns true when the personal-subscription Claude path is fully usable: * the setting is enabled AND both credentials (Claude OAuth token + gateway * key) are present. Mirrors the gate in registerFlantProviders so we never * generate `sub/` role assignments that cannot resolve to a real provider. */ export function isSubscriptionActive(settings?: FlantSettings): boolean { const s = settings ?? loadFlantSettings(); return s.subscription && !!readClaudeOAuthToken() && !!readGatewayApiKey(); } function modelSpec(modelId: string): string { if (modelId.startsWith("claude-")) { return `${SUB_PROVIDER}/${SUB_MODEL_PREFIX}${modelId}`; } return `pp-flant-openai/${modelId}`; } /** * The gateway is LiteLLM in front of OpenRouter, and neither pins a session to * one upstream by default: consecutive requests land on deployments holding * unrelated prompt-cache state, so a long session re-buys the tail of its own * prompt every turn. `x-session-id` is the header OpenRouter documents for * sticky routing and the shape LiteLLM's own session-affinity check reads. * Inert until the gateway forwards client headers, which costs nothing while * it does not. Not applied to the `sub/` provider, whose compat pi's Claude * catalog owns. */ const SESSION_AFFINITY_COMPAT = { sendSessionAffinityHeaders: true, sessionAffinityFormat: "openrouter" } as const; function buildProviderModelConfig( flantModelId: string, metadata: Record, ): ProviderModelConfig { const modelMeta = metadata[flantModelId]; const modality = (modelMeta?.modality ?? "text").toLowerCase(); return { id: flantModelId, name: modelMeta?.name ?? generateDisplayName(flantModelId), reasoning: true, input: modality.includes("image") ? ["text", "image"] : ["text"], cost: { input: toNumber(modelMeta?.pricing.prompt, 0) * 1_000_000, output: toNumber(modelMeta?.pricing.completion, 0) * 1_000_000, cacheRead: toNumber(modelMeta?.pricing.cacheRead, 0) * 1_000_000, cacheWrite: toNumber(modelMeta?.pricing.cacheWrite, 0) * 1_000_000, }, contextWindow: toNumber(modelMeta?.context_length, 200000), maxTokens: toNumber(modelMeta?.max_completion_tokens, 32000), }; } /** * Copy a Claude model's capability metadata from pi's own catalog entry for the * bare id onto its `sub/`-prefixed twin. Without this the sub provider carries * no `compat`, so pi-ai sends the LEGACY `thinking: {type: "enabled"}` form and * the gateway returns a thinking block whose text is empty; only the adaptive * form `forceAdaptiveThinking` selects yields readable summaries. The same gap * drops `thinkingLevelMap`, silently clamping xhigh/max down to high. * * `allowedFallbackModels` is deliberately dropped: pi-ai turns it into a * server-side `fallbacks` param naming a BARE claude id, which the gateway key * cannot access (it only sees `sub/*`) and which would fail the whole request. * * The window and output cap come from here too. Without them a model the * gateway's metadata does not cover falls back to a 200K window, and a 200K * window puts the fold ceiling at 150K: every request over that folds as hard * as it can, still does not fit, and rewrites the prompt's prefix for nothing. * * `supportsMidConvoEffort` is inherited again. It rides the * `mid-conversation-output-config-2026-07-01` and * `thinking-binding-controls-2026-08-01` betas, which the gateway used to strip * from `anthropic-beta` against a hardcoded dictionary — every turn on the * flagged models (claude-opus-5, claude-fable-5-1) then 400'd with `Extra inputs * are not permitted`. The gateway now forwards an unknown Anthropic beta as-is, * verified against both betas on both models. */ function applyClaudeModelCapabilities(config: ProviderModelConfig, bareModelId: string): ProviderModelConfig { const catalog = modelRegistryRef?.find?.("anthropic", bareModelId) ?? (typeof getModel === "function" ? getModel("anthropic" as never, bareModelId as never) : undefined); if (!catalog) return config; const { allowedFallbackModels: _droppedFallbacks, ...compat } = (catalog.compat ?? {}) as Record; const out: ProviderModelConfig = { ...config }; if (Object.keys(compat).length > 0) out.compat = compat as ProviderModelConfig["compat"]; if (catalog.thinkingLevelMap) out.thinkingLevelMap = catalog.thinkingLevelMap as ProviderModelConfig["thinkingLevelMap"]; if (typeof catalog.contextWindow === "number" && catalog.contextWindow > 0) out.contextWindow = catalog.contextWindow; if (typeof catalog.maxTokens === "number" && catalog.maxTokens > 0) out.maxTokens = catalog.maxTokens; return out; } export interface RegisterFlantOptions { /** Whether to also register the personal-subscription provider. Defaults to loadFlantSettings().subscription. */ subscription?: boolean; } export function registerFlantProviders( pi: ExtensionAPI, models: string[], metadata: Record, options: RegisterFlantOptions = {}, ): void { const log = getLogger(); const uniqueModels = [...new Set(models)]; // The gateway no longer serves bare claude-* models over the paid API — // Claude is available EXCLUSIVELY through the personal subscription's `sub/` // groups. pp-flant-anthropic is not registered at all; every claude id in the // catalog (bare from legacy caches, or `sub/`-prefixed from the live gateway) // is a subscription candidate. const subCandidates = [...new Set( uniqueModels .filter((m) => m.startsWith("claude-") || m.startsWith(`${SUB_MODEL_PREFIX}claude-`)) .map((m) => (m.startsWith(SUB_MODEL_PREFIX) ? m.slice(SUB_MODEL_PREFIX.length) : m)), )]; const openaiModels = uniqueModels.filter((m) => !m.startsWith("claude-") && !m.startsWith(SUB_MODEL_PREFIX)); unregisterFlantProviders(pi); const gatewayKey = readGatewayApiKey() ?? "$FLANT_API_KEY"; pi.registerProvider("pp-flant-openai", { api: "openai-completions", baseUrl: "https://llm-api.flant.ru/v1", apiKey: gatewayKey, models: openaiModels.map((m) => ({ ...buildProviderModelConfig(m, metadata), compat: { ...SESSION_AFFINITY_COMPAT } })), }); const availableSpecs = openaiModels.map((id) => `pp-flant-openai/${id}`); const subscription = options.subscription ?? loadFlantSettings().subscription; let subModels: string[] = []; if (subscription) { // Remember the models/metadata so refreshSubProvider can rebuild the // provider with a fresh OAuth token on each turn (see below). subProviderContext = { anthropicModels: subCandidates, metadata }; subModels = registerSubProvider(pi, subCandidates, metadata); availableSpecs.push(...subModels.map((id) => `${SUB_PROVIDER}/${id}`)); } else { subProviderContext = null; lastSubToken = null; } log.debug({ s: "flant", total: uniqueModels.length, sub: subModels.length, openai: openaiModels.length }, "registering flant providers"); // updateRegistryFromAvailableModels REPLACES the catalog. Keep the non-flant // specs (github-copilot, native anthropic/openai) that session_start already // registered, or tier resolution loses the Copilot fallback right after this // re-registration. const foreignSpecs = listRegisteredSpecs().filter((spec) => !spec.startsWith("pp-flant-anthropic/") && !spec.startsWith("pp-flant-openai/") && !spec.startsWith(`${SUB_PROVIDER}/`)); updateRegistryFromAvailableModels([...foreignSpecs, ...availableSpecs]); } /** * Register (or re-register) the personal-subscription Claude provider using the * OAuth access token currently persisted in auth.json. Returns the list of * `sub/`-prefixed model ids registered (empty when credentials are missing). * * The provider is registered with a literal `apiKey` (the resolved token), so * the value is a static snapshot for the lifetime of the registration. Because * the OAuth token expires within a few hours, refreshSubProvider must be called * periodically (on each turn) to rebuild the provider with a fresh token; * otherwise long-lived sessions send a stale token and the gateway returns 401. */ function registerSubProvider( pi: ExtensionAPI, anthropicModels: string[], metadata: Record, ): string[] { const log = getLogger(); const stored = readStoredClaudeOAuthToken(); const gatewayKey = readGatewayApiKey(); if (!stored || !gatewayKey) { log.debug({ s: "flant", hasOAuth: !!stored, hasGatewayKey: !!gatewayKey }, "subscription enabled but credentials missing; skipping sub provider"); lastSubToken = null; return []; } const oauthToken = stored.access; if (stored.expired) { log.debug({ s: "flant" }, "registering the sub provider with an expired token; a refresh re-registers it before the first request"); } pi.registerProvider(SUB_PROVIDER, { name: "Flant (personal Claude subscription)", api: "anthropic-messages", baseUrl: "https://llm-api.flant.ru", // The OAuth token (sk-ant-oat...) triggers pi-ai's Claude Code identity // headers and is forwarded as `Authorization: Bearer ...`. apiKey: oauthToken, // Gateway key travels in a side header so it does not clobber the OAuth auth. // The full-form Claude Code user-agent (item 10) overrides pi-ai's hardcoded // bare `claude-cli/`: mergeHeaders applies these model headers AFTER the // default UA, and Anthropic's plan-billing validation expects this form. Its // version must match the billing header's cc_version (both from CC_VERSION). // Best-effort — LiteLLM may replace it upstream. headers: { "x-litellm-api-key": `Bearer ${gatewayKey}`, "user-agent": buildUserAgent() }, // Model id carries the `sub/` prefix the gateway expects, while pricing/ // metadata is looked up by the bare claude-* id. models: anthropicModels.map((m) => { const cfg = applyClaudeModelCapabilities(buildProviderModelConfig(m, metadata), m); return { ...cfg, id: `${SUB_MODEL_PREFIX}${m}` }; }), }); lastSubToken = oauthToken; return anthropicModels.map((m) => `${SUB_MODEL_PREFIX}${m}`); } /** * Ensure the personal-subscription Claude provider is registered with a * non-expired OAuth token. Refreshes the token (persisting it to auth.json when * needed) and, when it changed since the last registration, re-registers the * provider so subsequent LLM calls pick up the fresh token. * * Called on each turn for the root session. No-op when subscription routing is * not active. Cheap when the token is unchanged (only a token read + compare). */ export async function refreshSubProvider(pi?: ExtensionAPI): Promise { const ctx = subProviderContext; if (!ctx) return; const api = pi ?? piRef; if (!api) return; await refreshClaudeOAuthToken(); const token = readClaudeOAuthToken(); const log = getLogger(); if (!token) { // The registration still carries the previous token as a literal apiKey, so // leaving it in place turns every Claude request into a 401 instead of // letting resolution fall to another tier. The cached context stays, so a // later refresh that succeeds re-registers the provider. api.unregisterProvider(SUB_PROVIDER); lastSubToken = null; log.debug({ s: "flant" }, "subscription oauth token unavailable; unregistered the sub provider"); return; } // Only re-register when the token actually changed; registerProvider takes // effect immediately, so re-registering every turn would be wasteful churn. if (token === lastSubToken) return; registerSubProvider(api, ctx.anthropicModels, ctx.metadata); log.debug({ s: "flant" }, "re-registered the sub provider with a refreshed oauth token"); } /** * Recover from a subscription request the gateway rejected as unauthenticated. * Rotates the OAuth credential past its unexpired-looking state and rebinds the * sub provider to the new token. Returns false when no usable credential could * be minted — the caller must then route away from the subscription, because * the registration keeps failing every request until the user re-logs in. * * `rejectedToken` is the credential the failed request carried; it defaults to * the one the provider is registered with, NOT to whatever is on disk now, * which may already be another instance's replacement. */ export async function recoverRejectedSubCredential( pi?: ExtensionAPI, rejectedToken?: string | null, ): Promise { const rejected = rejectedToken ?? lastSubToken ?? readClaudeOAuthToken(); const result = await forceRefreshClaudeOAuthToken(rejected); if (result.status !== "rotated") return result.status; const ctx = subProviderContext; const api = pi ?? piRef; if (ctx && api && result.token !== lastSubToken) { registerSubProvider(api, ctx.anthropicModels, ctx.metadata); getLogger().info({ s: "flant" }, "rebound the sub provider after a rejected credential"); } return "rotated"; } function pickLatest(models: string[]): string | null { if (models.length === 0) return null; return models .slice() .sort((a, b) => compareModelVersion(b, a))[0] ?? null; } function pickFastFallbackModel(models: string[]): string | null { const gptMini = pickLatest(models.filter((m) => /^gpt-5.*-mini$/.test(m))); if (gptMini) return gptMini; return pickLatest(models.filter((m) => /^claude-haiku-/.test(m))); } function makeVariant(modelId: string | null, fallbackModelId: string): { enabled: boolean; model: string; thinking: string } { if (!modelId) return { enabled: false, model: modelSpec(fallbackModelId), thinking: "high" }; return { enabled: true, model: modelSpec(modelId), thinking: "high" }; } export function generateFlantConfig(models: string[], subscriptionActive = false): Partial { const rawModels = [...new Set(models)]; if (rawModels.length === 0) return {}; // Claude is served EXCLUSIVELY through the subscription's `sub/` groups — the // paid gateway has no Claude anymore. Normalize `sub/claude-*` ids to bare // claude ids for the pickers, and drop Claude entirely when the subscription // is inactive so generated roles never point at an unroutable model. const uniqueModels = [...new Set(rawModels.map((m) => (m.startsWith(SUB_MODEL_PREFIX) ? m.slice(SUB_MODEL_PREFIX.length) : m)))] .filter((m) => (subscriptionActive || !m.startsWith("claude-")) && !m.startsWith("gemini-")); if (uniqueModels.length === 0) return {}; const latestOpus = pickLatest(uniqueModels.filter((m) => /^claude-opus-/.test(m))); const latestFable = pickLatest(uniqueModels.filter((m) => /^claude-fable-/.test(m))); const latestClaude = pickLatest(uniqueModels.filter((m) => /^claude-/.test(m))); // gpt-5.6+ tier pickers. Base and -pro are distinct SKUs within a tier: the // `-pro` regexes are end-anchored on `-pro`, the base regexes negative-look // ahead to exclude `-pro`, so gpt-6-astra and gpt-6-astra-pro never collide. const gptAstra = pickLatest(uniqueModels.filter((m) => /^gpt-[0-9.]+-astra$/.test(m))); const gptAstraPro = pickLatest(uniqueModels.filter((m) => /^gpt-[0-9.]+-astra-pro$/.test(m))); const gptSol = pickLatest(uniqueModels.filter((m) => /^gpt-[0-9.]+-sol$/.test(m))); const gptSolPro = pickLatest(uniqueModels.filter((m) => /^gpt-[0-9.]+-sol-pro$/.test(m))); const gptTerra = pickLatest(uniqueModels.filter((m) => /^gpt-[0-9.]+-terra$/.test(m))); const gptLuna = pickLatest(uniqueModels.filter((m) => /^gpt-[0-9.]+-luna$/.test(m))); // Legacy fallback for pre-5.6 gateways that expose a single gpt SKU: keep the // old "latest gpt-5, else any gpt" behavior so tier pickers degrade to it. const isLegacyGpt = (model: string) => !model.endsWith("-mini") && !model.endsWith("-codex") && !/-(?:sol|terra|luna)(?:-pro)?$/.test(model); const latestGpt5 = pickLatest(uniqueModels.filter((m) => /^gpt-5/.test(m) && isLegacyGpt(m))); const latestGptLegacy = latestGpt5 ?? pickLatest(uniqueModels.filter((m) => /^gpt-/.test(m) && isLegacyGpt(m))); // Per-role gpt selections with graceful degradation to the legacy single SKU // when the split tiers are absent (older catalogs / non-flant gateways). const gptSmartPro = gptSolPro ?? gptSol ?? latestGptLegacy; const gptSmart = gptSol ?? latestGptLegacy; // The GPT the on-demand pools consult: the strongest SKU the gateway serves, // Astra Pro first. The pools are the one place worth the top-end price — they // run on demand, for judgment the session cannot produce itself — so they do // NOT degrade to the Sol tier that main-line roles (debug, fast) run on. const gptTop = gptAstraPro ?? gptAstra ?? gptSmartPro; const gptBalanced = gptTerra ?? gptSmart; const gptFast = gptLuna ?? gptBalanced; const latestGpt = gptSmart; const latestDeepseek = pickLatest(uniqueModels.filter((m) => /^deepseek-/.test(m))); const latestGrok = pickLatest(uniqueModels.filter((m) => /^grok-/.test(m))); const fastFallback = pickFastFallbackModel(uniqueModels); const fallback = latestOpus ?? latestClaude ?? latestGpt ?? latestDeepseek ?? latestGrok ?? uniqueModels[0]; const mainModel = latestOpus ?? latestClaude ?? fallback; const debugModel = gptSmart ?? latestDeepseek ?? fallback; const taskModel = latestOpus ?? latestClaude ?? fallback; const fastModel = gptLuna ?? fastFallback ?? gptFast ?? debugModel; return { agents: { main: { model: modelSpec(mainModel), thinking: "high" }, maxConcurrentSubagents: getDefaultConfig().agents.maxConcurrentSubagents, subagents: { simple: { explore: { model: modelSpec(fastModel), thinking: "medium" }, librarian: { model: modelSpec(fastModel), thinking: "medium" }, task: { model: modelSpec(taskModel), thinking: "medium" }, }, pools: { // Every pool is the same pair: the top-end GPT for a genuinely foreign // read, the latest Fable for a same-vendor one. Roles differ in what // they are ASKED, not in which models answer. advisors: [ makeVariant(gptTop, fallback), makeVariant(latestFable, fallback), ], reviewers: [ makeVariant(gptTop, fallback), makeVariant(latestFable, fallback), ], deepDebuggers: [ makeVariant(gptTop, fallback), makeVariant(latestFable, fallback), ], }, }, }, }; } // Bare gateway id → the OpenRouter id a metadata fetch would look up. Only ids // a fetch already tried and failed to resolve belong here; pairing each with the // id that was tried lets a later mapping change retry it. function collectUnmappedModels(models: string[], metadata: Record): Record { const out: Record = {}; for (const modelId of models) { const bare = modelId.startsWith(SUB_MODEL_PREFIX) ? modelId.slice(SUB_MODEL_PREFIX.length) : modelId; const mapped = mapFlantToOpenRouterId(bare); if (mapped && !metadata[bare]) out[bare] = mapped; } return out; } function isCacheValid(settings: FlantSettings): boolean { if (!settings.lastUpdated || !settings.cachedFlantModels || !settings.cachedOpenRouterData) return false; // An empty metadata map carries no context window and no pricing, so it is // never worth serving from cache — including caches an earlier build stamped. if (Object.keys(settings.cachedOpenRouterData).length === 0) return false; // A cache an earlier build stamped can be missing a whole family whose ids it // could not yet map, which leaves those models on the fallback window for the // rest of the TTL. Any mappable model without an entry means refetch, unless // the last fetch already established that OpenRouter has nothing under that id. const metadata = settings.cachedOpenRouterData; const unmapped = settings.unmappedModels ?? {}; for (const modelId of settings.cachedFlantModels) { const bare = modelId.startsWith(SUB_MODEL_PREFIX) ? modelId.slice(SUB_MODEL_PREFIX.length) : modelId; const mapped = mapFlantToOpenRouterId(bare); if (mapped && !metadata[bare] && unmapped[bare] !== mapped) return false; } const updatedAt = new Date(settings.lastUpdated).getTime(); if (!Number.isFinite(updatedAt)) return false; const ttlMs = Math.max(1, settings.cacheTTLDays) * 24 * 60 * 60 * 1000; return Date.now() - updatedAt < ttlMs; } export function getFlantGeneratedConfig(): Partial | null { return generatedFlantConfig; } (globalThis as any)[Symbol.for("pi-pi:flant-config")] = getFlantGeneratedConfig; export interface UpdateFlantOptions { /** * Bypass the TTL cache and re-fetch the model list from the gateway even when * the cache is still fresh. Set by the explicit "Update now" action so it does * what its label promises; auto/startup callers omit it and stay cache-bound. */ force?: boolean; /** Project dir for binding project-scoped flant overrides; omit for global-only. */ cwd?: string; } export async function updateFlantInfra( pi: ExtensionAPI, options: UpdateFlantOptions = {}, ): Promise<{ ok: boolean; error?: string; models?: string[] }> { setPI(pi); const settings = loadFlantSettings(options.cwd); // Refresh the personal-subscription Claude OAuth token before (re)registering // providers so the sub provider is built with a valid, non-expired token. if (settings.subscription) { await refreshClaudeOAuthToken(); } const cacheOk = !options.force && isCacheValid(settings); let models = cacheOk ? settings.cachedFlantModels : null; let metadata = cacheOk ? settings.cachedOpenRouterData : null; let refreshed = false; if (!models || !metadata) { const apiKey = readGatewayApiKey(); if (!apiKey) { if (settings.cachedFlantModels && settings.cachedOpenRouterData) { models = settings.cachedFlantModels; metadata = settings.cachedOpenRouterData; } else { return { ok: false, error: "LLM_API_KEY (or FLANT_API_KEY) is not set" }; } } else { try { models = await discoverFlantModels(apiKey); let metadataFetched = true; try { metadata = await fetchOpenRouterMetadata(models); } catch { metadataFetched = false; metadata = settings.cachedOpenRouterData ?? {}; } settings.cachedFlantModels = models; settings.cachedOpenRouterData = metadata; // Only a fetch that actually reached OpenRouter can establish that a // model has no entry there; recording it after a failed one would // suppress the retry that fills in the models still missing metadata. if (metadataFetched) settings.unmappedModels = collectUnmappedModels(models, metadata); // Serving empty metadata pins every model to the fallback context window // and zero cost, so it must never look cache-valid: clear the timestamp // outright rather than merely declining to refresh it — a forced update // inside the TTL would otherwise keep the previous one alive. A fetch // that succeeds but matches nothing is just as empty as a failed one. settings.lastUpdated = Object.keys(metadata).length > 0 ? new Date().toISOString() : null; saveFlantSettings(settings); refreshed = true; } catch (err: any) { if (settings.cachedFlantModels && settings.cachedOpenRouterData) { models = settings.cachedFlantModels; metadata = settings.cachedOpenRouterData; } else { return { ok: false, error: err?.message ?? "Failed to update Flant infrastructure" }; } } } } if (!models || !metadata) { return { ok: false, error: "No Flant model data available" }; } try { // Forward the EFFECTIVE (cwd-scoped) subscription so a project override // binds, instead of letting registerFlantProviders fall back to a // global-only read. registerFlantProviders(pi, models, metadata, { subscription: settings.subscription }); generatedFlantConfig = generateFlantConfig(models, isSubscriptionActive(settings)); // Backfill a timestamp for a cache that predates stamping — but never for an // empty metadata map, which is exactly the state the stamp must not bless. if (!refreshed && settings.cachedFlantModels && !settings.lastUpdated && Object.keys(settings.cachedOpenRouterData ?? {}).length > 0) { settings.lastUpdated = new Date().toISOString(); saveFlantSettings(settings); } return { ok: true, models }; } catch (err: any) { return { ok: false, error: err?.message ?? "Failed to register Flant providers" }; } } /** Whether the Copilot provider tier is enabled and has environment or /login credentials. */ export function isCopilotTierActive(settings?: FlantSettings): boolean { const s = settings ?? loadFlantSettings(); return s.copilotEnabled && !!(process.env.COPILOT_GITHUB_TOKEN || readCopilotOAuthToken()); } /** * Push the current provider-tier enable flags from settings into the model * registry's resolver. flant-api is the always-on paid floor; flant-sub follows * the subscription toggle + credentials; copilot follows its toggle + token. */ export function syncProviderTiers(settings?: FlantSettings): void { const s = settings ?? loadFlantSettings(); // Respect a LIVE rate-limit sub-fallback: setSubscriptionFallbackActive(true) // disables flant-sub, and syncProviderTiers must NOT clobber that back on // (e.g. when the Copilot toggle triggers a resync) or resolutions would route // onto the still-limited subscription while the fallback timers believe it's // off. flant-sub is enabled only when the subscription is active AND no live // fallback is in effect. setTierEnabled({ "copilot": isCopilotTierActive(s), "flant-sub": isSubscriptionActive(s) && !isSubscriptionFallbackActive(), "flant-api": true, }); } // `cwd`, when supplied, binds project-scoped flant overrides (the root session's // project dir, shared with child sessions via ORCHESTRATOR_CWD_KEY). Omit for a // global-only read (e.g. before any cwd is known). export function initFlantSync(pi: ExtensionAPI, cwd?: string): void { setPI(pi); const settings = loadFlantSettings(cwd); const log = getLogger(); syncProviderTiers(settings); if (!settings.enabled) { log.debug({ s: "flant" }, "flant disabled"); generatedFlantConfig = null; return; } if (settings.cachedFlantModels && settings.cachedOpenRouterData) { // Forward the EFFECTIVE (cwd-scoped) subscription — registerFlantProviders // otherwise falls back to a GLOBAL-only read, missing a project override. registerFlantProviders(pi, settings.cachedFlantModels, settings.cachedOpenRouterData, { subscription: settings.subscription }); generatedFlantConfig = generateFlantConfig(settings.cachedFlantModels, isSubscriptionActive(settings)); } } // `cwd` is the root session's project dir (known at session_start); passing it // lets project-scoped flant overrides bind. Omit for the global-only read. export async function initFlantOnStartup(pi: ExtensionAPI, cwd?: string): Promise { setPI(pi); const settings = loadFlantSettings(cwd); if (settings.copilotEnabled && !process.env.COPILOT_GITHUB_TOKEN) await refreshCopilotOAuthToken(); // Refresh before reading tiers: a stored subscription token that merely // expired still refreshes, but the tier check reads it as absent, and nothing // re-syncs afterwards — every Claude spec would resolve away from the // subscription for the rest of the session. if (settings.subscription) await refreshClaudeOAuthToken(); // Sync tiers from the EFFECTIVE (cwd-scoped) settings so a project override // rebinds what initFlantSync (global-only, at extension init) computed. syncProviderTiers(settings); if (!settings.enabled) { // A project override may DISABLE flant even though initFlantSync already // registered the global-enabled providers — unregister them so the override // is honored rather than leaving stale registrations in place. unregisterFlantProviders(pi); generatedFlantConfig = null; return; } if (!settings.autoUpdate) { // autoUpdate=false means "do not REFRESH the cache", NOT "stay unregistered". // With effective (project-scoped) enabled=true, still register from the // cached models so a global-disabled + project-enabled + autoUpdate=false // session actually gets its providers (initFlantSync ran global-only and // registered nothing). If there is no cache yet, bootstrap once via update. if (settings.cachedFlantModels && settings.cachedOpenRouterData) { registerFlantProviders(pi, settings.cachedFlantModels, settings.cachedOpenRouterData, { subscription: settings.subscription }); generatedFlantConfig = generateFlantConfig(settings.cachedFlantModels, isSubscriptionActive(settings)); return; } await updateFlantInfra(pi, { cwd }); return; } await updateFlantInfra(pi, { cwd }); }