refactor(app): DRY pass — split large files, extract shared utils
S1: delete page.new.js (1724L abandoned) + remove dead getAntigravityProjectId
S2: split large files by natural seams
- usage.js → usage/{github,google,claude,codex,kiro,minimax,misc,shared}.js
- media-providers page → components/{Embedding,Tts,Generic,Stt}ExampleCard.js
- EndpointPageClient → endpointConstants.js + endpointPing.js + components/
- tokenRefresh.js → tokenRefresh/{dedup,providers}.js
- ProviderLimits/index.js: 16 pure fn + 9 constants → utils.js
- oauth/providers.js: 7 pure helpers → providerHelpers.js
S3: shared utils
- getModelKind(m, fallback) → shared/constants/models.js (replaces 20× m.kind||m.type)
- getStatusVariant → shared/utils/connectionStatus.js (dedup ConnectionRow/ConnectionsCard)
- sseChunk → open-sse/utils/sse.js (dedup grok-web/perplexity-web)
- fetchWithTimeout → usage/shared.js (replace 4× AbortController pattern in google.js)
fix: enableObservability2 field name in requestDetailsRepo
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
import { BaseExecutor } from "./base.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
import { SSE_DONE, SSE_HEADERS_NO_BUFFER } from "../utils/sseConstants.js";
|
||||
import { sseChunk } from "../utils/sse.js";
|
||||
|
||||
const GROK_CHAT_API = PROVIDERS["grok-web"].baseUrl;
|
||||
const GROK_USER_AGENT = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/136.0.0.0 Safari/537.36";
|
||||
@@ -131,10 +132,6 @@ async function* extractContent(eventStream, isThinkingModel, signal) {
|
||||
yield { done: true, fingerprint, responseId };
|
||||
}
|
||||
|
||||
function sseChunk(data) {
|
||||
return `data: ${JSON.stringify(data)}\n\n`;
|
||||
}
|
||||
|
||||
function buildStreamingResponse(eventStream, model, cid, created, isThinkingModel, signal) {
|
||||
const encoder = new TextEncoder();
|
||||
return new ReadableStream({
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { BaseExecutor } from "./base.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
import { SSE_DONE, SSE_HEADERS_NO_BUFFER } from "../utils/sseConstants.js";
|
||||
import { sseChunk } from "../utils/sse.js";
|
||||
|
||||
const PPLX_SSE_ENDPOINT = PROVIDERS["perplexity-web"].baseUrl;
|
||||
const PPLX_API_VERSION = "2.18";
|
||||
@@ -290,10 +291,6 @@ async function* extractContent(eventStream, signal) {
|
||||
yield { delta: "", answer: fullAnswer, backendUuid: backendUuid ?? undefined, done: true };
|
||||
}
|
||||
|
||||
function sseChunk(data) {
|
||||
return `data: ${JSON.stringify(data)}\n\n`;
|
||||
}
|
||||
|
||||
function buildStreamingResponse(eventStream, model, cid, created, history, currentMsg, signal) {
|
||||
const encoder = new TextEncoder();
|
||||
return new ReadableStream({
|
||||
|
||||
@@ -1,73 +1,35 @@
|
||||
import { PROVIDERS, PROVIDER_OAUTH } from "../config/providers.js";
|
||||
import { OAUTH_ENDPOINTS, GITHUB_COPILOT, REFRESH_LEAD_MS } from "../config/appConstants.js";
|
||||
import { proxyAwareFetch } from "../utils/proxyFetch.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
import { OAUTH_ENDPOINTS, REFRESH_LEAD_MS } from "../config/appConstants.js";
|
||||
import {
|
||||
refreshXaiToken,
|
||||
refreshAccessToken,
|
||||
refreshClaudeOAuthToken,
|
||||
refreshGoogleToken,
|
||||
refreshQwenToken,
|
||||
refreshCodexToken,
|
||||
refreshKiroToken,
|
||||
refreshIflowToken,
|
||||
refreshGitHubToken,
|
||||
refreshCopilotToken,
|
||||
classifyOAuthRefreshError,
|
||||
} from "./tokenRefresh/providers.js";
|
||||
|
||||
// xAI refresh — wraps the class method from src/lib/oauth/services/xai.js so
|
||||
// the token-refresh switches below can stay flat (one function per provider).
|
||||
let _xaiServiceSingleton = null;
|
||||
async function refreshXaiToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("xai", refreshToken, async () => {
|
||||
try {
|
||||
if (!_xaiServiceSingleton) {
|
||||
const mod = await import("../../src/lib/oauth/services/xai.js");
|
||||
_xaiServiceSingleton = new mod.XaiService();
|
||||
}
|
||||
const tokens = await _xaiServiceSingleton.refreshAccessToken(refreshToken);
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
idToken: tokens.id_token,
|
||||
};
|
||||
} catch (e) {
|
||||
log?.warn?.("TOKEN_REFRESH", `xai refresh failed: ${e?.message || e}`);
|
||||
const msg = String(e?.message || "");
|
||||
if (msg.includes("invalid_grant") || msg.includes("invalid_request")) {
|
||||
return { error: "invalid_grant" };
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
// Re-export all provider refresh functions (preserves public API for all consumers)
|
||||
export {
|
||||
refreshAccessToken,
|
||||
refreshClaudeOAuthToken,
|
||||
refreshGoogleToken,
|
||||
refreshQwenToken,
|
||||
refreshCodexToken,
|
||||
refreshKiroToken,
|
||||
refreshIflowToken,
|
||||
refreshGitHubToken,
|
||||
refreshCopilotToken,
|
||||
classifyOAuthRefreshError,
|
||||
};
|
||||
|
||||
// Default token expiry buffer (refresh if expires within 5 minutes)
|
||||
export const TOKEN_EXPIRY_BUFFER_MS = 5 * 60 * 1000;
|
||||
|
||||
// Dedup: cache in-flight promise + recent result to prevent refresh_token_reused (Auth0 family revoke)
|
||||
const REFRESH_RESULT_TTL_MS = 10_000;
|
||||
const refreshDedupCache = new Map();
|
||||
|
||||
async function dedupRefresh(provider, oldToken, fn, log) {
|
||||
if (!oldToken) return fn();
|
||||
const key = `${provider}:${oldToken}`;
|
||||
const hit = refreshDedupCache.get(key);
|
||||
if (hit) {
|
||||
if (hit.promise) {
|
||||
log?.info?.("TOKEN_REFRESH", `Reusing in-flight refresh for ${provider}`);
|
||||
return hit.promise;
|
||||
}
|
||||
if (hit.expiresAt > Date.now()) {
|
||||
log?.info?.("TOKEN_REFRESH", `Reusing recent refresh result for ${provider}`);
|
||||
return hit.result;
|
||||
}
|
||||
refreshDedupCache.delete(key);
|
||||
}
|
||||
const promise = (async () => {
|
||||
try {
|
||||
const result = await fn();
|
||||
refreshDedupCache.set(key, { result, expiresAt: Date.now() + REFRESH_RESULT_TTL_MS });
|
||||
return result;
|
||||
} catch (err) {
|
||||
refreshDedupCache.delete(key);
|
||||
throw err;
|
||||
}
|
||||
})();
|
||||
refreshDedupCache.set(key, { promise });
|
||||
return promise;
|
||||
}
|
||||
|
||||
// Check if refresh result indicates unrecoverable error (caller should stop retry, force re-auth)
|
||||
export function isUnrecoverableRefreshError(result) {
|
||||
return (
|
||||
result &&
|
||||
@@ -79,543 +41,82 @@ export function isUnrecoverableRefreshError(result) {
|
||||
);
|
||||
}
|
||||
|
||||
// Get provider-specific refresh lead time, falls back to default buffer
|
||||
export function getRefreshLeadMs(provider) {
|
||||
return REFRESH_LEAD_MS[provider] || TOKEN_EXPIRY_BUFFER_MS;
|
||||
}
|
||||
|
||||
/**
|
||||
* Refresh OAuth access token using refresh token
|
||||
*/
|
||||
export async function refreshAccessToken(provider, refreshToken, credentials, log) {
|
||||
const config = PROVIDERS[provider];
|
||||
|
||||
if (!config || !config.refreshUrl) {
|
||||
log?.warn?.("TOKEN_REFRESH", `No refresh URL configured for provider: ${provider}`);
|
||||
return null;
|
||||
}
|
||||
|
||||
if (!refreshToken) {
|
||||
log?.warn?.("TOKEN_REFRESH", `No refresh token available for provider: ${provider}`);
|
||||
return null;
|
||||
}
|
||||
|
||||
return dedupRefresh(provider, refreshToken, async () => {
|
||||
export function parseVertexSaJson(apiKey) {
|
||||
if (typeof apiKey !== "string") return null;
|
||||
try {
|
||||
const response = await fetch(config.refreshUrl, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: new URLSearchParams({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: config.clientId,
|
||||
client_secret: config.clientSecret,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", `Failed to refresh token for ${provider}`, {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
const parsed = JSON.parse(apiKey);
|
||||
if (parsed.type === "service_account" && parsed.client_email && parsed.private_key && parsed.project_id) {
|
||||
return parsed;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", `Successfully refreshed token for ${provider}`, {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
};
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Error refreshing token for ${provider}`, {
|
||||
error: error.message,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Specialized refresh for Claude OAuth tokens
|
||||
*/
|
||||
export async function refreshClaudeOAuthToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("claude", refreshToken, async () => {
|
||||
try {
|
||||
const response = await fetch(OAUTH_ENDPOINTS.anthropic.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: PROVIDERS.claude.clientId,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Claude OAuth token", { status: response.status, error: errorText });
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Claude OAuth token", { hasNewAccessToken: !!tokens.access_token, expiresIn: tokens.expires_in });
|
||||
return { accessToken: tokens.access_token, refreshToken: tokens.refresh_token || refreshToken, expiresIn: tokens.expires_in };
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Network error refreshing Claude token: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Specialized refresh for Google providers (Gemini, Antigravity)
|
||||
*/
|
||||
export async function refreshGoogleToken(refreshToken, clientId, clientSecret, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh(`google:${clientId}`, refreshToken, async () => {
|
||||
try {
|
||||
const response = await fetch(OAUTH_ENDPOINTS.google.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: new URLSearchParams({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: clientId,
|
||||
client_secret: clientSecret,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Google token", { status: response.status, error: errorText });
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Google token", { hasNewAccessToken: !!tokens.access_token, expiresIn: tokens.expires_in });
|
||||
return { accessToken: tokens.access_token, refreshToken: tokens.refresh_token || refreshToken, expiresIn: tokens.expires_in };
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Network error refreshing Google token: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Specialized refresh for Qwen OAuth tokens
|
||||
*/
|
||||
export async function refreshQwenToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("qwen", refreshToken, async () => {
|
||||
const endpoint = OAUTH_ENDPOINTS.qwen.token;
|
||||
|
||||
try {
|
||||
const response = await fetch(endpoint, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: new URLSearchParams({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: PROVIDERS.qwen.clientId,
|
||||
}),
|
||||
});
|
||||
|
||||
if (response.status === 200) {
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Qwen token", {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
providerSpecificData: tokens.resource_url
|
||||
? { resourceUrl: tokens.resource_url }
|
||||
: undefined,
|
||||
};
|
||||
} else {
|
||||
const errorText = await response.text().catch(() => "");
|
||||
log?.warn?.("TOKEN_REFRESH", `Error with Qwen endpoint`, {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
}
|
||||
} catch (error) {
|
||||
log?.warn?.("TOKEN_REFRESH", `Network error trying Qwen endpoint`, {
|
||||
error: error.message,
|
||||
});
|
||||
}
|
||||
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Qwen token");
|
||||
return null;
|
||||
}, log);
|
||||
}
|
||||
|
||||
export function classifyOAuthRefreshError(errorText = "", status = 0) {
|
||||
let parsed = null;
|
||||
try {
|
||||
parsed = errorText ? JSON.parse(errorText) : null;
|
||||
} catch {
|
||||
parsed = null;
|
||||
}
|
||||
|
||||
const code = parsed?.error?.code || parsed?.error || parsed?.error_code || "";
|
||||
const description = parsed?.error_description || parsed?.message || errorText || "";
|
||||
const combined = `${code} ${description}`.toLowerCase();
|
||||
const permanent = [
|
||||
"refresh_token_expired",
|
||||
"refresh_token_reused",
|
||||
"refresh_token_invalidated",
|
||||
"invalid_grant",
|
||||
].some((marker) => combined.includes(marker));
|
||||
|
||||
return { status, code, description, permanent };
|
||||
}
|
||||
|
||||
/**
|
||||
* Specialized refresh for Codex (OpenAI) OAuth tokens.
|
||||
* OpenAI uses rotating (one-time-use) refresh tokens.
|
||||
* Returns { error: 'unrecoverable_refresh_error' } when token already consumed/invalid,
|
||||
* so callers stop retrying and request re-authentication.
|
||||
*/
|
||||
export async function refreshCodexToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("codex", refreshToken, async () => {
|
||||
try {
|
||||
const response = await fetch(OAUTH_ENDPOINTS.openai.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
client_id: PROVIDERS.codex.clientId,
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
const failure = classifyOAuthRefreshError(errorText, response.status);
|
||||
if (failure.permanent) {
|
||||
log?.error?.("TOKEN_REFRESH", "Codex refresh token already used or invalid. Re-auth required.", {
|
||||
status: response.status,
|
||||
code: failure.code,
|
||||
});
|
||||
return { error: "unrecoverable_refresh_error", code: failure.code };
|
||||
}
|
||||
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Codex token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
code: failure.code,
|
||||
permanent: failure.permanent,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Codex token", {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
hasIdToken: !!tokens.id_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
idToken: tokens.id_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
};
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Network error refreshing Codex token: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Specialized refresh for Kiro (AWS CodeWhisperer) tokens
|
||||
* Supports both AWS SSO OIDC (Builder ID/IDC) and Social Auth (Google/GitHub)
|
||||
*/
|
||||
// Backfill missing Kiro profileArn on refresh so existing IDC connections self-heal
|
||||
async function resolveKiroProfileArnPatch(providerSpecificData, accessToken, refreshedArn) {
|
||||
if (providerSpecificData?.profileArn) return {};
|
||||
let profileArn = refreshedArn?.trim?.() || null;
|
||||
if (!profileArn) {
|
||||
const { fetchKiroProfileArn } = await import("../../src/lib/oauth/providers.js");
|
||||
profileArn = await fetchKiroProfileArn(accessToken);
|
||||
}
|
||||
return profileArn ? { providerSpecificData: { profileArn } } : {};
|
||||
}
|
||||
|
||||
export async function refreshKiroToken(refreshToken, providerSpecificData, log, proxyOptions = null) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("kiro", refreshToken, async () => {
|
||||
const authMethod = providerSpecificData?.authMethod;
|
||||
const clientId = providerSpecificData?.clientId;
|
||||
const clientSecret = providerSpecificData?.clientSecret;
|
||||
const region = providerSpecificData?.region;
|
||||
|
||||
// AWS SSO OIDC (Builder ID or IDC)
|
||||
// If clientId and clientSecret exist, assume AWS SSO OIDC (default to builder-id if authMethod not specified)
|
||||
if (clientId && clientSecret) {
|
||||
const isIDC = authMethod === "idc";
|
||||
const endpoint = isIDC && region
|
||||
? `https://oidc.${region}.amazonaws.com/token`
|
||||
: "https://oidc.us-east-1.amazonaws.com/token";
|
||||
|
||||
const response = await proxyAwareFetch(endpoint, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
clientId: clientId,
|
||||
clientSecret: clientSecret,
|
||||
refreshToken: refreshToken,
|
||||
grantType: "refresh_token",
|
||||
}),
|
||||
}, proxyOptions);
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Kiro AWS token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Kiro AWS token", {
|
||||
hasNewAccessToken: !!tokens.accessToken,
|
||||
expiresIn: tokens.expiresIn,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.accessToken,
|
||||
refreshToken: tokens.refreshToken || refreshToken,
|
||||
expiresIn: tokens.expiresIn,
|
||||
...(await resolveKiroProfileArnPatch(providerSpecificData, tokens.accessToken, tokens.profileArn)),
|
||||
};
|
||||
}
|
||||
|
||||
// Social Auth (Google/GitHub) - use Kiro's refresh endpoint
|
||||
const response = await proxyAwareFetch(PROVIDERS.kiro.tokenUrl, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
"User-Agent": "kiro-cli/1.0.0",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
refreshToken: refreshToken,
|
||||
}),
|
||||
}, proxyOptions);
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Kiro social token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Kiro social token", {
|
||||
hasNewAccessToken: !!tokens.accessToken,
|
||||
expiresIn: tokens.expiresIn,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.accessToken,
|
||||
refreshToken: tokens.refreshToken || refreshToken,
|
||||
expiresIn: tokens.expiresIn,
|
||||
...(await resolveKiroProfileArnPatch(providerSpecificData, tokens.accessToken, tokens.profileArn)),
|
||||
};
|
||||
}, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Specialized refresh for iFlow OAuth tokens
|
||||
*/
|
||||
export async function refreshIflowToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("iflow", refreshToken, async () => {
|
||||
const basicAuth = btoa(`${PROVIDERS.iflow.clientId}:${PROVIDERS.iflow.clientSecret}`);
|
||||
// Cache Vertex tokens keyed by service account email { token, expiresAt }
|
||||
const vertexTokenCache = new Map();
|
||||
|
||||
const response = await fetch(OAUTH_ENDPOINTS.iflow.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
Authorization: `Basic ${basicAuth}`,
|
||||
},
|
||||
body: new URLSearchParams({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: PROVIDERS.iflow.clientId,
|
||||
client_secret: PROVIDERS.iflow.clientSecret,
|
||||
}),
|
||||
});
|
||||
export async function refreshVertexToken(saJson, log) {
|
||||
const cacheKey = saJson.client_email;
|
||||
const cached = vertexTokenCache.get(cacheKey);
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh iFlow token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
if (cached && cached.expiresAt - Date.now() > 5 * 60 * 1000) {
|
||||
return { accessToken: cached.token, expiresAt: cached.expiresAt };
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed iFlow token", {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
};
|
||||
}, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Specialized refresh for GitHub Copilot OAuth tokens
|
||||
*/
|
||||
export async function refreshGitHubToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("github", refreshToken, async () => {
|
||||
const params = {
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: PROVIDERS.github.clientId,
|
||||
};
|
||||
if (PROVIDERS.github.clientSecret) {
|
||||
params.client_secret = PROVIDERS.github.clientSecret;
|
||||
}
|
||||
|
||||
const response = await fetch(OAUTH_ENDPOINTS.github.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: new URLSearchParams(params),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh GitHub token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed GitHub token", {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
};
|
||||
}, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Refresh GitHub Copilot token using GitHub access token
|
||||
*/
|
||||
export async function refreshCopilotToken(githubAccessToken, log) {
|
||||
if (!githubAccessToken) return null;
|
||||
return dedupRefresh("copilot", githubAccessToken, async () => {
|
||||
try {
|
||||
const response = await fetch(PROVIDER_OAUTH["github"]?.copilotTokenUrl, {
|
||||
headers: {
|
||||
"Authorization": `token ${githubAccessToken}`,
|
||||
"User-Agent": GITHUB_COPILOT.USER_AGENT,
|
||||
"Editor-Version": `vscode/${GITHUB_COPILOT.VSCODE_VERSION}`,
|
||||
"Editor-Plugin-Version": `copilot-chat/${GITHUB_COPILOT.COPILOT_CHAT_VERSION}`,
|
||||
"Accept": "application/json",
|
||||
"x-github-api-version": GITHUB_COPILOT.API_VERSION
|
||||
}
|
||||
const { SignJWT, importPKCS8 } = await import("jose");
|
||||
log?.debug?.("TOKEN_REFRESH", `Vertex minting token for ${saJson.client_email}`);
|
||||
const privateKey = await importPKCS8(saJson.private_key.replace(/\\n/g, "\n"), "RS256");
|
||||
const now = Math.floor(Date.now() / 1000);
|
||||
|
||||
const jwt = await new SignJWT({ scope: "https://www.googleapis.com/auth/cloud-platform" })
|
||||
.setProtectedHeader({ alg: "RS256" })
|
||||
.setIssuer(saJson.client_email)
|
||||
.setAudience(OAUTH_ENDPOINTS.google.token)
|
||||
.setIssuedAt(now)
|
||||
.setExpirationTime(now + 3600)
|
||||
.sign(privateKey);
|
||||
|
||||
const res = await fetch(OAUTH_ENDPOINTS.google.token, {
|
||||
method: "POST",
|
||||
headers: { "Content-Type": "application/x-www-form-urlencoded" },
|
||||
body: new URLSearchParams({
|
||||
grant_type: "urn:ietf:params:oauth:grant-type:jwt-bearer",
|
||||
assertion: jwt,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Copilot token", {
|
||||
status: response.status,
|
||||
error: errorText
|
||||
});
|
||||
if (!res.ok) {
|
||||
const err = await res.text();
|
||||
log?.error?.("TOKEN_REFRESH", `Vertex token mint failed: ${err}`);
|
||||
return null;
|
||||
}
|
||||
|
||||
const data = await response.json();
|
||||
const { access_token, expires_in } = await res.json();
|
||||
const expiresAt = Date.now() + (expires_in ?? 3600) * 1000;
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Copilot token", {
|
||||
hasToken: !!data.token,
|
||||
expiresAt: data.expires_at
|
||||
});
|
||||
vertexTokenCache.set(cacheKey, { token: access_token, expiresAt });
|
||||
log?.info?.("TOKEN_REFRESH", `Vertex token minted for ${saJson.client_email}`);
|
||||
|
||||
return {
|
||||
token: data.token,
|
||||
expiresAt: data.expires_at
|
||||
};
|
||||
return { accessToken: access_token, expiresAt };
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", "Error refreshing Copilot token", {
|
||||
error: error.message
|
||||
});
|
||||
log?.error?.("TOKEN_REFRESH", `Vertex token error: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
// Single source of per-provider refresh dispatch (logic stays in each refreshXxx fn).
|
||||
// Each handler: (credentials, log) => Promise<tokens|null>
|
||||
function vertexRefreshHandler(c, log) {
|
||||
const saJson = parseVertexSaJson(c.apiKey);
|
||||
if (!saJson) return null;
|
||||
return refreshVertexToken(saJson, log);
|
||||
}
|
||||
|
||||
const REFRESH_HANDLERS = {
|
||||
"gemini-cli": (c, log) => refreshGoogleToken(c.refreshToken, PROVIDERS["gemini-cli"].clientId, PROVIDERS["gemini-cli"].clientSecret, log),
|
||||
antigravity: (c, log) => refreshGoogleToken(c.refreshToken, PROVIDERS.antigravity.clientId, PROVIDERS.antigravity.clientSecret, log),
|
||||
@@ -630,28 +131,15 @@ const REFRESH_HANDLERS = {
|
||||
"vertex-partner": vertexRefreshHandler
|
||||
};
|
||||
|
||||
function vertexRefreshHandler(c, log) {
|
||||
const saJson = parseVertexSaJson(c.apiKey);
|
||||
if (!saJson) return null;
|
||||
return refreshVertexToken(saJson, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Get access token for a specific provider (with in-flight dedup).
|
||||
* If a refresh is already in-flight for same provider+token, share the promise
|
||||
* to prevent parallel OAuth requests → Auth0 'refresh_token_reused' family revoke.
|
||||
*/
|
||||
export async function getAccessToken(provider, credentials, log) {
|
||||
if (!credentials || !credentials.refreshToken || typeof credentials.refreshToken !== "string") {
|
||||
log?.warn?.("TOKEN_REFRESH", `No valid refresh token available for provider: ${provider}`);
|
||||
return null;
|
||||
}
|
||||
// Dedup is handled inside each refreshXxxToken function
|
||||
return _getAccessTokenInternal(provider, credentials, log);
|
||||
}
|
||||
|
||||
async function _getAccessTokenInternal(provider, credentials, log) {
|
||||
// "gemini" shares Google refresh here (unlike refreshTokenByProvider)
|
||||
if (provider === "gemini") {
|
||||
return refreshGoogleToken(credentials.refreshToken, PROVIDERS.gemini.clientId, PROVIDERS.gemini.clientSecret, log);
|
||||
}
|
||||
@@ -663,19 +151,12 @@ async function _getAccessTokenInternal(provider, credentials, log) {
|
||||
return handler(credentials, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Refresh token by provider type (helper for handlers)
|
||||
*/
|
||||
export async function refreshTokenByProvider(provider, credentials, log) {
|
||||
if (!credentials.refreshToken) return null;
|
||||
const handler = REFRESH_HANDLERS[provider];
|
||||
// default: generic refresh (note: "gemini" is NOT special-cased here, by design)
|
||||
return handler ? handler(credentials, log) : refreshAccessToken(provider, credentials.refreshToken, credentials, log);
|
||||
}
|
||||
|
||||
/**
|
||||
* Format credentials for provider
|
||||
*/
|
||||
export function formatProviderCredentials(provider, credentials, log) {
|
||||
const config = PROVIDERS[provider];
|
||||
if (!config) {
|
||||
@@ -725,9 +206,6 @@ export function formatProviderCredentials(provider, credentials, log) {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get all access tokens for a user
|
||||
*/
|
||||
export async function getAllAccessTokens(userInfo, log) {
|
||||
const results = {};
|
||||
|
||||
@@ -748,89 +226,6 @@ export async function getAllAccessTokens(userInfo, log) {
|
||||
return results;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse Vertex AI Service Account JSON from apiKey string
|
||||
*/
|
||||
export function parseVertexSaJson(apiKey) {
|
||||
if (typeof apiKey !== "string") return null;
|
||||
try {
|
||||
const parsed = JSON.parse(apiKey);
|
||||
if (parsed.type === "service_account" && parsed.client_email && parsed.private_key && parsed.project_id) {
|
||||
return parsed;
|
||||
}
|
||||
return null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
// Cache Vertex tokens keyed by service account email { token, expiresAt }
|
||||
const vertexTokenCache = new Map();
|
||||
|
||||
/**
|
||||
* Mint a short-lived OAuth2 Bearer token for Google Cloud Vertex AI
|
||||
* using Service Account JSON + jose (RS256 JWT assertion flow).
|
||||
* Token is cached until 5 minutes before expiry.
|
||||
*/
|
||||
export async function refreshVertexToken(saJson, log) {
|
||||
const cacheKey = saJson.client_email;
|
||||
const cached = vertexTokenCache.get(cacheKey);
|
||||
|
||||
// Return cached token if still valid (5-min buffer)
|
||||
if (cached && cached.expiresAt - Date.now() > 5 * 60 * 1000) {
|
||||
return { accessToken: cached.token, expiresAt: cached.expiresAt };
|
||||
}
|
||||
|
||||
try {
|
||||
const { SignJWT, importPKCS8 } = await import("jose");
|
||||
log?.debug?.("TOKEN_REFRESH", `Vertex minting token for ${saJson.client_email}`);
|
||||
const privateKey = await importPKCS8(saJson.private_key.replace(/\\n/g, "\n"), "RS256");
|
||||
const now = Math.floor(Date.now() / 1000);
|
||||
|
||||
const jwt = await new SignJWT({ scope: "https://www.googleapis.com/auth/cloud-platform" })
|
||||
.setProtectedHeader({ alg: "RS256" })
|
||||
.setIssuer(saJson.client_email)
|
||||
.setAudience(OAUTH_ENDPOINTS.google.token)
|
||||
.setIssuedAt(now)
|
||||
.setExpirationTime(now + 3600)
|
||||
.sign(privateKey);
|
||||
|
||||
const res = await fetch(OAUTH_ENDPOINTS.google.token, {
|
||||
method: "POST",
|
||||
headers: { "Content-Type": "application/x-www-form-urlencoded" },
|
||||
body: new URLSearchParams({
|
||||
grant_type: "urn:ietf:params:oauth:grant-type:jwt-bearer",
|
||||
assertion: jwt,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!res.ok) {
|
||||
const err = await res.text();
|
||||
log?.error?.("TOKEN_REFRESH", `Vertex token mint failed: ${err}`);
|
||||
return null;
|
||||
}
|
||||
|
||||
const { access_token, expires_in } = await res.json();
|
||||
const expiresAt = Date.now() + (expires_in ?? 3600) * 1000;
|
||||
|
||||
vertexTokenCache.set(cacheKey, { token: access_token, expiresAt });
|
||||
log?.info?.("TOKEN_REFRESH", `Vertex token minted for ${saJson.client_email}`);
|
||||
|
||||
return { accessToken: access_token, expiresAt };
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Vertex token error: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Refresh token with retry and exponential backoff
|
||||
* Retries on failure with increasing delay: 1s, 2s, 3s...
|
||||
* @param {function} refreshFn - Async function that returns token or null
|
||||
* @param {number} maxRetries - Max retry attempts (default 3)
|
||||
* @param {object} log - Logger instance (optional)
|
||||
* @returns {Promise<object|null>} Token result or null if all retries fail
|
||||
*/
|
||||
export async function refreshWithRetry(refreshFn, maxRetries = 3, log = null) {
|
||||
for (let attempt = 0; attempt < maxRetries; attempt++) {
|
||||
if (attempt > 0) {
|
||||
|
||||
31
open-sse/services/tokenRefresh/dedup.js
Normal file
31
open-sse/services/tokenRefresh/dedup.js
Normal file
@@ -0,0 +1,31 @@
|
||||
const REFRESH_RESULT_TTL_MS = 10_000;
|
||||
const refreshDedupCache = new Map();
|
||||
|
||||
export async function dedupRefresh(provider, oldToken, fn, log) {
|
||||
if (!oldToken) return fn();
|
||||
const key = `${provider}:${oldToken}`;
|
||||
const hit = refreshDedupCache.get(key);
|
||||
if (hit) {
|
||||
if (hit.promise) {
|
||||
log?.info?.("TOKEN_REFRESH", `Reusing in-flight refresh for ${provider}`);
|
||||
return hit.promise;
|
||||
}
|
||||
if (hit.expiresAt > Date.now()) {
|
||||
log?.info?.("TOKEN_REFRESH", `Reusing recent refresh result for ${provider}`);
|
||||
return hit.result;
|
||||
}
|
||||
refreshDedupCache.delete(key);
|
||||
}
|
||||
const promise = (async () => {
|
||||
try {
|
||||
const result = await fn();
|
||||
refreshDedupCache.set(key, { result, expiresAt: Date.now() + REFRESH_RESULT_TTL_MS });
|
||||
return result;
|
||||
} catch (err) {
|
||||
refreshDedupCache.delete(key);
|
||||
throw err;
|
||||
}
|
||||
})();
|
||||
refreshDedupCache.set(key, { promise });
|
||||
return promise;
|
||||
}
|
||||
526
open-sse/services/tokenRefresh/providers.js
Normal file
526
open-sse/services/tokenRefresh/providers.js
Normal file
@@ -0,0 +1,526 @@
|
||||
import { PROVIDERS, PROVIDER_OAUTH } from "../../config/providers.js";
|
||||
import { OAUTH_ENDPOINTS, GITHUB_COPILOT } from "../../config/appConstants.js";
|
||||
import { proxyAwareFetch } from "../../utils/proxyFetch.js";
|
||||
import { dedupRefresh } from "./dedup.js";
|
||||
|
||||
let _xaiServiceSingleton = null;
|
||||
export async function refreshXaiToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("xai", refreshToken, async () => {
|
||||
try {
|
||||
if (!_xaiServiceSingleton) {
|
||||
const mod = await import("../../../src/lib/oauth/services/xai.js");
|
||||
_xaiServiceSingleton = new mod.XaiService();
|
||||
}
|
||||
const tokens = await _xaiServiceSingleton.refreshAccessToken(refreshToken);
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
idToken: tokens.id_token,
|
||||
};
|
||||
} catch (e) {
|
||||
log?.warn?.("TOKEN_REFRESH", `xai refresh failed: ${e?.message || e}`);
|
||||
const msg = String(e?.message || "");
|
||||
if (msg.includes("invalid_grant") || msg.includes("invalid_request")) {
|
||||
return { error: "invalid_grant" };
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
export async function refreshAccessToken(provider, refreshToken, credentials, log) {
|
||||
const config = PROVIDERS[provider];
|
||||
|
||||
if (!config || !config.refreshUrl) {
|
||||
log?.warn?.("TOKEN_REFRESH", `No refresh URL configured for provider: ${provider}`);
|
||||
return null;
|
||||
}
|
||||
|
||||
if (!refreshToken) {
|
||||
log?.warn?.("TOKEN_REFRESH", `No refresh token available for provider: ${provider}`);
|
||||
return null;
|
||||
}
|
||||
|
||||
return dedupRefresh(provider, refreshToken, async () => {
|
||||
try {
|
||||
const response = await fetch(config.refreshUrl, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: new URLSearchParams({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: config.clientId,
|
||||
client_secret: config.clientSecret,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", `Failed to refresh token for ${provider}`, {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", `Successfully refreshed token for ${provider}`, {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
};
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Error refreshing token for ${provider}`, {
|
||||
error: error.message,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
export async function refreshClaudeOAuthToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("claude", refreshToken, async () => {
|
||||
try {
|
||||
const response = await fetch(OAUTH_ENDPOINTS.anthropic.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: PROVIDERS.claude.clientId,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Claude OAuth token", { status: response.status, error: errorText });
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Claude OAuth token", { hasNewAccessToken: !!tokens.access_token, expiresIn: tokens.expires_in });
|
||||
return { accessToken: tokens.access_token, refreshToken: tokens.refresh_token || refreshToken, expiresIn: tokens.expires_in };
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Network error refreshing Claude token: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
export async function refreshGoogleToken(refreshToken, clientId, clientSecret, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh(`google:${clientId}`, refreshToken, async () => {
|
||||
try {
|
||||
const response = await fetch(OAUTH_ENDPOINTS.google.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: new URLSearchParams({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: clientId,
|
||||
client_secret: clientSecret,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Google token", { status: response.status, error: errorText });
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Google token", { hasNewAccessToken: !!tokens.access_token, expiresIn: tokens.expires_in });
|
||||
return { accessToken: tokens.access_token, refreshToken: tokens.refresh_token || refreshToken, expiresIn: tokens.expires_in };
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Network error refreshing Google token: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
export async function refreshQwenToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("qwen", refreshToken, async () => {
|
||||
const endpoint = OAUTH_ENDPOINTS.qwen.token;
|
||||
|
||||
try {
|
||||
const response = await fetch(endpoint, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: new URLSearchParams({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: PROVIDERS.qwen.clientId,
|
||||
}),
|
||||
});
|
||||
|
||||
if (response.status === 200) {
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Qwen token", {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
providerSpecificData: tokens.resource_url
|
||||
? { resourceUrl: tokens.resource_url }
|
||||
: undefined,
|
||||
};
|
||||
} else {
|
||||
const errorText = await response.text().catch(() => "");
|
||||
log?.warn?.("TOKEN_REFRESH", `Error with Qwen endpoint`, {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
}
|
||||
} catch (error) {
|
||||
log?.warn?.("TOKEN_REFRESH", `Network error trying Qwen endpoint`, {
|
||||
error: error.message,
|
||||
});
|
||||
}
|
||||
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Qwen token");
|
||||
return null;
|
||||
}, log);
|
||||
}
|
||||
|
||||
export function classifyOAuthRefreshError(errorText = "", status = 0) {
|
||||
let parsed = null;
|
||||
try {
|
||||
parsed = errorText ? JSON.parse(errorText) : null;
|
||||
} catch {
|
||||
parsed = null;
|
||||
}
|
||||
|
||||
const code = parsed?.error?.code || parsed?.error || parsed?.error_code || "";
|
||||
const description = parsed?.error_description || parsed?.message || errorText || "";
|
||||
const combined = `${code} ${description}`.toLowerCase();
|
||||
const permanent = [
|
||||
"refresh_token_expired",
|
||||
"refresh_token_reused",
|
||||
"refresh_token_invalidated",
|
||||
"invalid_grant",
|
||||
].some((marker) => combined.includes(marker));
|
||||
|
||||
return { status, code, description, permanent };
|
||||
}
|
||||
|
||||
export async function refreshCodexToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("codex", refreshToken, async () => {
|
||||
try {
|
||||
const response = await fetch(OAUTH_ENDPOINTS.openai.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
client_id: PROVIDERS.codex.clientId,
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
const failure = classifyOAuthRefreshError(errorText, response.status);
|
||||
if (failure.permanent) {
|
||||
log?.error?.("TOKEN_REFRESH", "Codex refresh token already used or invalid. Re-auth required.", {
|
||||
status: response.status,
|
||||
code: failure.code,
|
||||
});
|
||||
return { error: "unrecoverable_refresh_error", code: failure.code };
|
||||
}
|
||||
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Codex token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
code: failure.code,
|
||||
permanent: failure.permanent,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Codex token", {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
hasIdToken: !!tokens.id_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
idToken: tokens.id_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
};
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", `Network error refreshing Codex token: ${error.message}`);
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
|
||||
async function resolveKiroProfileArnPatch(providerSpecificData, accessToken, refreshedArn) {
|
||||
if (providerSpecificData?.profileArn) return {};
|
||||
let profileArn = refreshedArn?.trim?.() || null;
|
||||
if (!profileArn) {
|
||||
const { fetchKiroProfileArn } = await import("../../../src/lib/oauth/providers.js");
|
||||
profileArn = await fetchKiroProfileArn(accessToken);
|
||||
}
|
||||
return profileArn ? { providerSpecificData: { profileArn } } : {};
|
||||
}
|
||||
|
||||
export async function refreshKiroToken(refreshToken, providerSpecificData, log, proxyOptions = null) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("kiro", refreshToken, async () => {
|
||||
const authMethod = providerSpecificData?.authMethod;
|
||||
const clientId = providerSpecificData?.clientId;
|
||||
const clientSecret = providerSpecificData?.clientSecret;
|
||||
const region = providerSpecificData?.region;
|
||||
|
||||
if (clientId && clientSecret) {
|
||||
const isIDC = authMethod === "idc";
|
||||
const endpoint = isIDC && region
|
||||
? `https://oidc.${region}.amazonaws.com/token`
|
||||
: "https://oidc.us-east-1.amazonaws.com/token";
|
||||
|
||||
const response = await proxyAwareFetch(endpoint, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
clientId: clientId,
|
||||
clientSecret: clientSecret,
|
||||
refreshToken: refreshToken,
|
||||
grantType: "refresh_token",
|
||||
}),
|
||||
}, proxyOptions);
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Kiro AWS token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Kiro AWS token", {
|
||||
hasNewAccessToken: !!tokens.accessToken,
|
||||
expiresIn: tokens.expiresIn,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.accessToken,
|
||||
refreshToken: tokens.refreshToken || refreshToken,
|
||||
expiresIn: tokens.expiresIn,
|
||||
...(await resolveKiroProfileArnPatch(providerSpecificData, tokens.accessToken, tokens.profileArn)),
|
||||
};
|
||||
}
|
||||
|
||||
const response = await proxyAwareFetch(PROVIDERS.kiro.tokenUrl, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
"User-Agent": "kiro-cli/1.0.0",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
refreshToken: refreshToken,
|
||||
}),
|
||||
}, proxyOptions);
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Kiro social token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Kiro social token", {
|
||||
hasNewAccessToken: !!tokens.accessToken,
|
||||
expiresIn: tokens.expiresIn,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.accessToken,
|
||||
refreshToken: tokens.refreshToken || refreshToken,
|
||||
expiresIn: tokens.expiresIn,
|
||||
...(await resolveKiroProfileArnPatch(providerSpecificData, tokens.accessToken, tokens.profileArn)),
|
||||
};
|
||||
}, log);
|
||||
}
|
||||
|
||||
export async function refreshIflowToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("iflow", refreshToken, async () => {
|
||||
const basicAuth = btoa(`${PROVIDERS.iflow.clientId}:${PROVIDERS.iflow.clientSecret}`);
|
||||
|
||||
const response = await fetch(OAUTH_ENDPOINTS.iflow.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
Authorization: `Basic ${basicAuth}`,
|
||||
},
|
||||
body: new URLSearchParams({
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: PROVIDERS.iflow.clientId,
|
||||
client_secret: PROVIDERS.iflow.clientSecret,
|
||||
}),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh iFlow token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed iFlow token", {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
};
|
||||
}, log);
|
||||
}
|
||||
|
||||
export async function refreshGitHubToken(refreshToken, log) {
|
||||
if (!refreshToken) return null;
|
||||
return dedupRefresh("github", refreshToken, async () => {
|
||||
const params = {
|
||||
grant_type: "refresh_token",
|
||||
refresh_token: refreshToken,
|
||||
client_id: PROVIDERS.github.clientId,
|
||||
};
|
||||
if (PROVIDERS.github.clientSecret) {
|
||||
params.client_secret = PROVIDERS.github.clientSecret;
|
||||
}
|
||||
|
||||
const response = await fetch(OAUTH_ENDPOINTS.github.token, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
Accept: "application/json",
|
||||
},
|
||||
body: new URLSearchParams(params),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh GitHub token", {
|
||||
status: response.status,
|
||||
error: errorText,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const tokens = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed GitHub token", {
|
||||
hasNewAccessToken: !!tokens.access_token,
|
||||
hasNewRefreshToken: !!tokens.refresh_token,
|
||||
expiresIn: tokens.expires_in,
|
||||
});
|
||||
|
||||
return {
|
||||
accessToken: tokens.access_token,
|
||||
refreshToken: tokens.refresh_token || refreshToken,
|
||||
expiresIn: tokens.expires_in,
|
||||
};
|
||||
}, log);
|
||||
}
|
||||
|
||||
export async function refreshCopilotToken(githubAccessToken, log) {
|
||||
if (!githubAccessToken) return null;
|
||||
return dedupRefresh("copilot", githubAccessToken, async () => {
|
||||
try {
|
||||
const response = await fetch(PROVIDER_OAUTH["github"]?.copilotTokenUrl, {
|
||||
headers: {
|
||||
"Authorization": `token ${githubAccessToken}`,
|
||||
"User-Agent": GITHUB_COPILOT.USER_AGENT,
|
||||
"Editor-Version": `vscode/${GITHUB_COPILOT.VSCODE_VERSION}`,
|
||||
"Editor-Plugin-Version": `copilot-chat/${GITHUB_COPILOT.COPILOT_CHAT_VERSION}`,
|
||||
"Accept": "application/json",
|
||||
"x-github-api-version": GITHUB_COPILOT.API_VERSION
|
||||
}
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
log?.error?.("TOKEN_REFRESH", "Failed to refresh Copilot token", {
|
||||
status: response.status,
|
||||
error: errorText
|
||||
});
|
||||
return null;
|
||||
}
|
||||
|
||||
const data = await response.json();
|
||||
|
||||
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Copilot token", {
|
||||
hasToken: !!data.token,
|
||||
expiresAt: data.expires_at
|
||||
});
|
||||
|
||||
return {
|
||||
token: data.token,
|
||||
expiresAt: data.expires_at
|
||||
};
|
||||
} catch (error) {
|
||||
log?.error?.("TOKEN_REFRESH", "Error refreshing Copilot token", {
|
||||
error: error.message
|
||||
});
|
||||
return null;
|
||||
}
|
||||
}, log);
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
134
open-sse/services/usage/claude.js
Normal file
134
open-sse/services/usage/claude.js
Normal file
@@ -0,0 +1,134 @@
|
||||
/**
|
||||
* Claude usage handler
|
||||
*/
|
||||
|
||||
import { proxyAwareFetch } from "../../utils/proxyFetch.js";
|
||||
import { ANTHROPIC_API_VERSION } from "../../providers/shared.js";
|
||||
import { U, parseResetTime } from "./shared.js";
|
||||
|
||||
// Claude API config (urls from registry, apiVersion is header logic kept here)
|
||||
const CLAUDE_CONFIG = {
|
||||
oauthUsageUrl: U("claude").oauthUrl,
|
||||
usageUrl: U("claude").orgUrl,
|
||||
settingsUrl: U("claude").settingsUrl,
|
||||
apiVersion: ANTHROPIC_API_VERSION,
|
||||
};
|
||||
|
||||
/**
|
||||
* Claude Usage - Primary: OAuth endpoint, Fallback: legacy settings/org endpoint
|
||||
*/
|
||||
export async function getClaudeUsage(accessToken, proxyOptions = null) {
|
||||
try {
|
||||
// Primary: OAuth usage endpoint (Claude Code consumer OAuth tokens)
|
||||
const oauthResponse = await proxyAwareFetch(CLAUDE_CONFIG.oauthUsageUrl, {
|
||||
method: "GET",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"anthropic-beta": "oauth-2025-04-20",
|
||||
"anthropic-version": CLAUDE_CONFIG.apiVersion,
|
||||
},
|
||||
}, proxyOptions);
|
||||
|
||||
if (oauthResponse.ok) {
|
||||
const data = await oauthResponse.json();
|
||||
const quotas = {};
|
||||
|
||||
// utilization = % USED (e.g. 87 means 87% used, 13% remaining)
|
||||
const hasUtilization = (window) =>
|
||||
window && typeof window === "object" && typeof window.utilization === "number";
|
||||
|
||||
const createQuotaObject = (window) => {
|
||||
const used = window.utilization;
|
||||
const remaining = Math.max(0, 100 - used);
|
||||
return {
|
||||
used,
|
||||
total: 100,
|
||||
remaining,
|
||||
remainingPercentage: remaining,
|
||||
resetAt: parseResetTime(window.resets_at),
|
||||
unlimited: false,
|
||||
};
|
||||
};
|
||||
|
||||
if (hasUtilization(data.five_hour)) {
|
||||
quotas["session (5h)"] = createQuotaObject(data.five_hour);
|
||||
}
|
||||
|
||||
if (hasUtilization(data.seven_day)) {
|
||||
quotas["weekly (7d)"] = createQuotaObject(data.seven_day);
|
||||
}
|
||||
|
||||
// Parse model-specific weekly windows (e.g. seven_day_sonnet, seven_day_opus)
|
||||
for (const [key, value] of Object.entries(data)) {
|
||||
if (key.startsWith("seven_day_") && key !== "seven_day" && hasUtilization(value)) {
|
||||
const modelName = key.replace("seven_day_", "");
|
||||
quotas[`weekly ${modelName} (7d)`] = createQuotaObject(value);
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
plan: "Claude Code",
|
||||
extraUsage: data.extra_usage ?? null,
|
||||
quotas,
|
||||
};
|
||||
}
|
||||
|
||||
// Fallback: legacy settings + org usage endpoint
|
||||
console.warn(`[Claude Usage] OAuth endpoint returned ${oauthResponse.status}, falling back to legacy`);
|
||||
return await getClaudeUsageLegacy(accessToken, proxyOptions);
|
||||
} catch (error) {
|
||||
return { message: `Claude connected. Unable to fetch usage: ${error.message}` };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Legacy Claude usage for API key / org admin users
|
||||
*/
|
||||
async function getClaudeUsageLegacy(accessToken, proxyOptions = null) {
|
||||
try {
|
||||
const settingsResponse = await proxyAwareFetch(CLAUDE_CONFIG.settingsUrl, {
|
||||
method: "GET",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"anthropic-version": CLAUDE_CONFIG.apiVersion,
|
||||
},
|
||||
}, proxyOptions);
|
||||
|
||||
if (settingsResponse.ok) {
|
||||
const settings = await settingsResponse.json();
|
||||
|
||||
if (settings.organization_id) {
|
||||
const usageResponse = await proxyAwareFetch(
|
||||
CLAUDE_CONFIG.usageUrl.replace("{org_id}", settings.organization_id),
|
||||
{
|
||||
method: "GET",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"anthropic-version": CLAUDE_CONFIG.apiVersion,
|
||||
},
|
||||
},
|
||||
proxyOptions
|
||||
);
|
||||
|
||||
if (usageResponse.ok) {
|
||||
const usage = await usageResponse.json();
|
||||
return {
|
||||
plan: settings.plan || "Unknown",
|
||||
organization: settings.organization_name,
|
||||
quotas: usage,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
plan: settings.plan || "Unknown",
|
||||
organization: settings.organization_name,
|
||||
message: "Claude connected. Usage details require admin access.",
|
||||
};
|
||||
}
|
||||
|
||||
return { message: "Claude connected. Usage API requires admin permissions." };
|
||||
} catch (error) {
|
||||
return { message: `Claude connected. Unable to fetch usage: ${error.message}` };
|
||||
}
|
||||
}
|
||||
99
open-sse/services/usage/codex.js
Normal file
99
open-sse/services/usage/codex.js
Normal file
@@ -0,0 +1,99 @@
|
||||
/**
|
||||
* Codex (OpenAI) usage handler
|
||||
*/
|
||||
|
||||
import { proxyAwareFetch } from "../../utils/proxyFetch.js";
|
||||
import { U, parseResetTime, toFiniteNumber } from "./shared.js";
|
||||
|
||||
// Codex (OpenAI) API config
|
||||
const CODEX_CONFIG = {
|
||||
usageUrl: U("codex").url,
|
||||
};
|
||||
|
||||
function getCodexRateLimitBody(snapshot) {
|
||||
if (!snapshot || typeof snapshot !== "object" || Array.isArray(snapshot)) return null;
|
||||
return snapshot.rate_limit && typeof snapshot.rate_limit === "object"
|
||||
? snapshot.rate_limit
|
||||
: snapshot;
|
||||
}
|
||||
|
||||
function formatCodexWindow(window) {
|
||||
const used = Math.max(0, Math.min(100, toFiniteNumber(window?.used_percent ?? window?.percent_used, 0)));
|
||||
return {
|
||||
used,
|
||||
total: 100,
|
||||
remaining: Math.max(0, 100 - used),
|
||||
resetAt: parseResetTime(window?.reset_at ?? window?.resets_at ?? window?.resetAt ?? null),
|
||||
unlimited: false,
|
||||
};
|
||||
}
|
||||
|
||||
function appendCodexQuotaWindows(quotas, prefix, snapshot) {
|
||||
const rateLimit = getCodexRateLimitBody(snapshot);
|
||||
if (!rateLimit) return false;
|
||||
|
||||
const primary = rateLimit.primary_window || rateLimit.primary || snapshot.primary_window || snapshot.primary;
|
||||
const secondary = rateLimit.secondary_window || rateLimit.secondary || snapshot.secondary_window || snapshot.secondary;
|
||||
let added = false;
|
||||
|
||||
if (primary) {
|
||||
quotas[prefix ? `${prefix}_session` : "session"] = formatCodexWindow(primary);
|
||||
added = true;
|
||||
}
|
||||
if (secondary) {
|
||||
quotas[prefix ? `${prefix}_weekly` : "weekly"] = formatCodexWindow(secondary);
|
||||
added = true;
|
||||
}
|
||||
|
||||
return added;
|
||||
}
|
||||
|
||||
function getCodexReviewRateLimit(data) {
|
||||
if (data.code_review_rate_limit || data.review_rate_limit) {
|
||||
return data.code_review_rate_limit || data.review_rate_limit;
|
||||
}
|
||||
|
||||
const byLimitId = data.rate_limits_by_limit_id;
|
||||
if (byLimitId && typeof byLimitId === "object" && !Array.isArray(byLimitId)) {
|
||||
return byLimitId.code_review || byLimitId.codex_review || byLimitId.review || null;
|
||||
}
|
||||
|
||||
const additional = Array.isArray(data.additional_rate_limits) ? data.additional_rate_limits : [];
|
||||
return additional.find((entry) => {
|
||||
const id = String(entry?.limit_name || entry?.metered_feature || entry?.id || "").toLowerCase();
|
||||
return id === "code_review" || id === "codex_review" || id === "review" || id.includes("review");
|
||||
}) || null;
|
||||
}
|
||||
|
||||
export async function getCodexUsage(accessToken, proxyOptions = null) {
|
||||
try {
|
||||
const response = await proxyAwareFetch(CODEX_CONFIG.usageUrl, {
|
||||
method: "GET",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"Accept": "application/json",
|
||||
},
|
||||
}, proxyOptions);
|
||||
|
||||
if (!response.ok) {
|
||||
return { message: `Codex connected. Usage API temporarily unavailable (${response.status}).` };
|
||||
}
|
||||
|
||||
const data = await response.json();
|
||||
const normalRateLimit = data.rate_limit || data.rate_limits || data.rate_limits_by_limit_id?.codex || {};
|
||||
const reviewRateLimit = getCodexReviewRateLimit(data);
|
||||
const quotas = {};
|
||||
|
||||
appendCodexQuotaWindows(quotas, "", normalRateLimit);
|
||||
appendCodexQuotaWindows(quotas, "review", reviewRateLimit);
|
||||
|
||||
return {
|
||||
plan: data.plan_type || data.summary?.plan || "unknown",
|
||||
limitReached: getCodexRateLimitBody(normalRateLimit)?.limit_reached || false,
|
||||
reviewLimitReached: getCodexRateLimitBody(reviewRateLimit)?.limit_reached || false,
|
||||
quotas,
|
||||
};
|
||||
} catch (error) {
|
||||
throw new Error(`Failed to fetch Codex usage: ${error.message}`);
|
||||
}
|
||||
}
|
||||
100
open-sse/services/usage/github.js
Normal file
100
open-sse/services/usage/github.js
Normal file
@@ -0,0 +1,100 @@
|
||||
/**
|
||||
* GitHub Copilot usage handler
|
||||
*/
|
||||
|
||||
import { proxyAwareFetch } from "../../utils/proxyFetch.js";
|
||||
import { PROVIDER_OAUTH } from "../../providers/index.js";
|
||||
import { U, parseResetTime } from "./shared.js";
|
||||
|
||||
// GitHub API config — single source from registry oauth block
|
||||
const GITHUB_CONFIG = {
|
||||
apiVersion: PROVIDER_OAUTH.github?.apiVersion,
|
||||
userAgent: PROVIDER_OAUTH.github?.userAgent,
|
||||
};
|
||||
|
||||
/**
|
||||
* GitHub Copilot Usage
|
||||
* Uses GitHub accessToken (not copilotToken) to call copilot_internal/user API
|
||||
*/
|
||||
export async function getGitHubUsage(accessToken, providerSpecificData, proxyOptions = null) {
|
||||
try {
|
||||
if (!accessToken) {
|
||||
throw new Error("No GitHub access token available. Please re-authorize the connection.");
|
||||
}
|
||||
|
||||
// copilot_internal/user API requires GitHub OAuth token, not copilotToken
|
||||
const response = await proxyAwareFetch(U("github").url, {
|
||||
headers: {
|
||||
"Authorization": `token ${accessToken}`,
|
||||
"Accept": "application/json",
|
||||
"X-GitHub-Api-Version": GITHUB_CONFIG.apiVersion,
|
||||
"User-Agent": GITHUB_CONFIG.userAgent,
|
||||
"Editor-Version": "vscode/1.100.0",
|
||||
"Editor-Plugin-Version": "copilot-chat/0.26.7",
|
||||
},
|
||||
}, proxyOptions);
|
||||
|
||||
if (!response.ok) {
|
||||
const error = await response.text();
|
||||
throw new Error(`GitHub API error: ${error}`);
|
||||
}
|
||||
|
||||
const data = await response.json();
|
||||
|
||||
// Handle different response formats (paid vs free)
|
||||
if (data.quota_snapshots) {
|
||||
// Paid plan format
|
||||
const snapshots = data.quota_snapshots;
|
||||
const resetAt = parseResetTime(data.quota_reset_date);
|
||||
|
||||
return {
|
||||
plan: data.copilot_plan,
|
||||
resetDate: data.quota_reset_date,
|
||||
quotas: {
|
||||
chat: { ...formatGitHubQuotaSnapshot(snapshots.chat), resetAt },
|
||||
completions: { ...formatGitHubQuotaSnapshot(snapshots.completions), resetAt },
|
||||
premium_interactions: { ...formatGitHubQuotaSnapshot(snapshots.premium_interactions), resetAt },
|
||||
},
|
||||
};
|
||||
} else if (data.monthly_quotas || data.limited_user_quotas) {
|
||||
// Free/limited plan format
|
||||
const monthlyQuotas = data.monthly_quotas || {};
|
||||
const usedQuotas = data.limited_user_quotas || {};
|
||||
const resetAt = parseResetTime(data.limited_user_reset_date);
|
||||
|
||||
return {
|
||||
plan: data.copilot_plan || data.access_type_sku,
|
||||
resetDate: data.limited_user_reset_date,
|
||||
quotas: {
|
||||
chat: {
|
||||
used: usedQuotas.chat || 0,
|
||||
total: monthlyQuotas.chat || 0,
|
||||
unlimited: false,
|
||||
resetAt,
|
||||
},
|
||||
completions: {
|
||||
used: usedQuotas.completions || 0,
|
||||
total: monthlyQuotas.completions || 0,
|
||||
unlimited: false,
|
||||
resetAt,
|
||||
},
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
return { message: "GitHub Copilot connected. Unable to parse quota data." };
|
||||
} catch (error) {
|
||||
throw new Error(`Failed to fetch GitHub usage: ${error.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
function formatGitHubQuotaSnapshot(quota) {
|
||||
if (!quota) return { used: 0, total: 0, unlimited: true };
|
||||
|
||||
return {
|
||||
used: quota.entitlement - quota.remaining,
|
||||
total: quota.entitlement,
|
||||
remaining: quota.remaining,
|
||||
unlimited: quota.unlimited || false,
|
||||
};
|
||||
}
|
||||
240
open-sse/services/usage/google.js
Normal file
240
open-sse/services/usage/google.js
Normal file
@@ -0,0 +1,240 @@
|
||||
/**
|
||||
* Google usage handlers (Gemini CLI + Antigravity)
|
||||
*/
|
||||
|
||||
import { CLIENT_METADATA, getPlatformUserAgent } from "../../config/appConstants.js";
|
||||
import { ANTIGRAVITY_OAUTH_CLIENT } from "../../providers/shared.js";
|
||||
import { U, parseResetTime, normalizeCloudCodeProjectId, fetchWithTimeout } from "./shared.js";
|
||||
|
||||
// Antigravity API config (from Quotio) — urls from registry, oauth client + dynamic UA kept here
|
||||
const ANTIGRAVITY_CONFIG = {
|
||||
...U("antigravity"),
|
||||
...ANTIGRAVITY_OAUTH_CLIENT,
|
||||
userAgent: getPlatformUserAgent(),
|
||||
};
|
||||
|
||||
/**
|
||||
* Gemini CLI Usage — fetch per-model quota via Cloud Code Assist API.
|
||||
* Uses retrieveUserQuota (same endpoint as `gemini /stats`) returning
|
||||
* per-model buckets with remainingFraction + resetTime.
|
||||
*/
|
||||
export async function getGeminiUsage(accessToken, providerSpecificData, proxyOptions = null) {
|
||||
if (!accessToken) {
|
||||
return { plan: "Free", message: "Gemini CLI access token not available." };
|
||||
}
|
||||
|
||||
try {
|
||||
// Resolve project id: prefer connection-stored id, else loadCodeAssist lookup.
|
||||
// #1271: OAuth save stores projectId on the connection, not providerSpecificData.
|
||||
let projectId = normalizeCloudCodeProjectId(providerSpecificData?.projectId);
|
||||
let plan = "Free";
|
||||
|
||||
if (!projectId) {
|
||||
const subInfo = await getGeminiSubscriptionInfo(accessToken, proxyOptions);
|
||||
projectId = normalizeCloudCodeProjectId(subInfo?.cloudaicompanionProject);
|
||||
plan = subInfo?.currentTier?.name || plan;
|
||||
}
|
||||
|
||||
if (!projectId) {
|
||||
return {
|
||||
plan,
|
||||
message: "Gemini CLI project ID not available. Reconnect Gemini CLI, or configure a Google Cloud project with Gemini Code Assist access before checking quota.",
|
||||
};
|
||||
}
|
||||
|
||||
const response = await fetchWithTimeout(
|
||||
U("gemini-cli").quotaUrl,
|
||||
{
|
||||
method: "POST",
|
||||
headers: {
|
||||
Authorization: `Bearer ${accessToken}`,
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
body: JSON.stringify({ project: projectId }),
|
||||
},
|
||||
10000,
|
||||
proxyOptions
|
||||
);
|
||||
|
||||
if (!response.ok) {
|
||||
return { plan, message: `Gemini CLI quota error (${response.status}).` };
|
||||
}
|
||||
|
||||
const data = await response.json();
|
||||
const quotas = {};
|
||||
|
||||
if (Array.isArray(data.buckets)) {
|
||||
for (const bucket of data.buckets) {
|
||||
if (!bucket.modelId || bucket.remainingFraction == null) continue;
|
||||
|
||||
const remainingFraction = Number(bucket.remainingFraction) || 0;
|
||||
const total = 1000; // Normalized base, matches antigravity convention
|
||||
const remaining = Math.round(total * remainingFraction);
|
||||
const used = Math.max(0, total - remaining);
|
||||
|
||||
quotas[bucket.modelId] = {
|
||||
used,
|
||||
total,
|
||||
resetAt: parseResetTime(bucket.resetTime),
|
||||
remainingPercentage: remainingFraction * 100,
|
||||
unlimited: false,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
return { plan, quotas };
|
||||
} catch (error) {
|
||||
return { message: `Gemini CLI error: ${error.message}` };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get Gemini CLI subscription info via loadCodeAssist
|
||||
*/
|
||||
async function getGeminiSubscriptionInfo(accessToken, proxyOptions = null) {
|
||||
try {
|
||||
const response = await fetchWithTimeout(
|
||||
U("gemini-cli").loadCodeAssistUrl,
|
||||
{
|
||||
method: "POST",
|
||||
headers: {
|
||||
Authorization: `Bearer ${accessToken}`,
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
body: JSON.stringify({ metadata: CLIENT_METADATA }),
|
||||
},
|
||||
10000,
|
||||
proxyOptions
|
||||
);
|
||||
if (!response.ok) return null;
|
||||
return await response.json();
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Antigravity Usage - Fetch quota from Google Cloud Code API
|
||||
*/
|
||||
export async function getAntigravityUsage(accessToken, providerSpecificData, proxyOptions = null) {
|
||||
try {
|
||||
// Fetch subscription info once — reuse for both projectId and plan
|
||||
const subscriptionInfo = await getAntigravitySubscriptionInfo(accessToken, proxyOptions);
|
||||
const projectId = subscriptionInfo?.cloudaicompanionProject || null;
|
||||
|
||||
const response = await fetchWithTimeout(ANTIGRAVITY_CONFIG.quotaApiUrl, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"User-Agent": ANTIGRAVITY_CONFIG.userAgent,
|
||||
"Content-Type": "application/json",
|
||||
"X-Client-Name": "antigravity",
|
||||
"X-Client-Version": "1.107.0",
|
||||
"x-request-source": "local", // MITM bypass
|
||||
},
|
||||
body: JSON.stringify({
|
||||
...(projectId ? { project: projectId } : {})
|
||||
}),
|
||||
}, 10000, proxyOptions);
|
||||
|
||||
if (response.status === 403) {
|
||||
return {
|
||||
message: "Antigravity quota API access forbidden. Chat may still work.",
|
||||
quotas: {}
|
||||
};
|
||||
}
|
||||
|
||||
if (response.status === 401) {
|
||||
return {
|
||||
message: "Antigravity quota API authentication expired. Chat may still work.",
|
||||
quotas: {}
|
||||
};
|
||||
}
|
||||
|
||||
if (!response.ok) {
|
||||
throw new Error(`Antigravity API error: ${response.status}`);
|
||||
}
|
||||
|
||||
const data = await response.json();
|
||||
const quotas = {};
|
||||
|
||||
// Parse model quotas (inspired by vscode-antigravity-cockpit)
|
||||
if (data.models) {
|
||||
// Filter only recommended/important models (must match PROVIDER_MODELS ag ids)
|
||||
const importantModels = [
|
||||
'gemini-3-flash-agent',
|
||||
'gemini-3.5-flash-low',
|
||||
'gemini-3.5-flash-extra-low',
|
||||
'gemini-pro-agent',
|
||||
'gemini-3.1-pro-low',
|
||||
'claude-sonnet-4-6',
|
||||
'claude-opus-4-6-thinking',
|
||||
'gpt-oss-120b-medium',
|
||||
'gemini-3-flash',
|
||||
];
|
||||
|
||||
for (const [modelKey, info] of Object.entries(data.models)) {
|
||||
// Skip models without quota info
|
||||
if (!info.quotaInfo) {
|
||||
continue;
|
||||
}
|
||||
|
||||
// Skip internal models and non-important models
|
||||
if (info.isInternal || !importantModels.includes(modelKey)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
const remainingFraction = info.quotaInfo.remainingFraction || 0;
|
||||
const remainingPercentage = remainingFraction * 100;
|
||||
|
||||
// Convert percentage to used/total for UI compatibility
|
||||
const total = 1000; // Normalized base
|
||||
const remaining = Math.round(total * remainingFraction);
|
||||
const used = total - remaining;
|
||||
|
||||
// Use modelKey as key (matches PROVIDER_MODELS id)
|
||||
quotas[modelKey] = {
|
||||
used,
|
||||
total,
|
||||
resetAt: parseResetTime(info.quotaInfo.resetTime),
|
||||
remainingPercentage,
|
||||
unlimited: false,
|
||||
displayName: info.displayName || modelKey,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
plan: subscriptionInfo?.currentTier?.name || "Unknown",
|
||||
quotas,
|
||||
subscriptionInfo,
|
||||
};
|
||||
} catch (error) {
|
||||
console.error("[Antigravity Usage] Error:", error.message, error.cause);
|
||||
return { message: `Antigravity error: ${error.message}` };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get Antigravity subscription info
|
||||
*/
|
||||
async function getAntigravitySubscriptionInfo(accessToken, proxyOptions = null) {
|
||||
try {
|
||||
const response = await fetchWithTimeout(ANTIGRAVITY_CONFIG.loadProjectApiUrl, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"User-Agent": ANTIGRAVITY_CONFIG.userAgent,
|
||||
"Content-Type": "application/json",
|
||||
"x-request-source": "local", // MITM bypass
|
||||
},
|
||||
body: JSON.stringify({ metadata: CLIENT_METADATA, mode: 1 }),
|
||||
}, 10000, proxyOptions);
|
||||
|
||||
if (!response.ok) return null;
|
||||
return await response.json();
|
||||
} catch (error) {
|
||||
console.error("[Antigravity Subscription] Error:", error.message);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
168
open-sse/services/usage/kiro.js
Normal file
168
open-sse/services/usage/kiro.js
Normal file
@@ -0,0 +1,168 @@
|
||||
/**
|
||||
* Kiro (AWS CodeWhisperer) usage handler
|
||||
*/
|
||||
|
||||
import { proxyAwareFetch } from "../../utils/proxyFetch.js";
|
||||
import { resolveDefaultProfileArn } from "../../config/kiroConstants.js";
|
||||
import { U, parseResetTime } from "./shared.js";
|
||||
|
||||
/**
|
||||
* Kiro (AWS CodeWhisperer) Usage
|
||||
*/
|
||||
function parseKiroQuotaData(data) {
|
||||
const usageList = data.usageBreakdownList || [];
|
||||
const quotaInfo = {};
|
||||
const resetAt = parseResetTime(data.nextDateReset || data.resetDate);
|
||||
|
||||
usageList.forEach((breakdown) => {
|
||||
const resourceType = breakdown.resourceType?.toLowerCase() || "unknown";
|
||||
const used = breakdown.currentUsageWithPrecision || 0;
|
||||
const total = breakdown.usageLimitWithPrecision || 0;
|
||||
|
||||
quotaInfo[resourceType] = {
|
||||
used,
|
||||
total,
|
||||
remaining: total - used,
|
||||
resetAt,
|
||||
unlimited: false,
|
||||
};
|
||||
|
||||
// Add free trial if available
|
||||
if (breakdown.freeTrialInfo) {
|
||||
const freeUsed = breakdown.freeTrialInfo.currentUsageWithPrecision || 0;
|
||||
const freeTotal = breakdown.freeTrialInfo.usageLimitWithPrecision || 0;
|
||||
|
||||
quotaInfo[`${resourceType}_freetrial`] = {
|
||||
used: freeUsed,
|
||||
total: freeTotal,
|
||||
remaining: freeTotal - freeUsed,
|
||||
resetAt: parseResetTime(breakdown.freeTrialInfo.freeTrialExpiry || resetAt),
|
||||
unlimited: false,
|
||||
};
|
||||
}
|
||||
});
|
||||
|
||||
return {
|
||||
plan: data.subscriptionInfo?.subscriptionTitle || "Kiro",
|
||||
quotas: quotaInfo,
|
||||
};
|
||||
}
|
||||
|
||||
export async function getKiroUsage(accessToken, providerSpecificData, proxyOptions = null) {
|
||||
const authMethod = providerSpecificData?.authMethod || "builder-id";
|
||||
const profileArn = providerSpecificData?.profileArn || resolveDefaultProfileArn(authMethod);
|
||||
|
||||
const getUsageParams = new URLSearchParams({
|
||||
isEmailRequired: "true",
|
||||
origin: "AI_EDITOR",
|
||||
resourceType: "AGENTIC_REQUEST",
|
||||
});
|
||||
|
||||
// For compatibility, try multiple known Kiro usage endpoints
|
||||
const attempts = [
|
||||
{
|
||||
name: "codewhisperer-get",
|
||||
run: async () => proxyAwareFetch(
|
||||
`${U("kiro").cwHost}${U("kiro").limitsPath}?${getUsageParams.toString()}`,
|
||||
{
|
||||
method: "GET",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"Accept": "application/json",
|
||||
"x-amz-user-agent": "aws-sdk-js/1.0.0 KiroIDE",
|
||||
"user-agent": "aws-sdk-js/1.0.0 KiroIDE",
|
||||
},
|
||||
},
|
||||
proxyOptions
|
||||
),
|
||||
},
|
||||
{
|
||||
name: "codewhisperer-post",
|
||||
run: async () => proxyAwareFetch(U("kiro").cwHost, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"Content-Type": "application/x-amz-json-1.0",
|
||||
"x-amz-target": "AmazonCodeWhispererService.GetUsageLimits",
|
||||
"Accept": "application/json",
|
||||
},
|
||||
body: JSON.stringify({
|
||||
origin: "AI_EDITOR",
|
||||
profileArn,
|
||||
resourceType: "AGENTIC_REQUEST",
|
||||
}),
|
||||
}, proxyOptions),
|
||||
},
|
||||
{
|
||||
name: "q-get",
|
||||
run: async () => {
|
||||
const params = new URLSearchParams({
|
||||
origin: "AI_EDITOR",
|
||||
profileArn,
|
||||
resourceType: "AGENTIC_REQUEST",
|
||||
});
|
||||
return proxyAwareFetch(`${U("kiro").qHost}${U("kiro").limitsPath}?${params}`, {
|
||||
method: "GET",
|
||||
headers: {
|
||||
"Authorization": `Bearer ${accessToken}`,
|
||||
"Accept": "application/json",
|
||||
},
|
||||
}, proxyOptions);
|
||||
},
|
||||
},
|
||||
];
|
||||
|
||||
let sawAuthError = false;
|
||||
const errors = [];
|
||||
|
||||
for (const attempt of attempts) {
|
||||
try {
|
||||
const response = await attempt.run();
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text().catch(() => "");
|
||||
if (response.status === 401 || response.status === 403) {
|
||||
sawAuthError = true;
|
||||
}
|
||||
errors.push(`${attempt.name}:${response.status}${errorText ? `:${errorText}` : ""}`);
|
||||
continue;
|
||||
}
|
||||
|
||||
const data = await response.json();
|
||||
return parseKiroQuotaData(data);
|
||||
} catch (error) {
|
||||
errors.push(`${attempt.name}:${error.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
if (sawAuthError && authMethod === "idc") {
|
||||
return {
|
||||
message: "Kiro quota API is unavailable for the current AWS IAM Identity Center session. Chat may still work. If this persists after renewing your session, reconnect Kiro.",
|
||||
quotas: {},
|
||||
};
|
||||
}
|
||||
|
||||
// Social auth (Google/GitHub) - these use a different token format that may not work with AWS CodeWhisperer quota APIs
|
||||
if (sawAuthError && (authMethod === "google" || authMethod === "github")) {
|
||||
return {
|
||||
message: "Kiro quota API authentication expired. Chat may still work.",
|
||||
quotas: {},
|
||||
};
|
||||
}
|
||||
|
||||
if (sawAuthError) {
|
||||
return {
|
||||
message: "Kiro quota API rejected the current token. Chat may still work.",
|
||||
quotas: {},
|
||||
};
|
||||
}
|
||||
|
||||
const fallbackMessage =
|
||||
errors.length > 0
|
||||
? `Unable to fetch Kiro usage right now. (${errors[errors.length - 1]})`
|
||||
: "Unable to fetch Kiro usage right now.";
|
||||
|
||||
return {
|
||||
message: fallbackMessage,
|
||||
quotas: {},
|
||||
};
|
||||
}
|
||||
234
open-sse/services/usage/minimax.js
Normal file
234
open-sse/services/usage/minimax.js
Normal file
@@ -0,0 +1,234 @@
|
||||
/**
|
||||
* MiniMax usage handler
|
||||
*/
|
||||
|
||||
import { proxyAwareFetch } from "../../utils/proxyFetch.js";
|
||||
import { U, parseResetTime } from "./shared.js";
|
||||
|
||||
// MiniMax usage endpoints (try in order, fallback on transient errors)
|
||||
const MINIMAX_USAGE_URLS = {
|
||||
minimax: U("minimax").urls,
|
||||
"minimax-cn": U("minimax-cn").urls,
|
||||
};
|
||||
|
||||
// ── MiniMax helpers ──────────────────────────────────────────────────────
|
||||
function getMiniMaxField(model, snakeKey, camelKey) {
|
||||
if (!model || typeof model !== "object") return null;
|
||||
return model[snakeKey] ?? model[camelKey] ?? null;
|
||||
}
|
||||
|
||||
function getMiniMaxModelName(model) {
|
||||
return String(getMiniMaxField(model, "model_name", "modelName") || "").trim();
|
||||
}
|
||||
|
||||
function formatMiniMaxQuotaName(model) {
|
||||
const rawName = getMiniMaxModelName(model);
|
||||
if (!rawName) return "MiniMax";
|
||||
|
||||
// M3+ shared quota pool: MiniMax reports M-series as a single wildcard
|
||||
// bucket ("MiniMax-M*"). Newer responses rename it to plain "general".
|
||||
// Render both as a friendly series label rather than leaking the
|
||||
// asterisk or the vague "general" word to the UI.
|
||||
if (rawName === "MiniMax-M*" || rawName === "general") return "M-series";
|
||||
|
||||
return rawName
|
||||
.replace(/[_-]+/g, " ")
|
||||
.replace(/\s+/g, " ")
|
||||
.trim()
|
||||
.replace(/\b\w/g, (ch) => ch.toUpperCase())
|
||||
.replace(/\bTo\b/g, "to")
|
||||
.replace(/\bTts\b/g, "TTS")
|
||||
.replace(/\bHd\b/g, "HD");
|
||||
}
|
||||
|
||||
function getMiniMaxProvidedPercent(model, snakeKey, camelKey) {
|
||||
if (!model || typeof model !== "object") return null;
|
||||
const raw = model[snakeKey] ?? model[camelKey];
|
||||
if (raw === null || raw === undefined) return null;
|
||||
const num = Number(raw);
|
||||
if (!Number.isFinite(num)) return null;
|
||||
return Math.max(0, Math.min(100, num));
|
||||
}
|
||||
|
||||
function getMiniMaxSessionTotal(model) {
|
||||
return Math.max(0, Number(getMiniMaxField(model, "current_interval_total_count", "currentIntervalTotalCount")) || 0);
|
||||
}
|
||||
|
||||
function getMiniMaxWeeklyTotal(model) {
|
||||
return Math.max(0, Number(getMiniMaxField(model, "current_weekly_total_count", "currentWeeklyTotalCount")) || 0);
|
||||
}
|
||||
|
||||
function hasMiniMaxQuota(model) {
|
||||
// Old format has real count totals; M3-era M-series buckets ship percent-only
|
||||
// (count fields are 0) so accept those too.
|
||||
if (getMiniMaxSessionTotal(model) > 0 || getMiniMaxWeeklyTotal(model) > 0) return true;
|
||||
if (getMiniMaxProvidedPercent(model, "current_interval_remaining_percent", "currentIntervalRemainingPercent") !== null) return true;
|
||||
if (getMiniMaxProvidedPercent(model, "current_weekly_remaining_percent", "currentWeeklyRemainingPercent") !== null) return true;
|
||||
return false;
|
||||
}
|
||||
|
||||
function getMiniMaxResetAt(model, capturedAtMs, remainsSnake, remainsCamel, endSnake, endCamel) {
|
||||
const remainsMs = Number(getMiniMaxField(model, remainsSnake, remainsCamel)) || 0;
|
||||
if (remainsMs > 0) return new Date(capturedAtMs + remainsMs).toISOString();
|
||||
return parseResetTime(getMiniMaxField(model, endSnake, endCamel));
|
||||
}
|
||||
|
||||
function buildMiniMaxQuota(total, count, resetAt, countMeansRemaining, providedPercent = null) {
|
||||
const safeTotal = Math.max(0, total);
|
||||
const used = countMeansRemaining ? Math.max(safeTotal - count, 0) : Math.min(Math.max(0, count), safeTotal);
|
||||
const remaining = Math.max(safeTotal - used, 0);
|
||||
// M-series buckets ship percent-only (count = 0). Prefer the upstream value
|
||||
// when present, otherwise fall back to the computed percentage. When the
|
||||
// quota is unbounded (no count) and no upstream percent is available, surface
|
||||
// the percent anyway as long as it is defined.
|
||||
const remainingPercentage = providedPercentage(providedPercent, remaining, safeTotal);
|
||||
return {
|
||||
used,
|
||||
total: safeTotal,
|
||||
remaining,
|
||||
remainingPercentage,
|
||||
resetAt,
|
||||
unlimited: false,
|
||||
};
|
||||
}
|
||||
|
||||
function providedPercentage(provided, remaining, total) {
|
||||
if (provided !== null && provided !== undefined && Number.isFinite(provided)) {
|
||||
return Math.max(0, Math.min(100, provided));
|
||||
}
|
||||
return total > 0 ? Math.max(0, Math.min(100, (remaining / total) * 100)) : 0;
|
||||
}
|
||||
|
||||
function addMiniMaxQuota(quotas, key, model, getTotal, countSnake, countCamel, percentSnake, percentCamel, resetArgs, countMeansRemaining) {
|
||||
const total = getTotal(model);
|
||||
const providedPercent = getMiniMaxProvidedPercent(model, percentSnake, percentCamel);
|
||||
if (total <= 0 && providedPercent === null) return;
|
||||
|
||||
const count = Math.max(0, Number(getMiniMaxField(model, countSnake, countCamel)) || 0);
|
||||
let effectiveTotal = total;
|
||||
let effectiveCount = count;
|
||||
if (total <= 0) {
|
||||
// M-series bucket: API only ships *_remaining_percent (count = 0). Normalize
|
||||
// to total=100. The downstream buildMiniMaxQuota treats the count as
|
||||
// "used" or "remaining" depending on countMeansRemaining, so the synthetic
|
||||
// count has to match that semantic — otherwise the UI flips the percentage.
|
||||
effectiveTotal = 100;
|
||||
const pct = providedPercent;
|
||||
effectiveCount = countMeansRemaining
|
||||
? Math.round(effectiveTotal * (pct / 100))
|
||||
: Math.round(effectiveTotal * (1 - pct / 100));
|
||||
}
|
||||
quotas[key] = buildMiniMaxQuota(
|
||||
effectiveTotal,
|
||||
effectiveCount,
|
||||
getMiniMaxResetAt(model, ...resetArgs),
|
||||
countMeansRemaining,
|
||||
providedPercent
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* MiniMax Token Plan / Coding Plan usage
|
||||
*/
|
||||
export async function getMiniMaxUsage(apiKey, provider, proxyOptions = null) {
|
||||
if (!apiKey) {
|
||||
return { message: "MiniMax API key not available." };
|
||||
}
|
||||
|
||||
const usageUrls = MINIMAX_USAGE_URLS[provider] || [];
|
||||
let lastErrorMessage = "";
|
||||
|
||||
for (let index = 0; index < usageUrls.length; index += 1) {
|
||||
const usageUrl = usageUrls[index];
|
||||
const canFallback = index < usageUrls.length - 1;
|
||||
|
||||
try {
|
||||
const response = await proxyAwareFetch(usageUrl, {
|
||||
method: "GET",
|
||||
headers: {
|
||||
Authorization: `Bearer ${apiKey}`,
|
||||
Accept: "application/json",
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
}, proxyOptions);
|
||||
|
||||
const rawText = await response.text();
|
||||
let payload = {};
|
||||
if (rawText) {
|
||||
try { payload = JSON.parse(rawText); } catch { payload = {}; }
|
||||
}
|
||||
|
||||
const baseResp = (payload?.base_resp ?? payload?.baseResp) || {};
|
||||
const apiStatusCode = Number(baseResp.status_code ?? baseResp.statusCode) || 0;
|
||||
const apiStatusMessage = String(baseResp.status_msg ?? baseResp.statusMsg ?? "").trim();
|
||||
const combined = `${apiStatusMessage} ${rawText}`.trim();
|
||||
const authLike = /token plan|coding plan|invalid api key|invalid key|unauthorized|inactive/i;
|
||||
|
||||
if (response.status === 401 || response.status === 403 || apiStatusCode === 1004 || authLike.test(combined)) {
|
||||
return { message: "MiniMax API key invalid or inactive. Use an active Token/Coding Plan key." };
|
||||
}
|
||||
|
||||
if (!response.ok) {
|
||||
lastErrorMessage = `MiniMax usage endpoint error (${response.status})`;
|
||||
if ((response.status === 404 || response.status === 405 || response.status >= 500) && canFallback) continue;
|
||||
return { message: `MiniMax connected. ${lastErrorMessage}` };
|
||||
}
|
||||
|
||||
if (apiStatusCode !== 0) {
|
||||
return { message: `MiniMax connected. ${apiStatusMessage || "Upstream quota API error"}` };
|
||||
}
|
||||
|
||||
const modelRemains = payload?.model_remains ?? payload?.modelRemains;
|
||||
const allModels = Array.isArray(modelRemains) ? modelRemains : [];
|
||||
const quotaModels = allModels.filter(hasMiniMaxQuota);
|
||||
|
||||
if (quotaModels.length === 0) {
|
||||
return { message: "MiniMax connected. No quota data was returned." };
|
||||
}
|
||||
|
||||
const capturedAtMs = Date.now();
|
||||
const countMeansRemaining = usageUrl.includes("/coding_plan/remains");
|
||||
const quotas = {};
|
||||
|
||||
for (const model of quotaModels) {
|
||||
const displayName = formatMiniMaxQuotaName(model);
|
||||
addMiniMaxQuota(
|
||||
quotas,
|
||||
`${displayName} (5h)`,
|
||||
model,
|
||||
getMiniMaxSessionTotal,
|
||||
"current_interval_usage_count",
|
||||
"currentIntervalUsageCount",
|
||||
"current_interval_remaining_percent",
|
||||
"currentIntervalRemainingPercent",
|
||||
[capturedAtMs, "remains_time", "remainsTime", "end_time", "endTime"],
|
||||
countMeansRemaining
|
||||
);
|
||||
|
||||
addMiniMaxQuota(
|
||||
quotas,
|
||||
`${displayName} (7d)`,
|
||||
model,
|
||||
getMiniMaxWeeklyTotal,
|
||||
"current_weekly_usage_count",
|
||||
"currentWeeklyUsageCount",
|
||||
"current_weekly_remaining_percent",
|
||||
"currentWeeklyRemainingPercent",
|
||||
[capturedAtMs, "weekly_remains_time", "weeklyRemainsTime", "weekly_end_time", "weeklyEndTime"],
|
||||
countMeansRemaining
|
||||
);
|
||||
}
|
||||
|
||||
if (Object.keys(quotas).length === 0) {
|
||||
return { message: "MiniMax connected. Unable to extract quota usage." };
|
||||
}
|
||||
|
||||
return { quotas };
|
||||
} catch (error) {
|
||||
lastErrorMessage = error.message;
|
||||
if (!canFallback) break;
|
||||
}
|
||||
}
|
||||
|
||||
return { message: lastErrorMessage ? `MiniMax connected. Unable to fetch usage: ${lastErrorMessage}` : "MiniMax connected. Unable to fetch usage." };
|
||||
}
|
||||
269
open-sse/services/usage/misc.js
Normal file
269
open-sse/services/usage/misc.js
Normal file
@@ -0,0 +1,269 @@
|
||||
/**
|
||||
* Misc usage handlers (Qwen, iFlow, Ollama, GLM, Vercel AI Gateway, Qoder)
|
||||
*/
|
||||
|
||||
import { proxyAwareFetch } from "../../utils/proxyFetch.js";
|
||||
import { U } from "./shared.js";
|
||||
|
||||
// GLM quota endpoints (region-aware) — url from registry transport.usage
|
||||
const GLM_QUOTA_URLS = {
|
||||
international: U("glm").url,
|
||||
china: U("glm-cn").url,
|
||||
};
|
||||
|
||||
// Vercel AI Gateway credits endpoint
|
||||
// Returns { balance: "95.50", total_used: "4.50" } (USD as decimal strings).
|
||||
const VERCEL_AI_GATEWAY_CREDITS_URL = U("vercel-ai-gateway").url;
|
||||
|
||||
/**
|
||||
* Qwen Usage
|
||||
*/
|
||||
export async function getQwenUsage(accessToken, providerSpecificData) {
|
||||
try {
|
||||
const resourceUrl = providerSpecificData?.resourceUrl;
|
||||
if (!resourceUrl) {
|
||||
return { message: "Qwen connected. No resource URL available." };
|
||||
}
|
||||
|
||||
// Qwen may have usage endpoint at resource URL
|
||||
return { message: "Qwen connected. Usage tracked per request." };
|
||||
} catch (error) {
|
||||
return { message: "Unable to fetch Qwen usage." };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* iFlow Usage
|
||||
*/
|
||||
export async function getIflowUsage(accessToken) {
|
||||
try {
|
||||
// iFlow may have usage endpoint
|
||||
return { message: "iFlow connected. Usage tracked per request." };
|
||||
} catch (error) {
|
||||
return { message: "Unable to fetch iFlow usage." };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Ollama Cloud Usage
|
||||
* Ollama Cloud uses an API key from ollama.com/settings/keys
|
||||
* and has no public usage API — free tier has light usage limits (resets every 5h & 7d).
|
||||
* This returns an informational message with the plan details.
|
||||
*/
|
||||
export async function getOllamaUsage(accessToken, providerSpecificData) {
|
||||
try {
|
||||
// Ollama Cloud does not expose a public quota/usage API.
|
||||
// The provider is configured as noAuth with a notice explaining limits.
|
||||
// We return a graceful message so the UI shows a friendly state instead of an error.
|
||||
const plan = providerSpecificData?.plan || "Free";
|
||||
return {
|
||||
plan,
|
||||
message: "Ollama Cloud uses a free tier with light usage limits (resets every 5h & 7d). For detailed usage tracking, visit ollama.com/settings/keys.",
|
||||
quotas: [],
|
||||
};
|
||||
} catch (error) {
|
||||
return { message: "Unable to fetch Ollama Cloud usage." };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* GLM Coding Plan usage (international + China regions)
|
||||
*/
|
||||
export async function getGlmUsage(apiKey, provider, proxyOptions = null) {
|
||||
if (!apiKey) {
|
||||
return { message: "GLM API key not available." };
|
||||
}
|
||||
|
||||
const region = provider === "glm-cn" ? "china" : "international";
|
||||
const quotaUrl = GLM_QUOTA_URLS[region];
|
||||
|
||||
try {
|
||||
const response = await proxyAwareFetch(quotaUrl, {
|
||||
headers: {
|
||||
Authorization: `Bearer ${apiKey}`,
|
||||
Accept: "application/json",
|
||||
},
|
||||
}, proxyOptions);
|
||||
|
||||
if (!response.ok) {
|
||||
if (response.status === 401) {
|
||||
return { message: "GLM API key invalid or expired." };
|
||||
}
|
||||
return { message: `GLM quota API error (${response.status}).` };
|
||||
}
|
||||
|
||||
const json = await response.json();
|
||||
const data = json?.data && typeof json.data === "object" ? json.data : {};
|
||||
const limits = Array.isArray(data.limits) ? data.limits : [];
|
||||
const quotas = {};
|
||||
|
||||
for (const limit of limits) {
|
||||
if (!limit || limit.type !== "TOKENS_LIMIT") continue;
|
||||
const usedPercent = Number(limit.percentage) || 0;
|
||||
const resetMs = Number(limit.nextResetTime) || 0;
|
||||
const remaining = Math.max(0, 100 - usedPercent);
|
||||
|
||||
quotas["session"] = {
|
||||
used: usedPercent,
|
||||
total: 100,
|
||||
remaining,
|
||||
remainingPercentage: remaining,
|
||||
resetAt: resetMs > 0 ? new Date(resetMs).toISOString() : null,
|
||||
unlimited: false,
|
||||
};
|
||||
}
|
||||
|
||||
const levelRaw = typeof data.level === "string" ? data.level : "";
|
||||
const plan = levelRaw
|
||||
? levelRaw.charAt(0).toUpperCase() + levelRaw.slice(1).toLowerCase()
|
||||
: "Unknown";
|
||||
|
||||
return { plan, quotas };
|
||||
} catch (error) {
|
||||
return { message: `GLM error: ${error.message}` };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Vercel AI Gateway usage — credit balance for the API key
|
||||
*
|
||||
* Calls GET /v1/credits which returns:
|
||||
* { "balance": "95.50", "total_used": "4.50" } (USD as decimal strings)
|
||||
*
|
||||
* We surface this as a single "Balance ($)" quota row so the existing
|
||||
* QuotaTable / progress-bar UI can render it. used = total_used,
|
||||
* total = balance + total_used (the original credit allotment), so the
|
||||
* remaining percentage equals balance / total.
|
||||
*
|
||||
* Docs: https://vercel.com/docs/ai-gateway/usage
|
||||
*/
|
||||
export async function getVercelAiGatewayUsage(apiKey, proxyOptions = null) {
|
||||
if (!apiKey) {
|
||||
return { message: "Vercel AI Gateway API key not available." };
|
||||
}
|
||||
|
||||
try {
|
||||
const response = await proxyAwareFetch(VERCEL_AI_GATEWAY_CREDITS_URL, {
|
||||
method: "GET",
|
||||
headers: {
|
||||
Authorization: `Bearer ${apiKey}`,
|
||||
Accept: "application/json",
|
||||
},
|
||||
}, proxyOptions);
|
||||
|
||||
if (response.status === 401 || response.status === 403) {
|
||||
return { message: "Vercel AI Gateway API key invalid or expired." };
|
||||
}
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text().catch(() => "");
|
||||
const trimmed = errorText ? `: ${errorText.slice(0, 200)}` : "";
|
||||
return { message: `Vercel AI Gateway credits API error (${response.status})${trimmed}` };
|
||||
}
|
||||
|
||||
const data = await response.json();
|
||||
|
||||
// Vercel returns numeric strings; coerce safely.
|
||||
const balance = Number(data?.balance) || 0;
|
||||
const totalUsed = Number(data?.total_used) || 0;
|
||||
|
||||
// Vercel gives $5/month free credit. The API doesn't return the
|
||||
// monthly allocation so we use the known constant as the denominator.
|
||||
const MONTHLY_CREDIT = 5;
|
||||
const remainingPercentage = (balance / MONTHLY_CREDIT) * 100;
|
||||
|
||||
if (balance <= 0 && totalUsed <= 0) {
|
||||
return {
|
||||
plan: "Pay-as-you-go",
|
||||
message: "Vercel AI Gateway connected. No credit allocation found (BYOK or unfunded account).",
|
||||
quotas: {},
|
||||
};
|
||||
}
|
||||
|
||||
// "Used (USD)": how much has been spent this month (no fixed cap → unlimited).
|
||||
// "Remaining (USD)": balance remaining out of the $5 monthly allocation.
|
||||
return {
|
||||
plan: "Pay-as-you-go",
|
||||
quotas: {
|
||||
"Used (USD)": {
|
||||
used: totalUsed,
|
||||
total: 0,
|
||||
remaining: 0,
|
||||
remainingPercentage: 100,
|
||||
unlimited: true,
|
||||
},
|
||||
"Remaining (USD)": {
|
||||
used: balance,
|
||||
total: MONTHLY_CREDIT,
|
||||
remaining: balance,
|
||||
remainingPercentage,
|
||||
unlimited: false,
|
||||
},
|
||||
},
|
||||
};
|
||||
} catch (error) {
|
||||
return { message: `Vercel AI Gateway error: ${error.message}` };
|
||||
}
|
||||
}
|
||||
|
||||
export async function getQoderUsage(accessToken, proxyOptions = null) {
|
||||
if (!accessToken) {
|
||||
return { message: "Qoder usage unavailable: no access token" };
|
||||
}
|
||||
try {
|
||||
const response = await proxyAwareFetch(
|
||||
U("qoder").url,
|
||||
{
|
||||
method: "GET",
|
||||
headers: {
|
||||
Authorization: `Bearer ${accessToken}`,
|
||||
Accept: "application/json",
|
||||
},
|
||||
},
|
||||
proxyOptions,
|
||||
);
|
||||
if (!response.ok) {
|
||||
return { message: `Qoder connected. Usage fetch returned ${response.status}.` };
|
||||
}
|
||||
const body = await response.json().catch(() => null);
|
||||
if (!body) {
|
||||
return { message: "Qoder connected. Usage response was not JSON." };
|
||||
}
|
||||
// Quota records live under `quotas`; scalar metadata
|
||||
// (totalUsagePercentage, isQuotaExceeded, expiresAt) are surfaced as
|
||||
// siblings so the dashboard parser doesn't try to render them as rows.
|
||||
const userQuota = body.userQuota || {};
|
||||
const orgQuota = body.orgResourcePackage || {};
|
||||
// Qoder publishes a single absolute reset timestamp (`expiresAt` in ms);
|
||||
// surface it on every quota record as ISO so the table can render
|
||||
// "resets at" alongside used/total.
|
||||
const expiresAtMs = Number.isFinite(Number(body.expiresAt)) && Number(body.expiresAt) > 0
|
||||
? Number(body.expiresAt)
|
||||
: null;
|
||||
const resetAt = expiresAtMs ? new Date(expiresAtMs).toISOString() : null;
|
||||
const quotas = {
|
||||
user: {
|
||||
total: Number(userQuota.total) || 0,
|
||||
used: Number(userQuota.used) || 0,
|
||||
remaining: Number(userQuota.remaining) || 0,
|
||||
unit: userQuota.unit || "credits",
|
||||
resetAt,
|
||||
},
|
||||
organization: {
|
||||
total: Number(orgQuota.total) || 0,
|
||||
used: Number(orgQuota.used) || 0,
|
||||
remaining: Number(orgQuota.remaining) || 0,
|
||||
unit: orgQuota.unit || "credits",
|
||||
resetAt,
|
||||
},
|
||||
};
|
||||
return {
|
||||
quotas,
|
||||
totalUsagePercentage: Number(body.totalUsagePercentage) || 0,
|
||||
isQuotaExceeded: !!body.isQuotaExceeded,
|
||||
expiresAt: expiresAtMs,
|
||||
};
|
||||
} catch (error) {
|
||||
return { message: `Qoder connected. Unable to fetch usage: ${error.message}` };
|
||||
}
|
||||
}
|
||||
70
open-sse/services/usage/shared.js
Normal file
70
open-sse/services/usage/shared.js
Normal file
@@ -0,0 +1,70 @@
|
||||
/**
|
||||
* Shared usage helpers (cross-provider)
|
||||
*/
|
||||
|
||||
import { PROVIDERS } from "../../providers/index.js";
|
||||
import { proxyAwareFetch } from "../../utils/proxyFetch.js";
|
||||
|
||||
// usage endpoints: single source from registry transport.usage
|
||||
export const U = (id) => PROVIDERS[id]?.usage || {};
|
||||
|
||||
/**
|
||||
* Parse reset date/time to ISO string
|
||||
* Handles multiple formats: Unix timestamp (ms), ISO date string, etc.
|
||||
*/
|
||||
export function parseResetTime(resetValue) {
|
||||
if (!resetValue) return null;
|
||||
|
||||
try {
|
||||
// If it's already a Date object
|
||||
if (resetValue instanceof Date) {
|
||||
return resetValue.toISOString();
|
||||
}
|
||||
|
||||
// Unix timestamps from provider APIs may be seconds or milliseconds.
|
||||
if (typeof resetValue === 'number') {
|
||||
return new Date(resetValue < 1e12 ? resetValue * 1000 : resetValue).toISOString();
|
||||
}
|
||||
|
||||
// If it's a numeric string, treat it like a Unix timestamp too.
|
||||
if (typeof resetValue === 'string') {
|
||||
if (/^\d+$/.test(resetValue)) {
|
||||
const timestamp = Number(resetValue);
|
||||
return new Date(timestamp < 1e12 ? timestamp * 1000 : timestamp).toISOString();
|
||||
}
|
||||
return new Date(resetValue).toISOString();
|
||||
}
|
||||
|
||||
return null;
|
||||
} catch (error) {
|
||||
console.warn(`Failed to parse reset time: ${resetValue}`, error);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
export function toFiniteNumber(value, fallback = 0) {
|
||||
if (typeof value === "number" && Number.isFinite(value)) return value;
|
||||
if (typeof value === "string" && value.trim()) {
|
||||
const parsed = Number(value);
|
||||
if (Number.isFinite(parsed)) return parsed;
|
||||
}
|
||||
return fallback;
|
||||
}
|
||||
|
||||
export function normalizeCloudCodeProjectId(project) {
|
||||
if (typeof project === "string") return project.trim() || null;
|
||||
if (project && typeof project === "object" && typeof project.id === "string") {
|
||||
return project.id.trim() || null;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
export async function fetchWithTimeout(url, opts, ms = 10000, proxyOptions = null) {
|
||||
const controller = new AbortController();
|
||||
const timeoutId = setTimeout(() => controller.abort(), ms);
|
||||
try {
|
||||
return await proxyAwareFetch(url, { ...opts, signal: controller.signal }, proxyOptions);
|
||||
} finally {
|
||||
clearTimeout(timeoutId);
|
||||
}
|
||||
}
|
||||
3
open-sse/utils/sse.js
Normal file
3
open-sse/utils/sse.js
Normal file
@@ -0,0 +1,3 @@
|
||||
export function sseChunk(data) {
|
||||
return `data: ${JSON.stringify(data)}\n\n`;
|
||||
}
|
||||
Reference in New Issue
Block a user