merge: resolve conflicts with origin/master - keep both local and remote features

This commit is contained in:
2026-06-22 11:40:16 +07:00
424 changed files with 24520 additions and 11731 deletions

View File

@@ -3,9 +3,9 @@ import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js";
import { OAUTH_ENDPOINTS, ANTIGRAVITY_HEADERS, INTERNAL_REQUEST_HEADER, AG_DEFAULT_TOOLS, AG_TOOL_SUFFIX } from "../config/appConstants.js";
import { HTTP_STATUS } from "../config/runtimeConfig.js";
import { deriveSessionId } from "../utils/sessionManager.js";
import { resolveSessionId } from "../utils/sessionManager.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js";
import { cleanJSONSchemaForAntigravity } from "../translator/helpers/geminiHelper.js";
import { cleanJSONSchemaForAntigravity } from "../translator/formats/gemini.js";
// Sanitize function name: Gemini requires [a-zA-Z_][a-zA-Z0-9_.:\-]{0,63}
function sanitizeFunctionName(name) {
@@ -18,6 +18,56 @@ function sanitizeFunctionName(name) {
const MAX_RETRY_AFTER_MS = 10000;
const MAX_ANTIGRAVITY_OUTPUT_TOKENS = 16384;
// Fields Google generateContent rejects (Claude/OpenAI/Qwen thinking fields set at body root by thinkingUnified.js)
const ANTIGRAVITY_REQUEST_BLACKLIST = [
"output_config",
"thinking",
"reasoning_effort",
"reasoning",
"enable_thinking",
"thinking_budget",
"thinkingConfig",
];
// Strip blacklisted fields from an object (used for both body.request and top-level body)
const stripBlacklisted = obj => {
for (const key of ANTIGRAVITY_REQUEST_BLACKLIST) delete obj[key];
};
// Image generation model name patterns
const IMAGE_MODEL_PATTERNS = [
/image/i,
/imagen/i,
/image-generation/i,
];
// Detect if a model is an image generation model
function isImageModel(model) {
if (!model) return false;
return IMAGE_MODEL_PATTERNS.some(p => p.test(model));
}
// Parse aspect ratio / resolution from model name suffixes
// e.g. "gemini-3.1-flash-image-16x9" -> { aspectRatio: "16:9" }
// e.g. "gemini-3.1-flash-image-1024x768" -> { aspectRatio: "4:3" }
function parseImageConfig(model) {
const config = { aspectRatio: "1:1" };
const resMatch = model.match(/(\d+)x(\d+)$/);
if (resMatch) {
const w = parseInt(resMatch[1]);
const h = parseInt(resMatch[2]);
if (w <= 16 && h <= 16) {
config.aspectRatio = `${w}:${h}`;
} else {
// Resolution like 1024x768 — derive aspect ratio
const gcd = (a, b) => b ? gcd(b, a % b) : a;
const d = gcd(w, h);
config.aspectRatio = `${w/d}:${h/d}`;
}
}
return config;
}
export class AntigravityExecutor extends BaseExecutor {
constructor() {
super("antigravity", PROVIDERS.antigravity);
@@ -26,17 +76,22 @@ export class AntigravityExecutor extends BaseExecutor {
buildUrl(model, stream, urlIndex = 0) {
const baseUrls = this.getBaseUrls();
const baseUrl = baseUrls[urlIndex] || baseUrls[0];
const action = stream ? "streamGenerateContent?alt=sse" : "generateContent";
// Image generation MUST use non-streaming generateContent
const forceNonStream = isImageModel(model);
const action = (stream && !forceNonStream) ? "streamGenerateContent?alt=sse" : "generateContent";
return `${baseUrl}/v1internal:${action}`;
}
// sessionId comes from transformRequest output; base.execute runs transformRequest before
// buildHeaders, so we read it from instance state cached there (fallback: explicit arg).
buildHeaders(credentials, stream = true, sessionId = null) {
const sid = sessionId || this._lastSessionId;
return {
"Content-Type": "application/json",
"Authorization": `Bearer ${credentials.accessToken}`,
"User-Agent": this.config.headers?.["User-Agent"] || ANTIGRAVITY_HEADERS["User-Agent"],
[INTERNAL_REQUEST_HEADER.name]: INTERNAL_REQUEST_HEADER.value,
...(sessionId && { "X-Machine-Session-Id": sessionId }),
...(sid && { "X-Machine-Session-Id": sid }),
"Accept": stream ? "text/event-stream" : "application/json"
};
}
@@ -44,6 +99,53 @@ export class AntigravityExecutor extends BaseExecutor {
transformRequest(model, body, stream, credentials) {
const projectId = credentials?.projectId || this.generateProjectId();
// ─── Image generation: completely different request structure ───
if (isImageModel(model)) {
const imageConfig = parseImageConfig(model);
// Strip model name suffixes for the actual API model name
const cleanModel = model.replace(/-(\d+)x(\d+)$/, "");
// Build simplified contents — text-only, merge all user messages
const contents = [];
const srcContents = body.request?.contents || body.contents || [];
for (const c of srcContents) {
const textParts = (c.parts || []).filter(p => p.text !== undefined).map(p => ({ text: p.text }));
if (textParts.length > 0) {
contents.push({ role: c.role || "user", parts: textParts });
}
}
const sessionId = resolveSessionId({
headers: credentials?.rawHeaders,
body,
connectionId: credentials?.email || credentials?.connectionId,
scope: "antigravity",
});
this._lastSessionId = sessionId;
return {
project: projectId,
model: cleanModel,
userAgent: "antigravity",
requestType: "image_gen",
requestId: `agent-${crypto.randomUUID()}`,
request: {
contents,
generationConfig: {
temperature: 1.0,
topP: 0.95,
topK: 40,
maxOutputTokens: 8192,
imageConfig,
},
sessionId,
// No tools, no systemInstruction, no safetySettings for image gen
},
};
}
// ─── Standard (non-image) request ───
// Fix contents for Claude models via Antigravity
const contents = body.request?.contents?.map(c => {
let role = c.role;
@@ -80,7 +182,9 @@ export class AntigravityExecutor extends BaseExecutor {
tools = allDeclarations.length > 0 ? [{ functionDeclarations: allDeclarations }] : [];
}
// Strip tools/toolConfig (handled separately) and blacklisted fields that Google rejects
const { tools: _originalTools, toolConfig: _originalToolConfig, ...requestWithoutTools } = body.request || {};
stripBlacklisted(requestWithoutTools);
const generationConfig = { ...(requestWithoutTools.generationConfig || {}) };
if (generationConfig.maxOutputTokens > MAX_ANTIGRAVITY_OUTPUT_TOKENS) {
generationConfig.maxOutputTokens = MAX_ANTIGRAVITY_OUTPUT_TOKENS;
@@ -91,11 +195,16 @@ export class AntigravityExecutor extends BaseExecutor {
generationConfig,
...(contents && { contents }),
...(tools && { tools }),
sessionId: body.request?.sessionId || deriveSessionId(credentials?.email || credentials?.connectionId),
sessionId: body.request?.sessionId || resolveSessionId({ headers: credentials?.rawHeaders, body, connectionId: credentials?.email || credentials?.connectionId, scope: "antigravity" }),
safetySettings: undefined,
...(tools?.length > 0 && { toolConfig: { functionCallingConfig: { mode: "VALIDATED" } } })
};
// Strip blacklisted thinking fields from top-level body (set by thinkingUnified.js at root, not body.request)
stripBlacklisted(body);
this._lastSessionId = transformedRequest.sessionId; // cached for buildHeaders (base.execute order)
return {
...body,
project: projectId,
@@ -196,98 +305,23 @@ export class AntigravityExecutor extends BaseExecutor {
return totalMs > 0 ? totalMs : null;
}
async execute({ model, body, stream, credentials, signal, log, proxyOptions = null }) {
const fallbackCount = this.getFallbackCount();
let lastError = null;
let lastStatus = 0;
const MAX_AUTO_RETRIES = 3;
const MAX_RETRY_AFTER_RETRIES = 3;
const retryAttemptsByUrl = {}; // Track retry attempts per URL
const retryAfterAttemptsByUrl = {}; // Track Retry-After retries per URL
for (let urlIndex = 0; urlIndex < fallbackCount; urlIndex++) {
const url = this.buildUrl(model, stream, urlIndex);
const transformedBody = this.transformRequest(model, body, stream, credentials);
const sessionId = transformedBody.request?.sessionId;
const headers = this.buildHeaders(credentials, stream, sessionId);
// Initialize retry counters for this URL
if (!retryAttemptsByUrl[urlIndex]) {
retryAttemptsByUrl[urlIndex] = 0;
}
if (!retryAfterAttemptsByUrl[urlIndex]) {
retryAfterAttemptsByUrl[urlIndex] = 0;
}
// Hook called by BaseExecutor.tryRetry: derive delay from Retry-After (header → body),
// cap at MAX_RETRY_AFTER_MS, else exponential backoff for 429. Return false to veto (fallback URL).
async computeRetryDelay(response, attempt) {
let retryMs = this.parseRetryHeaders(response.headers);
if (!retryMs) {
try {
const response = await proxyAwareFetch(url, {
method: "POST",
headers,
body: JSON.stringify(transformedBody),
signal
}, proxyOptions);
if (response.status === HTTP_STATUS.RATE_LIMITED || response.status === HTTP_STATUS.SERVICE_UNAVAILABLE) {
// Try to get retry time from headers first
let retryMs = this.parseRetryHeaders(response.headers);
// If no retry time in headers, try to parse from error message body
if (!retryMs) {
try {
const errorBody = await response.clone().text();
const errorJson = JSON.parse(errorBody);
const errorMessage = errorJson?.error?.message || errorJson?.message || "";
retryMs = this.parseRetryFromErrorMessage(errorMessage);
} catch (e) {
// Ignore parse errors, will fall back to exponential backoff
}
}
if (retryMs && retryMs <= MAX_RETRY_AFTER_MS && retryAfterAttemptsByUrl[urlIndex] < MAX_RETRY_AFTER_RETRIES) {
retryAfterAttemptsByUrl[urlIndex]++;
log?.debug?.("RETRY", `${response.status} with Retry-After: ${Math.ceil(retryMs / 1000)}s, waiting... (${retryAfterAttemptsByUrl[urlIndex]}/${MAX_RETRY_AFTER_RETRIES})`);
await new Promise(resolve => setTimeout(resolve, retryMs));
urlIndex--;
continue;
}
// Auto retry only for 429 when retryMs is 0 or undefined
if (response.status === HTTP_STATUS.RATE_LIMITED && (!retryMs || retryMs === 0) && retryAttemptsByUrl[urlIndex] < MAX_AUTO_RETRIES) {
retryAttemptsByUrl[urlIndex]++;
// Exponential backoff: 2s, 4s, 8s...
const backoffMs = Math.min(1000 * (2 ** retryAttemptsByUrl[urlIndex]), MAX_RETRY_AFTER_MS);
log?.debug?.("RETRY", `429 auto retry ${retryAttemptsByUrl[urlIndex]}/${MAX_AUTO_RETRIES} after ${backoffMs / 1000}s`);
await new Promise(resolve => setTimeout(resolve, backoffMs));
urlIndex--;
continue;
}
log?.debug?.("RETRY", `${response.status}, Retry-After ${retryMs ? `too long (${Math.ceil(retryMs / 1000)}s)` : 'missing'}, trying fallback`);
lastStatus = response.status;
if (urlIndex + 1 < fallbackCount) {
continue;
}
}
if (this.shouldRetry(response.status, urlIndex)) {
log?.debug?.("RETRY", `${response.status} on ${url}, trying fallback ${urlIndex + 1}`);
lastStatus = response.status;
continue;
}
return { response, url, headers, transformedBody };
} catch (error) {
lastError = error;
if (urlIndex + 1 < fallbackCount) {
log?.debug?.("RETRY", `Error on ${url}, trying fallback ${urlIndex + 1}`);
continue;
}
throw error;
const errorJson = JSON.parse(await response.clone().text());
retryMs = this.parseRetryFromErrorMessage(errorJson?.error?.message || errorJson?.message || "");
} catch {
// ignore parse errors → fall through to backoff
}
}
throw lastError || new Error(`All ${fallbackCount} URLs failed with status ${lastStatus}`);
if (retryMs) return retryMs <= MAX_RETRY_AFTER_MS ? retryMs : false;
if (response.status === HTTP_STATUS.RATE_LIMITED) {
return Math.min(1000 * (2 ** attempt), MAX_RETRY_AFTER_MS); // exponential backoff
}
return false;
}
/**

View File

@@ -2,6 +2,7 @@ import { HTTP_STATUS, RETRY_CONFIG, DEFAULT_RETRY_CONFIG, resolveRetryEntry, FET
import { shouldRefreshCredentials } from "../services/oauthCredentialManager.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js";
import { dbg } from "../utils/debugLog.js";
import { ANTHROPIC_API_VERSION, OPENAI_COMPAT_BASE, ANTHROPIC_COMPAT_BASE } from "../providers/shared.js";
/**
* BaseExecutor - Base class for provider executors
@@ -27,13 +28,13 @@ export class BaseExecutor {
buildUrl(model, stream, urlIndex = 0, credentials = null) {
if (this.provider?.startsWith?.("openai-compatible-")) {
const baseUrl = credentials?.providerSpecificData?.baseUrl || "https://api.openai.com/v1";
const baseUrl = credentials?.providerSpecificData?.baseUrl || OPENAI_COMPAT_BASE;
const normalized = baseUrl.replace(/\/$/, "");
const path = this.provider.includes("responses") ? "/responses" : "/chat/completions";
return `${normalized}${path}`;
}
if (this.provider?.startsWith?.("anthropic-compatible-")) {
const baseUrl = credentials?.providerSpecificData?.baseUrl || "https://api.anthropic.com/v1";
const baseUrl = credentials?.providerSpecificData?.baseUrl || ANTHROPIC_COMPAT_BASE;
const normalized = baseUrl.replace(/\/$/, "");
return `${normalized}/messages`;
}
@@ -55,7 +56,7 @@ export class BaseExecutor {
headers["Authorization"] = `Bearer ${credentials.accessToken}`;
}
if (!headers["anthropic-version"]) {
headers["anthropic-version"] = "2023-06-01";
headers["anthropic-version"] = ANTHROPIC_API_VERSION;
}
} else {
// Standard Bearer token auth for other providers
@@ -105,12 +106,20 @@ export class BaseExecutor {
const retryConfig = { ...DEFAULT_RETRY_CONFIG, ...this.config.retry };
// Schedule retry via retryConfig[statusKey]. Returns true when caller should `urlIndex--; continue`
const tryRetry = async (urlIndex, statusKey, reason) => {
// response (optional) lets a subclass hook compute a dynamic delay (e.g. antigravity Retry-After).
const tryRetry = async (urlIndex, statusKey, reason, response = null) => {
const { attempts, delayMs } = resolveRetryEntry(retryConfig[statusKey]);
if (attempts <= 0 || retryAttemptsByUrl[urlIndex] >= attempts) return false;
// Hook: subclass may derive delay from the response (headers/body). null → skip retry, use fallback.
let waitMs = delayMs;
if (response && this.computeRetryDelay) {
const dynamic = await this.computeRetryDelay(response, retryAttemptsByUrl[urlIndex] + 1, delayMs);
if (dynamic === false) return false; // hook vetoes retry (e.g. Retry-After too long)
if (dynamic != null) waitMs = dynamic;
}
retryAttemptsByUrl[urlIndex]++;
log?.debug?.("RETRY", `${reason} retry ${retryAttemptsByUrl[urlIndex]}/${attempts} after ${delayMs / 1000}s`);
await new Promise(resolve => setTimeout(resolve, delayMs));
log?.debug?.("RETRY", `${reason} retry ${retryAttemptsByUrl[urlIndex]}/${attempts} after ${waitMs / 1000}s`);
await new Promise(resolve => setTimeout(resolve, waitMs));
return true;
};
@@ -154,7 +163,7 @@ export class BaseExecutor {
const cl = response.headers?.get?.("content-length") || "?";
dbg("FETCH", `${this.provider.toUpperCase()} ← ${response.status} | ttft=${Date.now() - fetchT0}ms | ct=${ct} | cl=${cl}`);
if (await tryRetry(urlIndex, response.status, `status ${response.status}`)) { urlIndex--; continue; }
if (await tryRetry(urlIndex, response.status, `status ${response.status}`, response)) { urlIndex--; continue; }
if (this.shouldRetry(response.status, urlIndex)) {
log?.debug?.("RETRY", `${response.status} on ${url}, trying fallback ${urlIndex + 1}`);

View File

@@ -0,0 +1,36 @@
import { DefaultExecutor } from "./default.js";
/**
* CodeBuddyExecutor — talks to https://copilot.tencent.com/v2/chat/completions
*
* CodeBuddy is OpenAI-compatible but rejects non-stream chat requests
* (HTTP 400, code 11101 "Non-stream chat request is currently not supported").
* The same-format (openai→openai) translator path leaves body.stream as the
* client sent it, so we force it true here — 9router still re-aggregates the
* SSE into a JSON response for non-streaming clients.
*/
export class CodeBuddyExecutor extends DefaultExecutor {
constructor() {
super("codebuddy-cn");
}
transformRequest(model, body, stream, credentials) {
const transformed = super.transformRequest(model, body, stream, credentials);
transformed.stream = true;
// CodeBuddy only surfaces model reasoning when the request carries the CLI's
// OpenAI-style params: reasoning_effort + reasoning_summary:"auto". 9router's
// thinking pipeline sets reasoning_effort only when the client asks, and never
// sets reasoning_summary — so reasoning never shows. Mirror the CLI here.
const eff = transformed.reasoning_effort;
if (eff === "none" || eff === "off") {
delete transformed.reasoning_effort; // gateway has no "none" — just omit
} else {
if (!eff) transformed.reasoning_effort = "medium";
transformed.reasoning_summary = "auto";
}
return transformed;
}
}
export default CodeBuddyExecutor;

View File

@@ -1,4 +1,3 @@
import { createHash } from "crypto";
import { BaseExecutor } from "./base.js";
import { CODEX_DEFAULT_INSTRUCTIONS } from "../config/codexInstructions.js";
import { PROVIDERS } from "../config/providers.js";
@@ -6,30 +5,30 @@ import {
refreshProviderCredentials,
shouldRefreshCredentials,
} from "../services/oauthCredentialManager.js";
import { normalizeResponsesInput } from "../translator/helpers/responsesApiHelper.js";
import { fetchImageAsBase64 } from "../translator/helpers/imageHelper.js";
import { normalizeResponsesInput } from "../translator/formats/responsesApi.js";
import { fetchImageAsBase64 } from "../translator/concerns/image.js";
import { getModelUpstreamId } from "../config/providerModels.js";
import { getConsistentMachineId } from "../../src/shared/utils/machineId.js";
import { DEFAULT_RETRY_CONFIG, resolveRetryEntry } from "../config/runtimeConfig.js";
import { dbg } from "../utils/debugLog.js";
import { resolveSessionId } from "../utils/sessionManager.js";
// SSE error patterns inside 200-OK body that should trigger retry as if 503
const CODEX_SSE_OVERLOADED_PATTERNS = ["server_is_overloaded", "service_unavailable_error"];
const CODEX_SSE_PEEK_BYTES = 4096;
// In-memory map: hash(machineId + first assistant content) → { sessionId, lastUsed }
const SESSION_TTL_MS = 60 * 60 * 1000; // 1 hour
const assistantSessionMap = new Map();
// Server-generated item id prefixes that Codex /responses cannot resolve when store=false
const SERVER_ID_PATTERN = /^(rs|fc|resp|msg)_/;
// Hosted tool types that Codex/OpenAI Responses executes server-side
const CODEX_HOSTED_TOOL_TYPES = new Set([
"image_generation", "web_search", "web_search_preview", "file_search",
"computer", "computer_use_preview", "code_interpreter", "mcp", "local_shell"
"computer", "computer_use_preview", "code_interpreter", "mcp", "local_shell",
"tool_search"
]);
// Responses-native freeform tools carry a name plus format payload and must pass through intact.
const CODEX_PASSTHROUGH_TOOL_TYPES = new Set(["custom"]);
// Allowlist of fields accepted by Codex Responses API — anything else is stripped
const RESPONSES_API_ALLOWLIST = new Set([
"model", "input", "instructions", "tools", "tool_choice", "stream", "store",
@@ -76,6 +75,7 @@ function normalizeCodexTools(body) {
return true;
}
if (type !== "function") {
if (CODEX_PASSTHROUGH_TOOL_TYPES.has(type)) return true;
if (!type || tool.function || typeof tool.name === "string") return false;
return CODEX_HOSTED_TOOL_TYPES.has(type);
}
@@ -104,86 +104,17 @@ function normalizeCodexTools(body) {
}
}
// Cache machine ID at module level (resolved once)
let cachedMachineId = null;
getConsistentMachineId().then(id => { cachedMachineId = id; });
function hashContent(text) {
return createHash("sha256").update(text).digest("hex").slice(0, 16);
// Resolve prompt-cache session id: client session → assistant-text-hash → workspaceId → connection
function resolveCacheSessionId(body, credentials) {
return resolveSessionId({
headers: credentials?.rawHeaders,
body,
connectionId: credentials?.connectionId,
workspaceId: credentials?.providerSpecificData?.workspaceId,
scope: "codex"
});
}
function generateSessionId() {
return `sess_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 9)}`;
}
// Extract text content from an input item
function extractItemText(item) {
if (!item) return "";
if (typeof item.content === "string") return item.content;
if (Array.isArray(item.content)) {
return item.content.map(c => c.text || c.output || "").filter(Boolean).join("");
}
return "";
}
// Normalize a session id candidate (trim, length cap)
function normalizeSessionId(value) {
if (typeof value !== "string") return null;
const v = value.trim();
if (!v || v.length > 256) return null;
return v;
}
// Resolve prompt-cache session id with priority: body → assistant-text-hash → workspaceId → machineId
function resolveCacheSessionId(body, credentials, machineId) {
// 1. Client-provided session/conversation id (highest priority — stable per conversation)
const fromBody =
normalizeSessionId(body?.prompt_cache_key) ||
normalizeSessionId(body?.session_id) ||
normalizeSessionId(body?.conversation_id);
if (fromBody) return fromBody;
// 2. Hash accumulated assistant text (≥50 chars) — sticky session across turns
if (Array.isArray(body?.input) && body.input.length > 0) {
let text = "";
const MIN_LEN = 50;
const CAP_LEN = 200;
for (const item of body.input) {
if (item?.role !== "assistant") continue;
const t = extractItemText(item);
if (!t) continue;
text += t;
if (text.length >= CAP_LEN) break;
}
if (text.length >= MIN_LEN) {
const hash = hashContent((machineId || "") + text.slice(0, CAP_LEN));
const entry = assistantSessionMap.get(hash);
if (entry) {
entry.lastUsed = Date.now();
return entry.sessionId;
}
const sessionId = generateSessionId();
assistantSessionMap.set(hash, { sessionId, lastUsed: Date.now() });
return sessionId;
}
}
// 3. Account-wide fallback (workspaceId from connection)
const workspaceId = normalizeSessionId(credentials?.providerSpecificData?.workspaceId);
if (workspaceId) return workspaceId;
// 4. Last resort — stable per-machine id
return machineId ? `sess_${hashContent(machineId)}` : generateSessionId();
}
// Cleanup expired entries periodically
setInterval(() => {
const now = Date.now();
for (const [key, entry] of assistantSessionMap) {
if (now - entry.lastUsed > SESSION_TTL_MS) assistantSessionMap.delete(key);
}
}, 10 * 60 * 1000);
/**
* Codex Executor - handles OpenAI Codex API (Responses API format)
* Automatically injects default instructions if missing
@@ -377,7 +308,7 @@ export class CodexExecutor extends BaseExecutor {
this._isCompact = !!body._compact;
delete body._compact;
// Resolve conversation-stable session_id (priority: body → assistant-text → workspace → machine)
this._currentSessionId = resolveCacheSessionId(body, credentials, cachedMachineId);
this._currentSessionId = resolveCacheSessionId(body, credentials);
// Convert string input to array format (Codex API requires input as array)
const normalized = normalizeResponsesInput(body.input);
if (normalized) body.input = normalized;

View File

@@ -1,7 +1,8 @@
import { randomUUID } from "crypto";
import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js";
import { convertCommandCodeToOpenAI } from "../translator/response/commandcode-to-openai.js";
import { commandCodeToOpenAIResponse } from "../translator/response/commandcode-to-openai.js";
import { SSE_DONE } from "../utils/sseConstants.js";
/**
* CommandCodeExecutor — talks to https://api.commandcode.ai/alpha/generate
@@ -70,15 +71,15 @@ function wrapNdjsonAsOpenAISse(originalResponse, model) {
const trimmed = line.trim();
if (!trimmed) continue;
// Translate AI SDK v5 NDJSON line to one or more OpenAI chunks
emitChunks(convertCommandCodeToOpenAI(trimmed, state), controller);
emitChunks(commandCodeToOpenAIResponse(trimmed, state), controller);
}
},
flush(controller) {
const trimmed = buffer.trim();
if (trimmed) {
emitChunks(convertCommandCodeToOpenAI(trimmed, state), controller);
emitChunks(commandCodeToOpenAIResponse(trimmed, state), controller);
}
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
},
});

View File

@@ -8,6 +8,8 @@ import {
} from "../utils/cursorProtobuf.js";
import { buildCursorHeaders } from "../utils/cursorChecksum.js";
import { estimateUsage } from "../utils/usageTracking.js";
import { SSE_DONE, SSE_HEADERS } from "../utils/sseConstants.js";
import { chatChunkSse } from "../utils/sse.js";
import { FORMATS } from "../translator/formats.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js";
import zlib from "zlib";
@@ -98,6 +100,32 @@ function decompressPayload(payload, flags) {
return payload;
}
// Read one cursor protobuf frame: header + bounds + decompress. Returns status + payload + new offset.
function readCursorFrame(buffer, offset, frameNum, tag) {
if (offset + 5 > buffer.length) {
debugLog(`[CURSOR BUFFER${tag}] Reached end, offset=${offset}, remaining=${buffer.length - offset}`);
return { status: "done" };
}
const flags = buffer[offset];
const length = buffer.readUInt32BE(offset + 1);
debugLog(`[CURSOR BUFFER${tag}] Frame ${frameNum + 1}: flags=0x${flags.toString(16).padStart(2, "0")}, length=${length}`);
if (offset + 5 + length > buffer.length) {
debugLog(`[CURSOR BUFFER${tag}] Incomplete frame, offset=${offset}, length=${length}, buffer.length=${buffer.length}`);
return { status: "done" };
}
let payload = buffer.slice(offset + 5, offset + 5 + length);
const newOffset = offset + 5 + length;
payload = decompressPayload(payload, flags);
if (!payload) {
debugLog(`[CURSOR BUFFER${tag}] Frame ${frameNum + 1}: decompression failed, skipping`);
return { status: "skip", offset: newOffset };
}
return { status: "ok", payload, offset: newOffset };
}
function createErrorResponse(jsonError) {
const errorMsg = jsonError?.error?.details?.[0]?.debug?.details?.title
|| jsonError?.error?.details?.[0]?.debug?.details?.detail
@@ -141,7 +169,7 @@ export class CursorExecutor extends BaseExecutor {
transformRequest(model, body, stream, credentials) {
// Messages are already translated by chatCore (claude→openai→cursor)
// Do NOT call buildCursorRequest again — double-translation drops tool_results
// Do NOT call openaiToCursorRequest again — double-translation drops tool_results
const messages = body.messages || [];
const tools = body.tools || [];
const reasoningEffort = body.reasoning_effort || null;
@@ -286,36 +314,12 @@ export class CursorExecutor extends BaseExecutor {
debugLog(`[CURSOR BUFFER] Total length: ${buffer.length} bytes`);
while (offset < buffer.length) {
if (offset + 5 > buffer.length) {
debugLog(
`[CURSOR BUFFER] Reached end, offset=${offset}, remaining=${buffer.length - offset}`
);
break;
}
const flags = buffer[offset];
const length = buffer.readUInt32BE(offset + 1);
debugLog(
`[CURSOR BUFFER] Frame ${frameCount + 1}: flags=0x${flags.toString(16).padStart(2, "0")}, length=${length}`
);
if (offset + 5 + length > buffer.length) {
debugLog(
`[CURSOR BUFFER] Incomplete frame, offset=${offset}, length=${length}, buffer.length=${buffer.length}`
);
break;
}
let payload = buffer.slice(offset + 5, offset + 5 + length);
offset += 5 + length;
const frame = readCursorFrame(buffer, offset, frameCount, "");
if (frame.status === "done") break;
offset = frame.offset;
frameCount++;
payload = decompressPayload(payload, flags);
if (!payload) {
debugLog(`[CURSOR BUFFER] Frame ${frameCount}: decompression failed, skipping`);
continue;
}
if (frame.status === "skip") continue;
const payload = frame.payload;
// Check for JSON error frames (byte guard: skip toString on non-JSON frames)
if (payload.length > 0 && payload[0] === 0x7b) {
@@ -466,36 +470,12 @@ export class CursorExecutor extends BaseExecutor {
debugLog(`[CURSOR BUFFER SSE] Total length: ${buffer.length} bytes`);
while (offset < buffer.length) {
if (offset + 5 > buffer.length) {
debugLog(
`[CURSOR BUFFER SSE] Reached end, offset=${offset}, remaining=${buffer.length - offset}`
);
break;
}
const flags = buffer[offset];
const length = buffer.readUInt32BE(offset + 1);
debugLog(
`[CURSOR BUFFER SSE] Frame ${frameCount + 1}: flags=0x${flags.toString(16).padStart(2, "0")}, length=${length}`
);
if (offset + 5 + length > buffer.length) {
debugLog(
`[CURSOR BUFFER SSE] Incomplete frame, offset=${offset}, length=${length}, buffer.length=${buffer.length}`
);
break;
}
let payload = buffer.slice(offset + 5, offset + 5 + length);
offset += 5 + length;
const frame = readCursorFrame(buffer, offset, frameCount, " SSE");
if (frame.status === "done") break;
offset = frame.offset;
frameCount++;
payload = decompressPayload(payload, flags);
if (!payload) {
debugLog(`[CURSOR BUFFER SSE] Frame ${frameCount}: decompression failed, skipping`);
continue;
}
if (frame.status === "skip") continue;
const payload = frame.payload;
// Check for JSON error frames (byte-guard: only decode if starts with '{')
if (payload[0] === 0x7b) {
@@ -542,21 +522,7 @@ export class CursorExecutor extends BaseExecutor {
const tc = result.toolCall;
if (chunks.length === 0) {
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: "chat.completion.chunk",
created,
model,
choices: [
{
index: 0,
delta: { role: "assistant", content: "" },
finish_reason: null
}
]
})}\n\n`
);
chunks.push(chatChunkSse({ id: responseId, created, model, delta: { role: "assistant", content: "" } }));
}
if (toolCallsMap.has(tc.id)) {
@@ -569,33 +535,22 @@ export class CursorExecutor extends BaseExecutor {
// Stream the delta arguments
if (tc.function.arguments) {
emittedToolCallIds.add(tc.id);
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: "chat.completion.chunk",
created,
model,
choices: [
chunks.push(chatChunkSse({
id: responseId, created, model,
delta: {
tool_calls: [
{
index: 0,
delta: {
tool_calls: [
{
index: existing.index,
id: tc.id,
type: "function",
function: {
name: tc.function.name,
arguments: tc.function.arguments
}
}
]
},
finish_reason: null
index: existing.index,
id: tc.id,
type: "function",
function: {
name: tc.function.name,
arguments: tc.function.arguments
}
}
]
})}\n\n`
);
}
}));
}
} else {
// New tool call - assign index and add to map
@@ -606,56 +561,34 @@ export class CursorExecutor extends BaseExecutor {
// Stream initial tool call with name
emittedToolCallIds.add(tc.id);
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: "chat.completion.chunk",
created,
model,
choices: [
chunks.push(chatChunkSse({
id: responseId, created, model,
delta: {
tool_calls: [
{
index: 0,
delta: {
tool_calls: [
{
index: toolCallIndex,
id: tc.id,
type: "function",
function: {
name: tc.function.name,
arguments: tc.function.arguments
}
}
]
},
finish_reason: null
index: toolCallIndex,
id: tc.id,
type: "function",
function: {
name: tc.function.name,
arguments: tc.function.arguments
}
}
]
})}\n\n`
);
}
}));
}
}
if (result.text) {
totalContent += result.text;
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: "chat.completion.chunk",
created,
model,
choices: [
{
index: 0,
delta:
chunks.length === 0 && toolCalls.length === 0
? { role: "assistant", content: result.text }
: { content: result.text },
finish_reason: null
}
]
})}\n\n`
);
chunks.push(chatChunkSse({
id: responseId, created, model,
delta:
chunks.length === 0 && toolCalls.length === 0
? { role: "assistant", content: result.text }
: { content: result.text }
}));
}
if (isComposerModel(model) && result.thinking) {
@@ -665,24 +598,13 @@ export class CursorExecutor extends BaseExecutor {
const deltaContent = visibleContent.slice(emittedComposerThinkingContentLength);
emittedComposerThinkingContentLength = visibleContent.length;
totalContent += deltaContent;
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: "chat.completion.chunk",
created,
model,
choices: [
{
index: 0,
delta:
chunks.length === 0 && toolCalls.length === 0
? { role: "assistant", content: deltaContent }
: { content: deltaContent },
finish_reason: null
}
]
})}\n\n`
);
chunks.push(chatChunkSse({
id: responseId, created, model,
delta:
chunks.length === 0 && toolCalls.length === 0
? { role: "assistant", content: deltaContent }
: { content: deltaContent }
}));
}
}
}
@@ -708,53 +630,28 @@ export class CursorExecutor extends BaseExecutor {
// Emit SSE chunk for the finalized tool call if not already emitted
if (!emittedToolCallIds.has(tc.id)) {
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: "chat.completion.chunk",
created,
model,
choices: [
chunks.push(chatChunkSse({
id: responseId, created, model,
delta: {
tool_calls: [
{
index: 0,
delta: {
tool_calls: [
{
index: toolCallIndex,
id: tc.id,
type: "function",
function: {
name: tc.function.name,
arguments: tc.function.arguments
}
}
]
},
finish_reason: null
index: toolCallIndex,
id: tc.id,
type: "function",
function: {
name: tc.function.name,
arguments: tc.function.arguments
}
}
]
})}\n\n`
);
}
}));
}
}
}
if (chunks.length === 0 && toolCalls.length === 0) {
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: "chat.completion.chunk",
created,
model,
choices: [
{
index: 0,
delta: { role: "assistant", content: "" },
finish_reason: null
}
]
})}\n\n`
);
chunks.push(chatChunkSse({ id: responseId, created, model, delta: { role: "assistant", content: "" } }));
}
const usage = estimateUsage(body, totalContent.length, FORMATS.OPENAI);
@@ -775,15 +672,11 @@ export class CursorExecutor extends BaseExecutor {
usage
})}\n\n`
);
chunks.push("data: [DONE]\n\n");
chunks.push(SSE_DONE);
return new Response(chunks.join(""), {
status: 200,
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"Connection": "keep-alive"
}
headers: { ...SSE_HEADERS }
});
}

View File

@@ -1,10 +1,80 @@
import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js";
import { PROVIDERS, PROVIDER_OAUTH } from "../config/providers.js";
import { ANTHROPIC_API_VERSION, OPENAI_COMPAT_BASE, ANTHROPIC_COMPAT_BASE } from "../providers/shared.js";
import { OAUTH_ENDPOINTS, buildKimiHeaders } from "../config/appConstants.js";
import { buildClineHeaders } from "../../src/shared/utils/clineAuth.js";
import { buildClineHeaders } from "../shared/clineAuth.js";
import { getCachedClaudeHeaders } from "../utils/claudeHeaderCache.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js";
import { injectReasoningContent } from "../utils/reasoningContentInjector.js";
import { stripUnsupportedParams } from "../translator/concerns/paramSupport.js";
// Auth header descriptors — derived from registry transport.auth, fallback to hardcoded defaults.
const BEARER = { combined: true, header: "Authorization", scheme: "bearer" };
const XAPIKEY = { combined: true, header: "x-api-key", scheme: "raw" };
const AUTH_DESCRIPTORS = Object.fromEntries(
Object.entries(PROVIDERS)
.filter(([, t]) => t.auth)
.map(([id, t]) => [id, t.auth])
);
// Apply a token to a header per scheme (matches legacy: combined always sets, even when undefined).
function setAuth(headers, spec, token) {
headers[spec.header] = spec.scheme === "bearer" ? `Bearer ${token}` : token;
}
// Resolve auth onto headers from a descriptor.
function applyAuth(headers, desc, credentials) {
if (desc.combined) {
// combined providers always set the header (legacy behavior, incl. noAuth → "Bearer undefined")
setAuth(headers, desc, credentials.apiKey || credentials.accessToken);
if (desc.anthropicVersion && !headers["anthropic-version"]) headers["anthropic-version"] = ANTHROPIC_API_VERSION;
return;
}
// split apiKey/oauth: set only the matching branch (legacy: anthropic-compatible skips when both absent)
if (credentials.apiKey) setAuth(headers, desc.apiKey, credentials.apiKey);
else if (credentials.accessToken) setAuth(headers, desc.oauth, credentials.accessToken);
if (desc.anthropicVersion && !headers["anthropic-version"]) headers["anthropic-version"] = ANTHROPIC_API_VERSION;
}
// Provider-specific header quirks kept as small hooks (not pure auth).
const HEADER_HOOKS = {
kimiHeaders: (h) => Object.assign(h, buildKimiHeaders()),
clineHeaders: (h, c) => Object.assign(h, buildClineHeaders(c.apiKey || c.accessToken)),
kilocodeOrg: (h, c) => { if (c.providerSpecificData?.orgId) h["X-Kilocode-OrganizationID"] = c.providerSpecificData.orgId; },
claudeOverlay: (h) => {
const cached = getCachedClaudeHeaders();
if (!cached) return;
for (const lcKey of Object.keys(cached)) {
const titleKey = lcKey.replace(/(^|-)([a-z])/g, (_, sep, ch) => sep + ch.toUpperCase());
if (lcKey === "anthropic-beta") {
const staticBetaStr = h[titleKey] || h[lcKey] || "";
const flags = new Set(staticBetaStr.split(",").map(f => f.trim()).filter(Boolean));
for (const f of cached[lcKey].split(",").map(f => f.trim()).filter(Boolean)) flags.add(f);
cached[lcKey] = Array.from(flags).join(",");
}
if (titleKey !== lcKey && h[titleKey] !== undefined) delete h[titleKey];
}
Object.assign(h, cached);
},
};
// Config-driven OAuth refresh grants — derived from registry oauth.refresh.
const REFRESH_GRANTS = Object.fromEntries(
Object.entries(PROVIDER_OAUTH)
.filter(([, o]) => o.refresh)
.map(([id, o]) => {
const tokenUrl = o.tokenUrl;
const encoding = o.refresh.encoding;
const extraParams = o.refresh.scope ? { scope: o.refresh.scope } : {};
return [id, {
encoding,
url: () => tokenUrl,
params: (ex) => id === "gemini"
? { client_id: ex.config.clientId, client_secret: ex.config.clientSecret, ...extraParams }
: { client_id: o.clientId, ...extraParams },
}];
})
);
export class DefaultExecutor extends BaseExecutor {
constructor(provider) {
@@ -15,9 +85,11 @@ export class DefaultExecutor extends BaseExecutor {
const transformed = this.applyJsonSchemaFallback(body);
if (transformed && typeof transformed === "object") {
if (this.provider === "cerebras" || this.provider === "mistral") {
// quirk: some openai-compatible providers reject Anthropic's client_metadata field
if (this.config.quirks?.dropClientMetadata) {
delete transformed.client_metadata;
}
stripUnsupportedParams(this.provider, model, transformed);
}
return injectReasoningContent({ provider: this.provider, model, body: transformed });
@@ -44,120 +116,57 @@ export class DefaultExecutor extends BaseExecutor {
}
buildUrl(model, stream, urlIndex = 0, credentials = null) {
// Runtime transport (multi-endpoint providers): use the sourceFormat-matched endpoint
const rt = credentials?.runtimeTransport;
if (rt?.baseUrl) {
return rt.urlSuffix ? `${rt.baseUrl}${rt.urlSuffix}` : rt.baseUrl;
}
if (this.provider?.startsWith?.("openai-compatible-")) {
const baseUrl = credentials?.providerSpecificData?.baseUrl || "https://api.openai.com/v1";
const baseUrl = credentials?.providerSpecificData?.baseUrl || OPENAI_COMPAT_BASE;
const normalized = baseUrl.replace(/\/$/, "");
const path = this.provider.includes("responses") ? "/responses" : "/chat/completions";
return `${normalized}${path}`;
}
if (this.provider?.startsWith?.("anthropic-compatible-")) {
const baseUrl = credentials?.providerSpecificData?.baseUrl || "https://api.anthropic.com/v1";
const baseUrl = credentials?.providerSpecificData?.baseUrl || ANTHROPIC_COMPAT_BASE;
const normalized = baseUrl.replace(/\/$/, "");
return `${normalized}/messages`;
}
switch (this.provider) {
case "claude":
case "glm":
case "kimi":
case "minimax":
case "minimax-cn":
return `${this.config.baseUrl}?beta=true`;
case "kimi-coding":
return `${this.config.baseUrl}?beta=true`;
case "gemini":
return `${this.config.baseUrl}/${model}:${stream ? "streamGenerateContent?alt=sse" : "generateContent"}`;
default: {
const url = this.config.baseUrl;
if (url?.includes("{accountId}")) {
const accountId = credentials?.providerSpecificData?.accountId;
if (!accountId) throw new Error(`${this.provider} requires accountId in providerSpecificData`);
return url.replace("{accountId}", accountId);
}
return url;
}
// gemini-format: build :streamGenerateContent / :generateContent path
if (this.config.format === "gemini") {
return `${this.config.baseUrl}/${model}:${stream ? "streamGenerateContent?alt=sse" : "generateContent"}`;
}
// urlSuffix (e.g. ?beta=true) declared per-provider in registry
if (this.config.urlSuffix) {
return `${this.config.baseUrl}${this.config.urlSuffix}`;
}
const url = this.config.baseUrl;
if (url?.includes("{accountId}")) {
const accountId = credentials?.providerSpecificData?.accountId;
if (!accountId) throw new Error(`${this.provider} requires accountId in providerSpecificData`);
return url.replace("{accountId}", accountId);
}
return url;
}
// Fallback descriptor for providers without an explicit entry in AUTH_DESCRIPTORS.
resolveAuthDescriptor() {
if (this.provider?.startsWith?.("anthropic-compatible-")) {
return { apiKey: { header: "x-api-key", scheme: "raw" }, oauth: { header: "Authorization", scheme: "bearer" }, anthropicVersion: true };
}
if (this.config?.format === "claude") {
return { ...XAPIKEY, anthropicVersion: true };
}
return BEARER;
}
buildHeaders(credentials, stream = true) {
const headers = { "Content-Type": "application/json", ...this.config.headers };
switch (this.provider) {
case "gemini":
credentials.apiKey ? headers["x-goog-api-key"] = credentials.apiKey : headers["Authorization"] = `Bearer ${credentials.accessToken}`;
break;
case "claude": {
// Overlay live cached headers from real Claude Code client over static defaults.
// Static headers (Title-Case) remain as cold-start fallback.
const cached = getCachedClaudeHeaders();
if (cached) {
// Remove Title-Case static keys that conflict with incoming lowercase cached keys
for (const lcKey of Object.keys(cached)) {
// Build the Title-Case equivalent: "anthropic-version" → "Anthropic-Version"
const titleKey = lcKey.replace(/(^|-)([a-z])/g, (_, sep, c) => sep + c.toUpperCase());
// Special handling for Anthropic-Beta to preserve required flags like OAuth
if (lcKey === "anthropic-beta") {
const staticBetaStr = headers[titleKey] || headers[lcKey] || "";
const staticFlags = new Set(staticBetaStr.split(",").map(f => f.trim()).filter(Boolean));
const cachedFlags = new Set(cached[lcKey].split(",").map(f => f.trim()).filter(Boolean));
// Merge all static flags (which contain oauth, thinking, etc) into the cached ones
for (const flag of staticFlags) {
cachedFlags.add(flag);
}
cached[lcKey] = Array.from(cachedFlags).join(",");
}
if (titleKey !== lcKey && headers[titleKey] !== undefined) {
delete headers[titleKey];
}
}
Object.assign(headers, cached);
}
credentials.apiKey
? (headers["x-api-key"] = credentials.apiKey)
: (headers["Authorization"] = `Bearer ${credentials.accessToken}`);
break;
}
case "glm":
case "kimi":
case "minimax":
case "minimax-cn":
case "kimi-coding":
headers["x-api-key"] = credentials.apiKey || credentials.accessToken;
if (this.provider === "kimi-coding") Object.assign(headers, buildKimiHeaders());
break;
default:
if (this.provider?.startsWith?.("anthropic-compatible-")) {
if (credentials.apiKey) {
headers["x-api-key"] = credentials.apiKey;
} else if (credentials.accessToken) {
headers["Authorization"] = `Bearer ${credentials.accessToken}`;
}
if (!headers["anthropic-version"]) {
headers["anthropic-version"] = "2023-06-01";
}
} else if (this.provider === "gitlab") {
// GitLab Duo uses Bearer token (PAT with ai_features scope, or OAuth access token)
headers["Authorization"] = `Bearer ${credentials.apiKey || credentials.accessToken}`;
} else if (this.provider === "codebuddy") {
headers["Authorization"] = `Bearer ${credentials.apiKey || credentials.accessToken}`;
} else if (this.provider === "kilocode") {
headers["Authorization"] = `Bearer ${credentials.apiKey || credentials.accessToken}`;
if (credentials.providerSpecificData?.orgId) {
headers["X-Kilocode-OrganizationID"] = credentials.providerSpecificData.orgId;
}
} else if (this.provider === "cline") {
Object.assign(headers, buildClineHeaders(credentials.apiKey || credentials.accessToken));
} else if (this.config?.format === "claude") {
// Generic claude-format provider (e.g. agentrouter): x-api-key + anthropic-version
headers["x-api-key"] = credentials.apiKey || credentials.accessToken;
if (!headers["anthropic-version"]) headers["anthropic-version"] = "2023-06-01";
} else {
headers["Authorization"] = `Bearer ${credentials.apiKey || credentials.accessToken}`;
}
}
const rt = credentials?.runtimeTransport;
const headers = { "Content-Type": "application/json", ...(rt ? rt.headers : this.config.headers) };
const desc = rt?.auth || AUTH_DESCRIPTORS[this.provider] || this.resolveAuthDescriptor();
// Hooks run BEFORE auth so dynamic overlays (claude cached headers) can't clobber the token.
for (const hook of desc.hooks || []) HEADER_HOOKS[hook]?.(headers, credentials);
applyAuth(headers, desc, credentials);
// Strip first-party Claude Code identity headers for non-Anthropic anthropic-compatible upstreams
if (this.provider?.startsWith?.("anthropic-compatible-")) {
@@ -196,15 +205,25 @@ export class DefaultExecutor extends BaseExecutor {
return headers;
}
// Generic OAuth refresh for the common {grant_type, refresh_token, client_id[, ...]} shape.
// grant = REFRESH_GRANTS[provider]; client creds resolved from PROVIDERS or this.config.
refreshFromGrant(credentials, proxyOptions) {
const grant = REFRESH_GRANTS[this.provider];
const params = { grant_type: "refresh_token", refresh_token: credentials.refreshToken, ...grant.params(this) };
return grant.encoding === "json"
? this.refreshWithJSON(grant.url(), params, proxyOptions)
: this.refreshWithForm(grant.url(), params, proxyOptions);
}
async refreshCredentials(credentials, log, proxyOptions = null) {
if (!credentials.refreshToken) return null;
const refreshers = {
claude: () => this.refreshWithJSON(OAUTH_ENDPOINTS.anthropic.token, { grant_type: "refresh_token", refresh_token: credentials.refreshToken, client_id: PROVIDERS.claude.clientId }, proxyOptions),
codex: () => this.refreshWithForm(OAUTH_ENDPOINTS.openai.token, { grant_type: "refresh_token", refresh_token: credentials.refreshToken, client_id: PROVIDERS.codex.clientId, scope: "openid profile email offline_access" }, proxyOptions),
claude: () => this.refreshFromGrant(credentials, proxyOptions),
codex: () => this.refreshFromGrant(credentials, proxyOptions),
qwen: () => this.refreshWithForm(OAUTH_ENDPOINTS.qwen.token, { grant_type: "refresh_token", refresh_token: credentials.refreshToken, client_id: PROVIDERS.qwen.clientId }, proxyOptions),
iflow: () => this.refreshIflow(credentials.refreshToken, proxyOptions),
gemini: () => this.refreshGoogle(credentials.refreshToken, proxyOptions),
gemini: () => this.refreshFromGrant(credentials, proxyOptions),
kiro: () => this.refreshKiro(credentials.refreshToken, proxyOptions),
cline: () => this.refreshCline(credentials.refreshToken, proxyOptions),
"kimi-coding": () => this.refreshKimiCoding(credentials.refreshToken, proxyOptions),
@@ -258,17 +277,6 @@ export class DefaultExecutor extends BaseExecutor {
return { accessToken: tokens.access_token, refreshToken: tokens.refresh_token || refreshToken, expiresIn: tokens.expires_in };
}
async refreshGoogle(refreshToken, proxyOptions = null) {
const response = await proxyAwareFetch(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: this.config.clientId, client_secret: this.config.clientSecret })
}, proxyOptions);
if (!response.ok) return null;
const tokens = await response.json();
return { accessToken: tokens.access_token, refreshToken: tokens.refresh_token || refreshToken, expiresIn: tokens.expires_in };
}
async refreshKiro(refreshToken, proxyOptions = null) {
const response = await proxyAwareFetch(PROVIDERS.kiro.tokenUrl, {
method: "POST",
@@ -281,37 +289,29 @@ export class DefaultExecutor extends BaseExecutor {
}
async refreshCline(refreshToken, proxyOptions = null) {
console.log('[DEBUG] Refreshing Cline token, refreshToken length:', refreshToken?.length);
const response = await proxyAwareFetch("https://api.cline.bot/api/v1/auth/refresh", {
const response = await proxyAwareFetch(PROVIDERS.cline.refreshUrl, {
method: "POST",
headers: { "Content-Type": "application/json", "Accept": "application/json" },
body: JSON.stringify({ refreshToken, grantType: "refresh_token", clientType: "extension" })
}, proxyOptions);
console.log('[DEBUG] Cline refresh response status:', response.status);
if (!response.ok) {
const errorText = await response.text();
console.log('[DEBUG] Cline refresh error:', errorText);
return null;
}
if (!response.ok) return null;
const payload = await response.json();
console.log('[DEBUG] Cline refresh payload:', JSON.stringify(payload).substring(0, 200));
const data = payload?.data || payload;
const expiresAtIso = data?.expiresAt;
const expiresIn = expiresAtIso ? Math.max(1, Math.floor((new Date(expiresAtIso).getTime() - Date.now()) / 1000)) : undefined;
console.log('[DEBUG] Cline refresh success, expiresIn:', expiresIn);
return { accessToken: data?.accessToken, refreshToken: data?.refreshToken || refreshToken, expiresIn };
}
async refreshKimiCoding(refreshToken, proxyOptions = null) {
const kimiHeaders = buildKimiHeaders();
const response = await proxyAwareFetch("https://auth.kimi.com/api/oauth/token", {
const response = await proxyAwareFetch(PROVIDERS["kimi-coding"].refreshUrl, {
method: "POST",
headers: {
"Content-Type": "application/x-www-form-urlencoded",
"Accept": "application/json",
...kimiHeaders
},
body: new URLSearchParams({ grant_type: "refresh_token", refresh_token: refreshToken, client_id: "17e5f671-d194-4dfb-9706-5516cb48c098" })
body: new URLSearchParams({ grant_type: "refresh_token", refresh_token: refreshToken, client_id: PROVIDERS["kimi-coding"].clientId })
}, proxyOptions);
if (!response.ok) return null;
const tokens = await response.json();

View File

@@ -7,6 +7,8 @@ import { openaiResponsesToOpenAIResponse } from "../translator/response/openai-r
import { initState } from "../translator/index.js";
import { parseSSELine, formatSSE } from "../utils/streamHelpers.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js";
import { stripUnsupportedParams } from "../translator/concerns/paramSupport.js";
import { SSE_DONE } from "../utils/sseConstants.js";
import crypto from "crypto";
export class GithubExecutor extends BaseExecutor {
@@ -108,54 +110,18 @@ export class GithubExecutor extends BaseExecutor {
return /gpt-5|o[134]-/i.test(model);
}
// Some models (like gpt-5.4) don't support the temperature parameter
supportsTemperature(model) {
// gpt-5.4 and similar newer models don't support temperature
return !/gpt-5\.4/i.test(model);
}
// GitHub Copilot /chat/completions rejects Claude-style thinking payloads
// (OpenClaw sends thinking: { type: "enabled" } → upstream 400).
// GPT-5 family on Copilot DOES honor reasoning_effort, so only strip for Claude. (#713)
supportsThinking(model) {
return !/claude/i.test(model);
}
// reasoning_effort works for GPT-5 family AND Claude Opus 4.6 / Sonnet 4.6
// on GitHub Copilot. Only strip for models that don't support it:
// Claude Haiku 4.5, Claude Opus 4.7 (rejected upstream).
supportsReasoningEffort(model) {
const m = model.toLowerCase();
// Claude models that DO support reasoning_effort
if (/claude.*opus.*4\.6/i.test(m) || /claude.*sonnet.*4\.6/i.test(m)) return true;
// All other Claude models: strip
if (/claude/i.test(model)) return false;
// GPT-5 family, Gemini, etc.: keep
return true;
}
transformRequest(model, body, stream, credentials) {
const transformed = { ...body };
if (this.requiresMaxCompletionTokens(model) && transformed.max_tokens !== undefined) {
transformed.max_completion_tokens = transformed.max_tokens;
delete transformed.max_tokens;
}
// Strip temperature for models that don't support it
if (!this.supportsTemperature(model) && transformed.temperature !== undefined) {
delete transformed.temperature;
}
// Always strip Claude-style thinking payload (Copilot doesn't understand it)
if (!this.supportsThinking(model)) {
delete transformed.thinking;
}
// "none" means no thinking — strip it so models that don't support "none" don't 400
if (transformed.reasoning_effort === "none") {
delete transformed.reasoning_effort;
}
// Strip reasoning_effort only for models that reject it
if (!this.supportsReasoningEffort(model) && transformed.reasoning_effort !== undefined) {
delete transformed.reasoning_effort;
}
// Config-driven strip of params unsupported by this provider/model
stripUnsupportedParams("github", model, transformed);
return transformed;
}
@@ -244,7 +210,7 @@ export class GithubExecutor extends BaseExecutor {
if (!parsed) continue;
if (parsed.done && stream === true) {
controller.enqueue(new TextEncoder().encode("data: [DONE]\n\n"));
controller.enqueue(new TextEncoder().encode(SSE_DONE));
continue;
}

View File

@@ -1,5 +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";
@@ -130,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({
@@ -175,13 +173,13 @@ function buildStreamingResponse(eventStream, model, cid, created, isThinkingMode
id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: fp || null,
choices: [{ index: 0, delta: {}, finish_reason: "stop", logprobs: null }],
})));
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
} catch (err) {
controller.enqueue(encoder.encode(sseChunk({
id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null,
choices: [{ index: 0, delta: { content: `[Stream error: ${err.message || String(err)}]` }, finish_reason: "stop", logprobs: null }],
})));
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
} finally {
controller.close();
}
@@ -333,7 +331,7 @@ export class GrokWebExecutor extends BaseExecutor {
const sseStream = buildStreamingResponse(response.body, model, cid, created, isThinking, signal);
finalResponse = new Response(sseStream, {
status: 200,
headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", "X-Accel-Buffering": "no" },
headers: { ...SSE_HEADERS_NO_BUFFER },
});
} else {
finalResponse = await buildNonStreamingResponse(response.body, model, cid, created, isThinking, signal);

View File

@@ -17,6 +17,7 @@ import { OllamaLocalExecutor } from "./ollama-local.js";
import { CommandCodeExecutor } from "./commandcode.js";
import { XiaomiTokenplanExecutor } from "./xiaomi-tokenplan.js";
import { MimoFreeExecutor } from "./mimo-free.js";
import { CodeBuddyExecutor } from "./codebuddy-cn.js";
import { DefaultExecutor } from "./default.js";
const executors = {
@@ -42,6 +43,7 @@ const executors = {
"xiaomi-tokenplan": new XiaomiTokenplanExecutor(),
"mimo-free": new MimoFreeExecutor(),
mmf: new MimoFreeExecutor(), // Alias for mimo-free
"codebuddy-cn": new CodeBuddyExecutor(),
};
const defaultCache = new Map();
@@ -77,3 +79,4 @@ export { OllamaLocalExecutor } from "./ollama-local.js";
export { CommandCodeExecutor } from "./commandcode.js";
export { XiaomiTokenplanExecutor } from "./xiaomi-tokenplan.js";
export { MimoFreeExecutor } from "./mimo-free.js";
export { CodeBuddyExecutor } from "./codebuddy-cn.js";

View File

@@ -2,6 +2,7 @@ import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js";
import { v4 as uuidv4 } from "uuid";
import { refreshKiroToken } from "../services/tokenRefresh.js";
import { SSE_DONE, SSE_HEADERS } from "../utils/sseConstants.js";
/**
* KiroExecutor - Executor for Kiro AI (AWS CodeWhisperer)
@@ -19,13 +20,50 @@ export class KiroExecutor extends BaseExecutor {
"Amz-Sdk-Invocation-Id": uuidv4()
};
if (credentials.accessToken) {
// API-key auth: the key is stored as accessToken and sent as a bearer token
// exactly like an OAuth access token, but with an extra `tokentype: API_KEY`
// header so CodeWhisperer treats it as a long-lived API key rather than an
// OIDC/social access token. Mirrors the Kiro IDE headless-auth behavior.
const isApiKey = credentials?.providerSpecificData?.authMethod === "api_key";
const apiKey = credentials?.apiKey || (isApiKey ? credentials?.accessToken : null);
if (isApiKey && apiKey) {
headers["Authorization"] = `Bearer ${apiKey}`;
headers["tokentype"] = "API_KEY";
} else if (credentials.accessToken) {
headers["Authorization"] = `Bearer ${credentials.accessToken}`;
}
return headers;
}
/**
* Auth-aware endpoint ordering.
*
* API-key Kiro connections store a raw CodeWhisperer credential (validated
* against codewhisperer.us-east-1.amazonaws.com via ListAvailableProfiles).
* The Kiro IDE gateway (runtime.*.kiro.dev) expects Kiro OIDC/social tokens
* and rejects an `tokentype: API_KEY` token with 401/403 — which
* BaseExecutor.execute() returns immediately (only 429 / network errors fall
* through to the next host). So for api-key auth we must try the *.amazonaws.com
* CodeWhisperer hosts FIRST, mirroring the Kiro-Go reference fork which never
* routes api-key traffic through kiro.dev. OAuth keeps the default order
* (kiro.dev first) since its token is what that gateway accepts.
*/
getOrderedBaseUrls(credentials) {
const baseUrls = this.getBaseUrls();
const isApiKey = credentials?.providerSpecificData?.authMethod === "api_key";
if (!isApiKey) return baseUrls;
const amazon = baseUrls.filter((u) => u.includes("amazonaws.com"));
const others = baseUrls.filter((u) => !u.includes("amazonaws.com"));
return amazon.length > 0 ? [...amazon, ...others] : baseUrls;
}
buildUrl(model, stream, urlIndex = 0, credentials = null) {
const baseUrls = this.getOrderedBaseUrls(credentials);
return baseUrls[urlIndex] || baseUrls[0] || this.config.baseUrl;
}
transformRequest(model, body, stream, credentials) {
return body;
}
@@ -37,6 +75,8 @@ export class KiroExecutor extends BaseExecutor {
* BaseExecutor.execute() walks config.baseUrls (runtime.us-east-1.kiro.dev →
* codewhisperer → q) advancing to the next host on 429 (shouldRetry) and on
* network/5xx errors, while tryRetry handles in-place retries per `retry: {429: 2}`.
* Note: api-key connections reorder these so the *.amazonaws.com hosts come
* first — see getOrderedBaseUrls/buildUrl above.
* Note: the baseUrls are alternate surfaces of one regional service, so rotation
* is edge-level failover — it does not grant fresh 429 quota. Per-account 429
* spreading is handled upstream by account rotation in sse/handlers/chat.js.
@@ -73,6 +113,8 @@ export class KiroExecutor extends BaseExecutor {
const transformStream = new TransformStream({
async transform(chunk, controller) {
// Track output so we can emit a keepalive if this frame yields no chunk.
const enqueueCountBefore = chunkIndex;
// Append to buffer
const newBuffer = new Uint8Array(buffer.length + chunk.length);
newBuffer.set(buffer);
@@ -96,7 +138,7 @@ export class KiroExecutor extends BaseExecutor {
if (!event) continue;
const eventType = event.headers[":event-type"] || "";
// Track total content length for token estimation
if (!state.totalContentLength) state.totalContentLength = 0;
if (!state.contextUsagePercentage) state.contextUsagePercentage = 0;
@@ -105,7 +147,7 @@ export class KiroExecutor extends BaseExecutor {
if (eventType === "assistantResponseEvent" && event.payload?.content) {
const content = event.payload.content;
state.totalContentLength += content.length;
const chunk = {
id: responseId,
object: "chat.completion.chunk",
@@ -292,7 +334,7 @@ export class KiroExecutor extends BaseExecutor {
if (metrics && typeof metrics === 'object') {
const inputTokens = metrics.inputTokens || 0;
const outputTokens = metrics.outputTokens || 0;
if (inputTokens > 0 || outputTokens > 0) {
state.usage = {
prompt_tokens: inputTokens,
@@ -306,27 +348,27 @@ export class KiroExecutor extends BaseExecutor {
// Emit final chunk only after receiving BOTH meteringEvent AND contextUsageEvent
if (state.hasMeteringEvent && state.hasContextUsage && !state.finishEmitted) {
state.finishEmitted = true;
// Estimate tokens if not available from events
if (!state.usage) {
// Estimate output tokens from content length
const estimatedOutputTokens = state.totalContentLength > 0
const estimatedOutputTokens = state.totalContentLength > 0
? Math.max(1, Math.floor(state.totalContentLength / 4))
: 0;
// Estimate input tokens from contextUsagePercentage
// Kiro models typically have 200k context window
const estimatedInputTokens = state.contextUsagePercentage > 0
? Math.floor(state.contextUsagePercentage * 200000 / 100)
: 0;
state.usage = {
prompt_tokens: estimatedInputTokens,
completion_tokens: estimatedOutputTokens,
total_tokens: estimatedInputTokens + estimatedOutputTokens
};
}
const finishChunk = {
id: responseId,
object: "chat.completion.chunk",
@@ -338,12 +380,12 @@ export class KiroExecutor extends BaseExecutor {
finish_reason: state.hasToolCalls ? "tool_calls" : "stop"
}]
};
// Include usage in final chunk if available
if (state.usage) {
finishChunk.usage = state.usage;
}
controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify(finishChunk)}\n\n`));
}
}
@@ -351,6 +393,12 @@ export class KiroExecutor extends BaseExecutor {
if (iterations >= maxIterations) {
console.warn("[Kiro] Max iterations reached in event parsing");
}
// No client chunk produced this frame — emit an SSE comment keepalive
// so the stall watchdog sees upstream activity (ignored by parser/client).
if (chunkIndex === enqueueCountBefore && !state.finishEmitted) {
controller.enqueue(new TextEncoder().encode(": ka\n\n"));
}
},
flush(controller) {
@@ -372,24 +420,20 @@ export class KiroExecutor extends BaseExecutor {
}
// Send final done message
controller.enqueue(new TextEncoder().encode("data: [DONE]\n\n"));
controller.enqueue(new TextEncoder().encode(SSE_DONE));
}
});
// Pipe response body through transform stream
if (!response.body) {
return new Response("data: [DONE]\n\n", { status: response.status, headers: { "Content-Type": "text/event-stream" } });
return new Response(SSE_DONE, { status: response.status, headers: { "Content-Type": "text/event-stream" } });
}
const transformedStream = response.body.pipeThrough(transformStream);
return new Response(transformedStream, {
status: response.status,
statusText: response.statusText,
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"Connection": "keep-alive"
}
headers: { ...SSE_HEADERS }
});
}

View File

@@ -5,13 +5,20 @@ import { createHash } from "crypto";
import os from "os";
const BOOTSTRAP_URL = "https://api.xiaomimimo.com/api/free-ai/bootstrap";
const CHAT_URL = "https://api.xiaomimimo.com/api/free-ai/openai/chat";
const CHAT_URL = PROVIDERS["mimo-free"].baseUrl;
const SESSION_AFFINITY_PREFIX = "ses_";
const SESSION_ID_LENGTH = 24;
const JWT_FALLBACK_TTL_SEC = 3000;
const JWT_EXPIRY_BUFFER_MS = 300000;
const SESSION_CHARS = "abcdefghijklmnopqrstuvwxyz0123456789";
// Anti-abuse gate: upstream rejects requests without a Chrome-like User-Agent with 403 "Illegal access"
const USER_AGENTS = [
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36",
"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36",
];
// Anti-abuse gate marker: the free chat endpoint returns 403 "Illegal access"
// unless a system message contains this exact MiMoCode signature substring.
export const MIMO_SYSTEM_MARKER =
@@ -76,7 +83,10 @@ async function bootstrapJwt(proxyOptions = null) {
const response = await proxyAwareFetch(BOOTSTRAP_URL, {
method: "POST",
headers: { "Content-Type": "application/json" },
headers: {
"Content-Type": "application/json",
"User-Agent": USER_AGENTS[Math.floor(Math.random() * USER_AGENTS.length)],
},
body: JSON.stringify({ client: generateFingerprint() }),
}, proxyOptions);
@@ -108,6 +118,7 @@ export class MimoFreeExecutor extends BaseExecutor {
return {
"Content-Type": "application/json",
"X-Mimo-Source": "mimocode-cli-free",
"User-Agent": USER_AGENTS[Math.floor(Math.random() * USER_AGENTS.length)],
"x-session-affinity": this.sessionId,
"Accept": stream ? "text/event-stream" : "application/json",
};

View File

@@ -1,9 +1,17 @@
import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js";
import { injectReasoningContent } from "../utils/reasoningContentInjector.js";
import { ANTHROPIC_API_VERSION } from "../providers/shared.js";
// Models that use /zen/go/v1/messages (Anthropic/Claude format + x-api-key auth)
const CLAUDE_FORMAT_MODELS = new Set(["minimax-m2.5", "minimax-m2.7"]);
const MESSAGES_FORMAT_MODELS = new Set([
"minimax-m3",
"minimax-m2.7",
"minimax-m2.5",
"qwen3.7-max",
"qwen3.7-plus",
"qwen3.6-plus",
]);
const BASE = "https://opencode.ai/zen/go/v1";
@@ -15,7 +23,7 @@ export class OpenCodeGoExecutor extends BaseExecutor {
// buildUrl runs before buildHeaders in BaseExecutor.execute, cache model here
buildUrl(model) {
this._lastModel = model;
return CLAUDE_FORMAT_MODELS.has(model)
return MESSAGES_FORMAT_MODELS.has(model)
? `${BASE}/messages`
: `${BASE}/chat/completions`;
}
@@ -24,9 +32,9 @@ export class OpenCodeGoExecutor extends BaseExecutor {
const key = credentials?.apiKey || credentials?.accessToken;
const headers = { "Content-Type": "application/json" };
if (CLAUDE_FORMAT_MODELS.has(this._lastModel)) {
if (MESSAGES_FORMAT_MODELS.has(this._lastModel)) {
headers["x-api-key"] = key;
headers["anthropic-version"] = "2023-06-01";
headers["anthropic-version"] = ANTHROPIC_API_VERSION;
} else {
headers["Authorization"] = `Bearer ${key}`;
}

View File

@@ -15,7 +15,7 @@ export class OpenCodeExecutor extends BaseExecutor {
}
buildUrl(model) {
const base = "https://opencode.ai";
const base = this.config.baseUrl;
return MESSAGES_MODELS.has(model)
? `${base}/zen/v1/messages`
: `${base}/zen/v1/chat/completions`;

View File

@@ -1,5 +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";
@@ -289,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({
@@ -340,7 +338,7 @@ function buildStreamingResponse(eventStream, model, cid, created, history, curre
id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null,
choices: [{ index: 0, delta: {}, finish_reason: "stop", logprobs: null }],
})));
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
sessionStore(history, currentMsg, cleanResponse(fullAnswer), respBackendUuid);
} catch (err) {
@@ -348,7 +346,7 @@ function buildStreamingResponse(eventStream, model, cid, created, history, curre
id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null,
choices: [{ index: 0, delta: { content: `[Stream error: ${err.message || String(err)}]` }, finish_reason: "stop", logprobs: null }],
})));
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
} finally {
controller.close();
}
@@ -493,7 +491,7 @@ export class PerplexityWebExecutor extends BaseExecutor {
const sseStream = buildStreamingResponse(response.body, model, cid, created, parsed.history, parsed.currentMsg, signal);
finalResponse = new Response(sseStream, {
status: 200,
headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", "X-Accel-Buffering": "no" },
headers: { ...SSE_HEADERS_NO_BUFFER },
});
} else {
finalResponse = await buildNonStreamingResponse(response.body, model, cid, created, parsed.history, parsed.currentMsg, signal);

View File

@@ -20,19 +20,20 @@
* different model upstream, so a missing entry is a hard error.
*/
import { qoderEncodeBody } from "@/lib/qoder/encoding.js";
import { buildCosyHeaders } from "@/lib/qoder/cosy.js";
import { qoderEncodeBody } from "../shared/qoder/encoding.js";
import { buildCosyHeaders } from "../shared/qoder/cosy.js";
import { v4 as uuidv4 } from "uuid";
import { createHash } from "crypto";
import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js";
import { SSE_DONE } from "../utils/sseConstants.js";
import { FETCH_CONNECT_TIMEOUT_MS } from "../config/runtimeConfig.js";
import {
QODER_CHAT_URL_ENCODED,
QODER_MODEL_MAP,
} from "@/lib/qoder/constants.js";
} from "../shared/qoder/constants.js";
import { getQoderModelConfig, resolveQoderModels } from "../services/qoderModels.js";
/**
@@ -356,7 +357,7 @@ async function wrapQoderSSE(response, model, midStreamError = {}) {
const data = trimmed.slice(5).trimStart();
if (data === "[DONE]") {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
doneEmitted = true;
return;
}
@@ -377,13 +378,13 @@ async function wrapQoderSSE(response, model, midStreamError = {}) {
choices: [{ index: 0, delta: { content: `\n\n${parsed.message}` }, finish_reason: "stop" }],
});
controller.enqueue(encoder.encode(`data: ${errChunk}\n\n`));
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
doneEmitted = true;
return;
}
if (!inner) return;
if (inner === "[DONE]") {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
doneEmitted = true;
return;
}
@@ -408,7 +409,7 @@ async function wrapQoderSSE(response, model, midStreamError = {}) {
buffer = "";
}
if (!doneEmitted) {
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
controller.enqueue(encoder.encode(SSE_DONE));
doneEmitted = true;
}
},

View File

@@ -1,18 +1,19 @@
import { DefaultExecutor } from "./default.js";
import { resolveXiaomiTokenplanBaseUrl } from "../config/providers.js";
import { getModelTargetFormat } from "../config/providerModels.js";
import { FORMATS } from "../translator/formats.js";
// import { getModelTargetFormat } from "../config/providerModels.js";
// import { FORMATS } from "../translator/formats.js";
export class XiaomiTokenplanExecutor extends DefaultExecutor {
constructor() {
super("xiaomi-tokenplan");
}
// Claude-native aliases route to the Anthropic-compatible messages endpoint
// Token Plan keys are region-specific. Route per sourceFormat-matched transport:
// claude → Anthropic /anthropic/v1/messages, openai → /chat/completions.
buildUrl(model, stream, urlIndex = 0, credentials = null) {
const baseUrl = resolveXiaomiTokenplanBaseUrl(credentials);
if (getModelTargetFormat(model, model) === FORMATS.CLAUDE) {
return `${baseUrl.replace(/\/v1\/?$/, "/anthropic/v1")}/messages`;
if (credentials?.runtimeTransport?.format === "claude") {
return `${baseUrl.replace(/\/v1\/?$/, "")}/anthropic/v1/messages`;
}
return `${baseUrl}/chat/completions`;
}