Merge remote-tracking branch 'origin/master' into gitea/new_feature
# Conflicts: # open-sse/executors/qoder.js # open-sse/handlers/chatCore.js # open-sse/handlers/chatCore/sseToJsonHandler.js # open-sse/providers/registry/commandcode.js # src/app/(dashboard)/dashboard/combos/page.js # src/app/api/v1/models/route.js # src/lib/db/repos/usageRepo.js # src/shared/components/UsageStats.js
This commit is contained in:
@@ -293,12 +293,19 @@ export class AntigravityExecutor extends BaseExecutor {
|
||||
|
||||
this._lastSessionId = transformedRequest.sessionId; // cached for buildHeaders (base.execute order)
|
||||
|
||||
// Official Antigravity client omits `requestType` entirely on the agent
|
||||
// (chat) path. Sending `requestType: "agent"` here (or leaking it through
|
||||
// from an upstream envelope via the ...body spread below) makes Google
|
||||
// bucket the request and return a detail-free 429 RESOURCE_EXHAUSTED even
|
||||
// with quota available. `image_gen` and
|
||||
// `search` buckets are unaffected and keep their own requestType.
|
||||
delete body.requestType;
|
||||
|
||||
return {
|
||||
...body,
|
||||
project: projectId,
|
||||
model: body.model || model,
|
||||
userAgent: "antigravity",
|
||||
requestType: "agent",
|
||||
requestId: buildIdeRequestId({ body, request: transformedRequest, credentials, model, requestType: "agent" }),
|
||||
request: transformedRequest
|
||||
};
|
||||
|
||||
@@ -7,13 +7,16 @@ import {
|
||||
wrapConnectRPCFrame,
|
||||
decodeMessage,
|
||||
parseConnectRPCFrame,
|
||||
extractTextFromResponse
|
||||
extractTextFromResponse,
|
||||
encodeMcpTools,
|
||||
decodeMcpArgs,
|
||||
} 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, sseChunk } from "../utils/sse.js";
|
||||
import { FORMATS } from "../translator/formats.js";
|
||||
import { ROLE, OPENAI_BLOCK } from "../translator/schema/index.js";
|
||||
import { proxyAwareFetch } from "../utils/proxyFetch.js";
|
||||
import zlib from "zlib";
|
||||
import crypto from "crypto";
|
||||
@@ -65,55 +68,74 @@ function textFromContent(content) {
|
||||
if (typeof content === "string") return content;
|
||||
if (!Array.isArray(content)) return "";
|
||||
return content
|
||||
.filter((part) => part?.type === "text" && typeof part.text === "string")
|
||||
.filter((part) => part?.type === OPENAI_BLOCK.TEXT && typeof part.text === "string")
|
||||
.map((part) => part.text)
|
||||
.join("\n");
|
||||
}
|
||||
|
||||
function isAgentTextRequest(body) {
|
||||
// Many compatible clients always attach their built-in tool schemas, even
|
||||
// for a normal text turn. Cursor's retired ChatService rejects those
|
||||
// requests; AgentService can still answer the text turn, so ignore schemas
|
||||
// here. A real tool-call/result conversation is kept on the legacy path
|
||||
// until its AgentService tool protocol is implemented.
|
||||
return Array.isArray(body?.messages) && body.messages.every((message) => {
|
||||
if (message?.tool_calls?.length || message?.role === "tool") return false;
|
||||
return typeof message?.content === "string"
|
||||
|| Array.isArray(message?.content) && message.content.every((part) => part?.type === "text");
|
||||
function isTextPart(part) {
|
||||
return !part || part.type === OPENAI_BLOCK.TEXT || typeof part === "string";
|
||||
}
|
||||
|
||||
export function isAgentCapableRequest(body) {
|
||||
// ChatService rejects auto/composer and most thinking variants. AgentService
|
||||
// can answer text turns (including declared tool schemas) and tool-call
|
||||
// history. Image parts still need the legacy protobuf path.
|
||||
if (!Array.isArray(body?.messages) || body.messages.length === 0) return false;
|
||||
return body.messages.every((message) => {
|
||||
if (Array.isArray(message?.content)) return message.content.every(isTextPart);
|
||||
return message?.content == null || typeof message.content === "string";
|
||||
});
|
||||
}
|
||||
|
||||
function encodeHistoryMessage(message) {
|
||||
const content = textFromContent(message?.content);
|
||||
if (!content) return null;
|
||||
const extras = [];
|
||||
if (message?.role === ROLE.ASSISTANT && message.tool_calls?.length) {
|
||||
for (const tc of message.tool_calls) {
|
||||
extras.push(`[tool_call id=${tc.id || ""} name=${tc.function?.name || "tool"} args=${tc.function?.arguments || "{}"}]`);
|
||||
}
|
||||
}
|
||||
if (message?.role === ROLE.TOOL) {
|
||||
extras.push(`[tool_result id=${message.tool_call_id || ""}]`);
|
||||
}
|
||||
const textBody = [content, ...extras].filter(Boolean).join("\n");
|
||||
if (!textBody) return null;
|
||||
|
||||
// ConversationHistoryMessage.user / .assistant -> repeated content -> text.
|
||||
const text = agentString(1, content);
|
||||
if (message.role === "assistant") {
|
||||
const text = agentString(1, textBody);
|
||||
if (message.role === ROLE.ASSISTANT) {
|
||||
return agentMessage(2, agentMessage(1, agentMessage(1, text)));
|
||||
}
|
||||
return agentMessage(1, agentMessage(1, agentMessage(1, text)));
|
||||
}
|
||||
|
||||
function buildAgentRunFrame(messages, model) {
|
||||
export function buildAgentRunFrame(messages, model, tools = []) {
|
||||
// custom_system_prompt (RunRequest field 8) makes AgentService return an
|
||||
// empty turn. Fold system text into the current user message instead.
|
||||
const system = messages
|
||||
.filter((message) => message?.role === "system")
|
||||
.filter((message) => message?.role === ROLE.SYSTEM)
|
||||
.map((message) => textFromContent(message.content))
|
||||
.filter(Boolean)
|
||||
.join("\n\n");
|
||||
const chatMessages = messages.filter((message) => message?.role !== "system");
|
||||
const currentIndex = [...chatMessages].map((message) => message?.role).lastIndexOf("user");
|
||||
const chatMessages = messages.filter((message) => message?.role !== ROLE.SYSTEM);
|
||||
const currentIndex = [...chatMessages].map((message) => message?.role).lastIndexOf(ROLE.USER);
|
||||
const current = currentIndex >= 0 ? chatMessages[currentIndex] : chatMessages.at(-1);
|
||||
const history = chatMessages
|
||||
.slice(0, currentIndex >= 0 ? currentIndex : -1)
|
||||
.map(encodeHistoryMessage)
|
||||
.filter(Boolean);
|
||||
const userText = textFromContent(current?.content) || "Continue.";
|
||||
const rawUser = textFromContent(current?.content) || "Continue.";
|
||||
const userText = system ? `${system}\n\n${rawUser}` : rawUser;
|
||||
|
||||
// agent.v1.UserMessageAction.user_message and its optional history.
|
||||
// selected_context (3) + mode=1 (4) match cursor-agent's wire format; without
|
||||
// them the server may accept the RPC and stream an empty turn.
|
||||
const userMessage = concatBuffers(
|
||||
agentString(1, userText),
|
||||
agentString(2, crypto.randomUUID()),
|
||||
agentMessage(3, new Uint8Array()),
|
||||
encodeField(4, PROTOBUF_VARINT, 1),
|
||||
);
|
||||
const conversationHistory = history.length
|
||||
? concatBuffers(...history.map((entry) => agentMessage(1, entry)))
|
||||
@@ -124,11 +146,20 @@ function buildAgentRunFrame(messages, model) {
|
||||
);
|
||||
const conversationAction = agentMessage(1, userAction);
|
||||
const requestedModel = concatBuffers(agentString(1, model), agentBool(7, true));
|
||||
// ModelDetails (field 3): thinking variants (Composer, Grok, *-thinking)
|
||||
// return an empty turn when only RequestedModel (field 9) is set.
|
||||
const modelDetails = concatBuffers(
|
||||
agentString(1, model),
|
||||
agentString(3, model),
|
||||
agentString(4, model),
|
||||
);
|
||||
const mcpTools = encodeMcpTools(tools);
|
||||
const runRequest = concatBuffers(
|
||||
// An empty ConversationStateStructure starts a fresh local agent session.
|
||||
agentMessage(1, new Uint8Array()),
|
||||
agentMessage(2, conversationAction),
|
||||
...(system ? [agentString(8, system)] : []),
|
||||
agentMessage(3, modelDetails),
|
||||
...(mcpTools.length ? [agentMessage(4, mcpTools)] : []),
|
||||
agentMessage(9, requestedModel),
|
||||
);
|
||||
|
||||
@@ -157,13 +188,51 @@ function decodeAgentFrames(buffer, onFrame) {
|
||||
return pending;
|
||||
}
|
||||
|
||||
function createRequestContextResponse() {
|
||||
// AgentService asks every run for client context. 9router has no IDE file
|
||||
// context, so acknowledge with an empty RequestContext.
|
||||
function execIds(execRequest) {
|
||||
const id = Number(execRequest?.get(1)?.[0]?.value || 0);
|
||||
const execId = extractAgentString(execRequest, 15);
|
||||
return { id, execId };
|
||||
}
|
||||
|
||||
function wrapExecClientMessage(execMsgId, execId, resultField, resultPayload) {
|
||||
const parts = [];
|
||||
if (execMsgId) parts.push(encodeField(1, PROTOBUF_VARINT, execMsgId));
|
||||
parts.push(agentString(15, execId || ""));
|
||||
parts.push(encodeField(resultField, PROTOBUF_LEN, resultPayload || new Uint8Array()));
|
||||
return wrapConnectRPCFrame(agentMessage(2, concatBuffers(...parts)));
|
||||
}
|
||||
|
||||
function createRequestContextResponse(execRequest) {
|
||||
// Tools already go out on AgentRunRequest.mcp_tools. Echoing them again on
|
||||
// this ack makes AgentService stall silently (0 SSE bytes until abort).
|
||||
const { id, execId } = execIds(execRequest);
|
||||
const requestContextSuccess = agentMessage(1, new Uint8Array());
|
||||
const requestContextResult = agentMessage(1, requestContextSuccess);
|
||||
const execClientMessage = agentMessage(10, requestContextResult);
|
||||
return wrapConnectRPCFrame(agentMessage(2, execClientMessage));
|
||||
return wrapExecClientMessage(id, execId, 10, requestContextResult);
|
||||
}
|
||||
|
||||
// ExecServerMessage variant → ExecClientMessage result field (same numbers).
|
||||
const EXEC_RESULT_FIELD = {
|
||||
2: 2, 3: 3, 4: 4, 5: 5, 7: 7, 8: 8, 9: 9, 16: 16, 20: 20, 23: 23,
|
||||
};
|
||||
|
||||
function rejectExecRequest(execRequest) {
|
||||
const { id, execId } = execIds(execRequest);
|
||||
const variant = [...(execRequest?.keys?.() || [])].find((field) => field !== 1 && field !== 15);
|
||||
const resultField = EXEC_RESULT_FIELD[variant];
|
||||
if (!resultField) return null;
|
||||
// Diagnostics has no rejected variant — empty success unblocks the stream.
|
||||
if (variant === 9) return wrapExecClientMessage(id, execId, 9, new Uint8Array());
|
||||
const rejected = agentMessage(2, agentString(2, "Tool not available in this environment. Use the MCP tools provided instead."));
|
||||
return wrapExecClientMessage(id, execId, resultField, rejected);
|
||||
}
|
||||
|
||||
function encodeKvClientMessage(kvId, resultField, resultPayload, metadata) {
|
||||
const parts = [];
|
||||
if (kvId) parts.push(encodeField(1, PROTOBUF_VARINT, kvId));
|
||||
parts.push(encodeField(resultField, PROTOBUF_LEN, resultPayload || new Uint8Array()));
|
||||
if (metadata && metadata.length) parts.push(encodeField(4, PROTOBUF_LEN, metadata));
|
||||
return wrapConnectRPCFrame(agentMessage(3, concatBuffers(...parts)));
|
||||
}
|
||||
|
||||
const CURSOR_STREAM_DEBUG = process.env.CURSOR_STREAM_DEBUG === "1";
|
||||
@@ -479,7 +548,7 @@ export class CursorExecutor extends BaseExecutor {
|
||||
};
|
||||
}
|
||||
|
||||
async executeAgent({ model, body, stream, credentials, signal }) {
|
||||
async executeAgent({ model, body, stream, credentials, signal, log }) {
|
||||
const agentEndpoint = PROVIDER_OAUTH.cursor?.agentEndpoint;
|
||||
if (!agentEndpoint) throw new Error("Cursor AgentService endpoint is not configured");
|
||||
|
||||
@@ -491,9 +560,10 @@ export class CursorExecutor extends BaseExecutor {
|
||||
}
|
||||
|
||||
let session;
|
||||
const tools = body.tools || [];
|
||||
try {
|
||||
session = this.openAgentHttp2Stream(url, headers, requestController.signal);
|
||||
session.write(buildAgentRunFrame(body.messages || [], model));
|
||||
session.write(buildAgentRunFrame(body.messages || [], model, tools));
|
||||
} catch (error) {
|
||||
throw new Error(`Cursor AgentService request failed: ${error.message}`);
|
||||
}
|
||||
@@ -533,8 +603,23 @@ export class CursorExecutor extends BaseExecutor {
|
||||
// so strict clients such as Claude Code accept the completed stream.
|
||||
const responseId = `chatcmpl-msg_${Date.now()}`;
|
||||
const created = Math.floor(Date.now() / 1000);
|
||||
const composerModel = isComposerModel(model);
|
||||
let pending = Buffer.alloc(0);
|
||||
let finished = false;
|
||||
let thinkingAcc = "";
|
||||
let emittedVisible = 0;
|
||||
let emittedText = false;
|
||||
|
||||
const flushThinkingFallback = (onEvent) => {
|
||||
if (emittedText || !thinkingAcc) return;
|
||||
const fallback = composerModel
|
||||
? visibleComposerContentFromThinking(thinkingAcc)
|
||||
: thinkingAcc.trim();
|
||||
if (fallback) {
|
||||
emittedText = true;
|
||||
onEvent({ type: "text", value: fallback });
|
||||
}
|
||||
};
|
||||
|
||||
const consume = async (onEvent) => {
|
||||
try {
|
||||
@@ -553,32 +638,87 @@ export class CursorExecutor extends BaseExecutor {
|
||||
const update = decodeMessage(serverMessage.get(1)[0].value);
|
||||
if (update.has(1)) {
|
||||
const textDelta = extractAgentString(decodeMessage(update.get(1)[0].value), 1);
|
||||
if (textDelta) onEvent({ type: "text", value: textDelta });
|
||||
if (textDelta) {
|
||||
emittedText = true;
|
||||
onEvent({ type: "text", value: textDelta });
|
||||
}
|
||||
}
|
||||
// Cursor's AgentService emits internal reasoning without the
|
||||
// cryptographic signature required by Anthropic thinking blocks.
|
||||
// Forwarding it makes strict Anthropic clients (Claude Code)
|
||||
// discard or wait on an otherwise complete response. Keep the
|
||||
// reasoning upstream-only and emit the normal answer text.
|
||||
// thinking_delta (field 4). Composer (and some Grok variants) put
|
||||
// the visible answer after </think> here and never send text_delta.
|
||||
if (update.has(4)) {
|
||||
const thinkingDelta = extractAgentString(decodeMessage(update.get(4)[0].value), 1);
|
||||
if (thinkingDelta) {
|
||||
thinkingAcc += thinkingDelta;
|
||||
if (composerModel) {
|
||||
const visible = visibleComposerContentFromThinking(thinkingAcc);
|
||||
if (visible.length > emittedVisible) {
|
||||
const deltaContent = visible.slice(emittedVisible);
|
||||
emittedVisible = visible.length;
|
||||
emittedText = true;
|
||||
onEvent({ type: "text", value: deltaContent });
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
// Keep unsigned reasoning upstream-only for Anthropic clients.
|
||||
if (update.has(14)) {
|
||||
flushThinkingFallback(onEvent);
|
||||
finished = true;
|
||||
onEvent({ type: "done" });
|
||||
}
|
||||
}
|
||||
|
||||
// KvServerMessage (field 4): get/set blob. Ack so the stream proceeds.
|
||||
if (serverMessage.has(4)) {
|
||||
const kv = decodeMessage(serverMessage.get(4)[0].value);
|
||||
const kvId = kv.get(1)?.[0]?.value || 0;
|
||||
const metadata = kv.get(4)?.[0]?.value || null;
|
||||
if (kv.has(2)) {
|
||||
session.write(encodeKvClientMessage(kvId, 2, agentMessage(1, new Uint8Array()), metadata));
|
||||
} else if (kv.has(3)) {
|
||||
session.write(encodeKvClientMessage(kvId, 3, new Uint8Array(), metadata));
|
||||
}
|
||||
}
|
||||
|
||||
// AgentService requests IDE context before producing a response.
|
||||
// Return an empty context; 9router is not coupled to an editor.
|
||||
if (serverMessage.has(2)) {
|
||||
const execRequest = decodeMessage(serverMessage.get(2)[0].value);
|
||||
if (execRequest.has(10)) {
|
||||
session.write(createRequestContextResponse());
|
||||
log?.info?.("CURSOR", "AgentService request_context ack");
|
||||
session.write(createRequestContextResponse(execRequest));
|
||||
} else if (execRequest.has(11)) {
|
||||
const mcp = decodeMcpArgs(execRequest.get(11)[0].value);
|
||||
const name = mcp.toolName || mcp.name;
|
||||
if (name) {
|
||||
log?.info?.("CURSOR", `AgentService MCP tool_call ${name}`);
|
||||
finished = true;
|
||||
onEvent({
|
||||
type: "tool_call",
|
||||
value: {
|
||||
id: mcp.toolCallId || `call_${crypto.randomUUID()}`,
|
||||
name,
|
||||
arguments: JSON.stringify(mcp.args || {}),
|
||||
},
|
||||
});
|
||||
onEvent({ type: "done", finishReason: "tool_calls" });
|
||||
} else {
|
||||
debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`);
|
||||
finished = true;
|
||||
onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" });
|
||||
}
|
||||
} else {
|
||||
// Every other ExecServerMessage variant is an editor-backed tool
|
||||
// (shell, read, write, …) that 9router cannot service. Fail the
|
||||
// turn rather than narrating protocol state as assistant text.
|
||||
debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`);
|
||||
finished = true;
|
||||
onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" });
|
||||
// Auto/Composer often probe IDE builtins (shell/read/…). Reject
|
||||
// them so the model can continue with MCP tools or a text answer
|
||||
// instead of stalling the h2 stream.
|
||||
const rejection = rejectExecRequest(execRequest);
|
||||
if (rejection) {
|
||||
log?.info?.("CURSOR", `AgentService rejected IDE exec fields=${[...execRequest.keys()].join(",")}`);
|
||||
session.write(rejection);
|
||||
} else {
|
||||
debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`);
|
||||
finished = true;
|
||||
onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" });
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -586,7 +726,10 @@ export class CursorExecutor extends BaseExecutor {
|
||||
} finally {
|
||||
try { session.end(); } catch {}
|
||||
try { session.close(); } catch {}
|
||||
if (!finished) onEvent({ type: "done" });
|
||||
if (!finished) {
|
||||
flushThinkingFallback(onEvent);
|
||||
onEvent({ type: "done" });
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
@@ -594,10 +737,21 @@ export class CursorExecutor extends BaseExecutor {
|
||||
let content = "";
|
||||
let reasoning = "";
|
||||
let agentError = null;
|
||||
const toolCalls = [];
|
||||
let finishReason = "stop";
|
||||
await consume((event) => {
|
||||
if (event.type === "text") content += event.value;
|
||||
else if (event.type === "thinking") reasoning += event.value;
|
||||
else if (event.type === "tool_call") {
|
||||
toolCalls.push({
|
||||
id: event.value.id,
|
||||
type: "function",
|
||||
function: { name: event.value.name, arguments: event.value.arguments },
|
||||
});
|
||||
finishReason = "tool_calls";
|
||||
}
|
||||
else if (event.type === "error") agentError = event.value;
|
||||
else if (event.type === "done" && event.finishReason) finishReason = event.finishReason;
|
||||
});
|
||||
if (agentError) {
|
||||
return {
|
||||
@@ -611,13 +765,19 @@ export class CursorExecutor extends BaseExecutor {
|
||||
responseFormat: FORMATS.OPENAI,
|
||||
};
|
||||
}
|
||||
const message = {
|
||||
role: "assistant",
|
||||
content: content || null,
|
||||
...(reasoning ? { reasoning_content: reasoning } : {}),
|
||||
...(toolCalls.length ? { tool_calls: toolCalls } : {}),
|
||||
};
|
||||
return {
|
||||
response: new Response(JSON.stringify({
|
||||
id: responseId,
|
||||
object: "chat.completion",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, message: { role: "assistant", content: content || null, ...(reasoning ? { reasoning_content: reasoning } : {}) }, finish_reason: "stop" }],
|
||||
choices: [{ index: 0, message, finish_reason: finishReason }],
|
||||
usage: estimateUsage(body, content.length, FORMATS.OPENAI),
|
||||
}), { headers: { "Content-Type": "application/json" } }),
|
||||
url,
|
||||
@@ -635,6 +795,18 @@ export class CursorExecutor extends BaseExecutor {
|
||||
controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: { content: event.value } })));
|
||||
} else if (event.type === "thinking") {
|
||||
controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: { reasoning_content: event.value } })));
|
||||
} else if (event.type === "tool_call") {
|
||||
controller.enqueue(encoder.encode(chatChunkSse({
|
||||
id: responseId, created, model,
|
||||
delta: {
|
||||
tool_calls: [{
|
||||
index: 0,
|
||||
id: event.value.id,
|
||||
type: "function",
|
||||
function: { name: event.value.name, arguments: event.value.arguments },
|
||||
}],
|
||||
},
|
||||
})));
|
||||
} else if (event.type === "error") {
|
||||
// An SSE error frame, not a content delta: a protocol failure must not
|
||||
// be rendered to the user as the assistant's reply, and downstream
|
||||
@@ -643,7 +815,10 @@ export class CursorExecutor extends BaseExecutor {
|
||||
controller.enqueue(encoder.encode(SSE_DONE));
|
||||
controller.close();
|
||||
} else if (event.type === "done") {
|
||||
controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: {}, finishReason: "stop" })));
|
||||
controller.enqueue(encoder.encode(chatChunkSse({
|
||||
id: responseId, created, model, delta: {},
|
||||
finishReason: event.finishReason || "stop",
|
||||
})));
|
||||
controller.enqueue(encoder.encode(SSE_DONE));
|
||||
controller.close();
|
||||
}
|
||||
@@ -664,9 +839,9 @@ export class CursorExecutor extends BaseExecutor {
|
||||
}
|
||||
|
||||
async execute({ model, body, stream, credentials, signal, log, proxyOptions = null }) {
|
||||
if (isAgentTextRequest(body)) {
|
||||
if (isAgentCapableRequest(body)) {
|
||||
try {
|
||||
return await this.executeAgent({ model, body, stream, credentials, signal });
|
||||
return await this.executeAgent({ model, body, stream, credentials, signal, log });
|
||||
} catch (error) {
|
||||
return {
|
||||
response: new Response(JSON.stringify({
|
||||
|
||||
@@ -11,6 +11,7 @@ import { CursorExecutor } from "./cursor.js";
|
||||
import { VertexExecutor } from "./vertex.js";
|
||||
import { OpenCodeExecutor } from "./opencode.js";
|
||||
import { OpenCodeGoExecutor } from "./opencode-go.js";
|
||||
import { OpenCodeZenExecutor } from "./opencode-zen.js";
|
||||
import { GrokWebExecutor } from "./grok-web.js";
|
||||
import { GrokCliExecutor } from "./grok-cli.js";
|
||||
import { PerplexityWebExecutor } from "./perplexity-web.js";
|
||||
@@ -34,6 +35,7 @@ const executors = {
|
||||
github: new GithubExecutor(),
|
||||
iflow: new IFlowExecutor(),
|
||||
qoder: new QoderExecutor(),
|
||||
"qoder-cn": new QoderExecutor("qoder-cn"),
|
||||
kiro: new KiroExecutor(),
|
||||
kimchi: new KimchiExecutor(),
|
||||
codex: new CodexExecutor(),
|
||||
@@ -43,6 +45,7 @@ const executors = {
|
||||
"vertex-partner": new VertexExecutor("vertex-partner"),
|
||||
opencode: new OpenCodeExecutor(),
|
||||
"opencode-go": new OpenCodeGoExecutor(),
|
||||
"opencode-zen": new OpenCodeZenExecutor(),
|
||||
"grok-web": new GrokWebExecutor(),
|
||||
"grok-cli": new GrokCliExecutor(),
|
||||
gcli: new GrokCliExecutor(), // Alias
|
||||
@@ -89,6 +92,7 @@ export { VertexExecutor } from "./vertex.js";
|
||||
export { DefaultExecutor } from "./default.js";
|
||||
export { OpenCodeExecutor } from "./opencode.js";
|
||||
export { OpenCodeGoExecutor } from "./opencode-go.js";
|
||||
export { OpenCodeZenExecutor } from "./opencode-zen.js";
|
||||
export { GrokWebExecutor } from "./grok-web.js";
|
||||
export { GrokCliExecutor } from "./grok-cli.js";
|
||||
export { PerplexityWebExecutor } from "./perplexity-web.js";
|
||||
|
||||
315
open-sse/executors/opencode-zen.js
Normal file
315
open-sse/executors/opencode-zen.js
Normal file
@@ -0,0 +1,315 @@
|
||||
import crypto from "node:crypto";
|
||||
import { DefaultExecutor } from "./default.js";
|
||||
import { resolveSessionId } from "../utils/sessionManager.js";
|
||||
import { isMuseSparkModel } from "../providers/models/helpers.js";
|
||||
import {
|
||||
normalizeResponsesInput,
|
||||
clampResponsesCallId,
|
||||
coerceResponsesArguments,
|
||||
coerceResponsesOutput,
|
||||
} from "../translator/formats/responsesApi.js";
|
||||
|
||||
const SESSION_HEADER = "x-opencode-session";
|
||||
const SESSION_FIELD = "_opencodeZenSession";
|
||||
const MAX_SESSION_LENGTH = 256;
|
||||
|
||||
const RESPONSES_BASE_URL = "https://opencode.ai/zen/v1/responses";
|
||||
const MAX_TOOL_NAME_LEN = 128;
|
||||
const OPENCODE_UA = "opencode/1.18.31";
|
||||
export const OPENCODE_SESSION_RE = /^ses_[0-9a-f]{12}[0-9A-Za-z]{14}$/;
|
||||
const BASE62_CHARS = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz";
|
||||
// Free-tier fingerprint (mirrors opencode executor, PR #4132): upstream 403s
|
||||
// requests without the file-search quartet and without stream:true.
|
||||
const OPENCODE_FINGERPRINT_TOOLS = ["bash", "glob", "grep", "read"];
|
||||
|
||||
function hasValidOpencodeVersion(ua) {
|
||||
const m = String(ua || "").match(/opencode\/(\d+)\.(\d+)(?:\.(\d+))?/i);
|
||||
if (!m) return false;
|
||||
const major = parseInt(m[1], 10);
|
||||
const minor = parseInt(m[2], 10);
|
||||
return major > 1 || (major === 1 && minor >= 17);
|
||||
}
|
||||
|
||||
function unstableRandom() {
|
||||
const bytes = crypto.randomBytes(14);
|
||||
let randomPart = "";
|
||||
for (let i = 0; i < 14; i++) {
|
||||
randomPart += BASE62_CHARS[bytes[i] % 62];
|
||||
}
|
||||
return randomPart;
|
||||
}
|
||||
|
||||
export function generateSessionId(timestamp = Date.now()) {
|
||||
const current = BigInt(timestamp) * 0x1000n + 1n;
|
||||
const value = ~current;
|
||||
const time = Array.from({ length: 6 }, (_, index) =>
|
||||
Number((value >> BigInt(40 - 8 * index)) & 0xffn)
|
||||
.toString(16)
|
||||
.padStart(2, "0")
|
||||
).join("");
|
||||
return `ses_${time}${unstableRandom()}`;
|
||||
}
|
||||
|
||||
export function generateRequestId(timestamp = Date.now()) {
|
||||
const current = BigInt(timestamp) * 0x1000n + 1n;
|
||||
const value = current;
|
||||
const time = Array.from({ length: 6 }, (_, index) =>
|
||||
Number((value >> BigInt(40 - 8 * index)) & 0xffn)
|
||||
.toString(16)
|
||||
.padStart(2, "0")
|
||||
).join("");
|
||||
return `msg_${time}${unstableRandom()}`;
|
||||
}
|
||||
|
||||
export function translateSessionId(sessionId, clientTool = "") {
|
||||
if (typeof sessionId === "string" && OPENCODE_SESSION_RE.test(sessionId.trim())) {
|
||||
return sessionId.trim();
|
||||
}
|
||||
const digest = crypto
|
||||
.createHash("sha256")
|
||||
.update(`opencode\0${clientTool || "generic"}\0${sessionId || ""}`)
|
||||
.digest();
|
||||
const timeHex = digest.subarray(0, 6).toString("hex");
|
||||
let randomPart = "";
|
||||
for (let i = 6; i < 20; i++) {
|
||||
randomPart += BASE62_CHARS[digest[i] % 62];
|
||||
}
|
||||
return `ses_${timeHex}${randomPart}`;
|
||||
}
|
||||
|
||||
function toolNameOf(tool) {
|
||||
if (!tool || typeof tool !== "object" || Array.isArray(tool)) return "";
|
||||
const fn = tool.function && typeof tool.function === "object" && !Array.isArray(tool.function) ? tool.function : null;
|
||||
const raw = typeof tool.name === "string" ? tool.name : (typeof fn?.name === "string" ? fn.name : "");
|
||||
return raw.trim();
|
||||
}
|
||||
|
||||
function ensureChatFingerprintTools(body) {
|
||||
if (!body || typeof body !== "object") return;
|
||||
const present = new Set();
|
||||
if (Array.isArray(body.tools)) {
|
||||
for (const tool of body.tools) {
|
||||
const name = toolNameOf(tool);
|
||||
if (name) present.add(name);
|
||||
}
|
||||
} else {
|
||||
body.tools = [];
|
||||
}
|
||||
for (const name of OPENCODE_FINGERPRINT_TOOLS) {
|
||||
if (present.has(name)) continue;
|
||||
body.tools.push({
|
||||
type: "function",
|
||||
function: {
|
||||
name,
|
||||
description: `OpenCode built-in ${name} tool`,
|
||||
parameters: { type: "object", properties: {} },
|
||||
},
|
||||
});
|
||||
present.add(name);
|
||||
}
|
||||
}
|
||||
|
||||
function ensureResponsesFingerprintTools(body) {
|
||||
if (!body || typeof body !== "object") return;
|
||||
const present = new Set();
|
||||
if (Array.isArray(body.tools)) {
|
||||
for (const tool of body.tools) {
|
||||
const name = toolNameOf(tool);
|
||||
if (name) present.add(name);
|
||||
}
|
||||
} else {
|
||||
body.tools = [];
|
||||
}
|
||||
for (const name of OPENCODE_FINGERPRINT_TOOLS) {
|
||||
if (present.has(name)) continue;
|
||||
body.tools.push({
|
||||
type: "function",
|
||||
name,
|
||||
description: `OpenCode built-in ${name} tool`,
|
||||
parameters: { type: "object", properties: {} },
|
||||
});
|
||||
present.add(name);
|
||||
}
|
||||
}
|
||||
|
||||
function normalizeSession(value) {
|
||||
if (typeof value !== "string") return null;
|
||||
const normalized = value.trim();
|
||||
if (!normalized || normalized.length > MAX_SESSION_LENGTH) return null;
|
||||
return normalized;
|
||||
}
|
||||
|
||||
function nativeSession(headers) {
|
||||
if (!headers || typeof headers !== "object") return null;
|
||||
for (const [key, value] of Object.entries(headers)) {
|
||||
if (key.toLowerCase() === SESSION_HEADER) {
|
||||
const normalized = normalizeSession(value);
|
||||
if (normalized && OPENCODE_SESSION_RE.test(normalized)) return normalized;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function translatedSession(sessionId, clientTool) {
|
||||
return translateSessionId(sessionId, clientTool);
|
||||
}
|
||||
|
||||
// Strip the thinking suffix "model(level)" so checks hit the base id.
|
||||
function baseModelId(model) {
|
||||
return String(model || "").replace(/\([^()]+\)\s*$/, "").trim();
|
||||
}
|
||||
|
||||
function isResponsesModel(model) {
|
||||
return isMuseSparkModel(baseModelId(model));
|
||||
}
|
||||
|
||||
// Flatten Chat Completions tool declarations into the Responses flat shape and
|
||||
// drop hosted/nameless tools the /responses endpoint rejects.
|
||||
function normalizeResponsesTools(body) {
|
||||
if (!Array.isArray(body.tools)) return;
|
||||
const validNames = new Set();
|
||||
body.tools = body.tools.filter((tool) => {
|
||||
if (!tool || typeof tool !== "object" || Array.isArray(tool)) return false;
|
||||
const fn = tool.function && typeof tool.function === "object" && !Array.isArray(tool.function) ? tool.function : null;
|
||||
const rawName = typeof tool.name === "string" ? tool.name : (typeof fn?.name === "string" ? fn.name : "");
|
||||
const name = rawName.trim();
|
||||
if (!name) return false;
|
||||
const description = typeof tool.description === "string" ? tool.description : (typeof fn?.description === "string" ? fn.description : "");
|
||||
let parameters = (tool.parameters && typeof tool.parameters === "object" && !Array.isArray(tool.parameters))
|
||||
? tool.parameters
|
||||
: (fn?.parameters && typeof fn.parameters === "object" && !Array.isArray(fn.parameters) ? fn.parameters : { type: "object", properties: {} });
|
||||
// Mirror the request translator: {type:"object"} without properties is rejected
|
||||
// by strict Responses backends, so fill in the empty properties map.
|
||||
if (parameters.type === "object" && !parameters.properties) parameters = { ...parameters, properties: {} };
|
||||
for (const k of Object.keys(tool)) delete tool[k];
|
||||
tool.type = "function";
|
||||
tool.name = name.slice(0, MAX_TOOL_NAME_LEN);
|
||||
if (description) tool.description = description;
|
||||
tool.parameters = parameters;
|
||||
validNames.add(tool.name);
|
||||
return true;
|
||||
});
|
||||
if (body.tool_choice && typeof body.tool_choice === "object" && !Array.isArray(body.tool_choice)) {
|
||||
if (body.tool_choice.type === "function") {
|
||||
const n = typeof body.tool_choice.name === "string" ? body.tool_choice.name.trim() : "";
|
||||
if (!n || !validNames.has(n)) delete body.tool_choice;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Last line of defense for native Responses clients (sourceFormat === targetFormat
|
||||
// skips translation): coerce items in place so malformed tool payloads 400 here
|
||||
// with a clear shape instead of upstream as InputValidationError.
|
||||
function sanitizeResponsesItems(body) {
|
||||
if (!Array.isArray(body.input)) return;
|
||||
body.input = body.input.filter((item) => {
|
||||
if (!item || typeof item !== "object" || Array.isArray(item)) return true;
|
||||
// Strip prior-turn reasoning items: Muse Spark contributor models route to
|
||||
// an upstream Console backend where encrypted_content cannot be validated across
|
||||
// rotated accounts or sessions, causing 400 "reasoning encrypted_content was not issued to this caller".
|
||||
if (item.type === "reasoning") return false;
|
||||
delete item.encrypted_content;
|
||||
delete item.reasoning_encrypted_content;
|
||||
if (item.type === "function_call") {
|
||||
if (!item.name || typeof item.name !== "string" || item.name.trim() === "") return false;
|
||||
item.name = item.name.trim().slice(0, MAX_TOOL_NAME_LEN);
|
||||
item.call_id = clampResponsesCallId(item.call_id);
|
||||
item.arguments = coerceResponsesArguments(item.arguments);
|
||||
return true;
|
||||
}
|
||||
if (item.type === "function_call_output") {
|
||||
item.call_id = clampResponsesCallId(item.call_id);
|
||||
item.output = coerceResponsesOutput(item.output);
|
||||
return true;
|
||||
}
|
||||
return true;
|
||||
});
|
||||
}
|
||||
|
||||
export class OpenCodeZenExecutor extends DefaultExecutor {
|
||||
constructor() {
|
||||
super("opencode-zen");
|
||||
}
|
||||
|
||||
buildUrl(model, stream, urlIndex = 0, credentials = null) {
|
||||
// Muse Spark lives on /responses even when a stale runtimeTransport leaks in.
|
||||
if (isResponsesModel(model)) return RESPONSES_BASE_URL;
|
||||
return super.buildUrl(model, stream, urlIndex, credentials);
|
||||
}
|
||||
|
||||
prepareRequestCredentials({ body, credentials, providerSessionId, clientTool } = {}) {
|
||||
const sourceCredentials = credentials || {};
|
||||
const native = nativeSession(sourceCredentials.rawHeaders);
|
||||
const resolved = normalizeSession(providerSessionId) || resolveSessionId({
|
||||
headers: sourceCredentials.rawHeaders,
|
||||
body,
|
||||
connectionId: sourceCredentials.connectionId,
|
||||
scope: "opencode-zen",
|
||||
});
|
||||
|
||||
return {
|
||||
...sourceCredentials,
|
||||
[SESSION_FIELD]: native || translatedSession(resolved, clientTool),
|
||||
};
|
||||
}
|
||||
|
||||
async execute(args) {
|
||||
const credentials = this.prepareRequestCredentials(args);
|
||||
return super.execute({ ...args, credentials });
|
||||
}
|
||||
|
||||
buildHeaders(credentials, stream = true, url, model) {
|
||||
const headers = super.buildHeaders(credentials || {}, stream, url, model);
|
||||
const raw = credentials?.rawHeaders || {};
|
||||
const lower = {};
|
||||
for (const [k, v] of Object.entries(raw)) lower[k.toLowerCase()] = v;
|
||||
const downstreamUa = lower["user-agent"] || "";
|
||||
// Free-tier gate: spoof the official client UA.
|
||||
headers["User-Agent"] = hasValidOpencodeVersion(downstreamUa) ? downstreamUa : OPENCODE_UA;
|
||||
headers["x-opencode-client"] = lower["x-opencode-client"] || "desktop";
|
||||
const prepared = credentials?.[SESSION_FIELD];
|
||||
if (prepared) {
|
||||
headers[SESSION_HEADER] = prepared;
|
||||
return headers;
|
||||
}
|
||||
|
||||
const fallback = this.prepareRequestCredentials({ credentials });
|
||||
headers[SESSION_HEADER] = fallback[SESSION_FIELD];
|
||||
return headers;
|
||||
}
|
||||
|
||||
transformRequest(model, body, stream, credentials) {
|
||||
const out = super.transformRequest(model, body);
|
||||
// Free-tier gate: upstream 403s stream:false even when everything else is valid.
|
||||
if (out && typeof out === "object") out.stream = true;
|
||||
if (!isResponsesModel(model || body?.model)) {
|
||||
ensureChatFingerprintTools(out);
|
||||
return out;
|
||||
}
|
||||
const normalized = normalizeResponsesInput(out.input);
|
||||
if (normalized) out.input = normalized;
|
||||
if (!Array.isArray(out.input) || out.input.length === 0) {
|
||||
out.input = [{ type: "message", role: "user", content: [{ type: "input_text", text: "..." }] }];
|
||||
}
|
||||
// Responses names the output cap max_output_tokens, not max_tokens.
|
||||
if (out.max_output_tokens === undefined) {
|
||||
if (out.max_completion_tokens !== undefined) out.max_output_tokens = out.max_completion_tokens;
|
||||
else if (out.max_tokens !== undefined) out.max_output_tokens = out.max_tokens;
|
||||
}
|
||||
delete out.max_tokens;
|
||||
delete out.max_completion_tokens;
|
||||
if (out.reasoning_effort !== undefined && out.reasoning === undefined) {
|
||||
out.reasoning = { effort: out.reasoning_effort, summary: "auto" };
|
||||
}
|
||||
if (out.reasoning && typeof out.reasoning === "object" && !Array.isArray(out.reasoning)) {
|
||||
if (!out.reasoning.summary) out.reasoning.summary = "auto";
|
||||
}
|
||||
delete out.reasoning_effort;
|
||||
out.stream = true;
|
||||
out.store = false;
|
||||
ensureResponsesFingerprintTools(out);
|
||||
normalizeResponsesTools(out);
|
||||
sanitizeResponsesItems(out);
|
||||
return out;
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,7 @@ import { getThinkingLevels } from "../providers/thinkingLevels.js";
|
||||
import { injectReasoningContent } from "../utils/reasoningContentInjector.js";
|
||||
import { resolveSessionId } from "../utils/sessionManager.js";
|
||||
import { isMuseSparkModel } from "../providers/models/helpers.js";
|
||||
import { applyFingerprintTools } from "../utils/opencodeFingerprint.js";
|
||||
import { ANTHROPIC_API_VERSION } from "../providers/shared.js";
|
||||
import {
|
||||
normalizeResponsesInput,
|
||||
@@ -24,68 +25,6 @@ export const OPENCODE_SESSION_RE = /^ses_[0-9a-f]{12}[0-9A-Za-z]{14}$/;
|
||||
export const OPENCODE_REQUEST_RE = /^msg_[0-9a-f]{12}[0-9A-Za-z]{14}$/;
|
||||
const BASE62_CHARS = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz";
|
||||
|
||||
// OpenCode free tier requires both 'bash' and 'read' in tools payload.
|
||||
// Injected as cloaked decoy tools so external CLI tools (e.g. Claude Code's Bash/Read)
|
||||
// take precedence while satisfying upstream verification.
|
||||
const OPENCODE_DECOY_CHAT_TOOLS = [
|
||||
{
|
||||
type: "function",
|
||||
function: {
|
||||
name: "bash",
|
||||
description: "This tool is currently unavailable and must not be used.",
|
||||
parameters: { type: "object", properties: {} },
|
||||
},
|
||||
},
|
||||
{
|
||||
type: "function",
|
||||
function: {
|
||||
name: "read",
|
||||
description: "This tool is currently unavailable and must not be used.",
|
||||
parameters: { type: "object", properties: {} },
|
||||
},
|
||||
},
|
||||
];
|
||||
|
||||
const OPENCODE_DECOY_RESPONSES_TOOLS = [
|
||||
{
|
||||
type: "function",
|
||||
name: "bash",
|
||||
description: "This tool is currently unavailable and must not be used.",
|
||||
parameters: { type: "object", properties: {} },
|
||||
},
|
||||
{
|
||||
type: "function",
|
||||
name: "read",
|
||||
description: "This tool is currently unavailable and must not be used.",
|
||||
parameters: { type: "object", properties: {} },
|
||||
},
|
||||
];
|
||||
|
||||
function cloakOpencodeTools(body, isResponses) {
|
||||
if (!body || typeof body !== "object") return;
|
||||
if (isResponses) {
|
||||
if (!Array.isArray(body.tools)) body.tools = [];
|
||||
const names = new Set(body.tools.map((t) => t.name || t.function?.name));
|
||||
for (const tool of OPENCODE_DECOY_RESPONSES_TOOLS) {
|
||||
if (!names.has(tool.name)) body.tools.push({ ...tool });
|
||||
}
|
||||
if (!body.tool_choice) body.tool_choice = "auto";
|
||||
} else {
|
||||
const hasTools = Array.isArray(body.tools) && body.tools.length > 0;
|
||||
if (!hasTools) {
|
||||
body.tools = OPENCODE_DECOY_CHAT_TOOLS.map((t) => ({ ...t, function: { ...t.function } }));
|
||||
if (!body.tool_choice) body.tool_choice = "none";
|
||||
} else {
|
||||
const names = new Set(body.tools.map((t) => t.function?.name || t.name));
|
||||
for (const tool of OPENCODE_DECOY_CHAT_TOOLS) {
|
||||
if (!names.has(tool.function.name)) {
|
||||
body.tools.push({ ...tool, function: { ...tool.function } });
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function hasValidOpencodeVersion(ua) {
|
||||
const m = String(ua || "").match(/opencode\/(\d+)\.(\d+)(?:\.(\d+))?/i);
|
||||
if (!m) return false;
|
||||
@@ -499,11 +438,12 @@ export class OpenCodeExecutor extends BaseExecutor {
|
||||
body.store = false;
|
||||
normalizeResponsesTools(body);
|
||||
sanitizeResponsesItems(body);
|
||||
if (!Array.isArray(body.tools) || body.tools.length === 0) {
|
||||
cloakOpencodeTools(body, true);
|
||||
}
|
||||
// Free-tier fingerprint tools are required even when an agent client
|
||||
// already supplied tools. ZCode/Claude Code requests normally have
|
||||
// non-empty tool arrays; skipping cloaking here triggers 403 FreeTierError.
|
||||
applyFingerprintTools(body, true);
|
||||
} else if (body && typeof body === "object") {
|
||||
cloakOpencodeTools(body, false);
|
||||
applyFingerprintTools(body, false);
|
||||
}
|
||||
return injectReasoningContent({ provider: this.provider, model, body });
|
||||
}
|
||||
|
||||
@@ -29,7 +29,7 @@ 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 { FETCH_CONNECT_TIMEOUT_MS, HTTP_STATUS } from "../config/runtimeConfig.js";
|
||||
import { resolveProviderTimeoutMs } from "../services/providerTimeout.js";
|
||||
import {
|
||||
QODER_CHAT_SIG_PATH,
|
||||
@@ -208,16 +208,16 @@ function truncate(s, n) {
|
||||
/**
|
||||
* Map the OpenAI-style request body into the exact shape Qoder expects.
|
||||
*/
|
||||
async function buildQoderRequestBody({ model, body, credentials, log, proxyOptions, signal, uploadFn = null }) {
|
||||
async function buildQoderRequestBody({ model, body, credentials, log, proxyOptions, signal, uploadFn = null, region = "intl" }) {
|
||||
const qoderKey = String(model || "").replace(/^qoder\//, "");
|
||||
|
||||
|
||||
// Fetch model config from dynamic API instead of relying on static QODER_MODEL_MAP.
|
||||
// This allows support for new Qoder models (e.g., qmodel_latest) without code changes.
|
||||
let modelConfig = await getQoderModelConfig(credentials, qoderKey, { log, proxyOptions, signal });
|
||||
let modelConfig = await getQoderModelConfig(credentials, qoderKey, { log, proxyOptions, signal, region });
|
||||
if (!modelConfig) {
|
||||
// Try a forced refresh once before giving up — the cache may simply
|
||||
// not be populated yet on first ever call for this credential.
|
||||
const refreshed = await resolveQoderModels(credentials, { forceRefresh: true, log, proxyOptions, signal });
|
||||
const refreshed = await resolveQoderModels(credentials, { forceRefresh: true, log, proxyOptions, signal, region });
|
||||
const retried = refreshed?.rawConfigs.get(qoderKey);
|
||||
if (!retried) {
|
||||
throw new Error(
|
||||
@@ -337,47 +337,65 @@ async function buildQoderRequestBody({ model, body, credentials, log, proxyOptio
|
||||
|
||||
/**
|
||||
* Check if a qoder error message indicates a billing/quota block.
|
||||
* Signatures: code 112 (quota exhausted), code 10605 (queue throttle), pricingUrl field.
|
||||
* Signatures: code 110 (billing daily count exceeded), code 112 (quota
|
||||
* exhausted), code 10605 (queue throttle), pricingUrl field.
|
||||
*/
|
||||
function isBillingBlock(inner) {
|
||||
if (!inner || typeof inner !== "string") return false;
|
||||
const lowerMsg = inner.toLowerCase();
|
||||
// Match: {"code":"112",...}, {"code":"10605",...}, or pricingUrl field
|
||||
return /\"code\"\s*:\s*\"(112|10605)\"/.test(inner) || lowerMsg.includes("pricingurl");
|
||||
if (lowerMsg.includes("pricingurl")) return true;
|
||||
// Parsed code preferred over regex: matches numeric or string "110"/"112"/"10605".
|
||||
try {
|
||||
const parsed = JSON.parse(inner);
|
||||
const code = String(parsed?.code ?? "");
|
||||
if (code === "110" || code === "112" || code === "10605") return true;
|
||||
} catch { /* not JSON — fall through to legacy shape match */ }
|
||||
// Match legacy exact shapes: {"code":"112",...}, {"code":"10605",...}.
|
||||
return /"code"\s*:\s*"(112|10605)"/.test(inner);
|
||||
}
|
||||
|
||||
/**
|
||||
* Peek the first SSE frame to detect billing errors before piping.
|
||||
* Returns { isBilling, statusVal, message, consumed } — `consumed` is every
|
||||
* Peek the first SSE data line to detect upstream errors before piping.
|
||||
* Returns { isError, isBilling, statusVal, message, consumed } — `consumed` is every
|
||||
* byte read so far (including the peeked line) so the caller can re-process
|
||||
* it and nothing is dropped from the stream.
|
||||
*/
|
||||
async function peekFirstQoderFrame(reader, decoder) {
|
||||
let consumed = "";
|
||||
let offset = 0;
|
||||
let upstreamDone = false;
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) return { isBilling: false, consumed, upstreamDone: true };
|
||||
let nl = consumed.indexOf("\n", offset);
|
||||
if (nl === -1 && !upstreamDone) {
|
||||
const { done, value } = await reader.read();
|
||||
upstreamDone = done;
|
||||
consumed += done ? decoder.decode() : decoder.decode(value, { stream: true });
|
||||
continue;
|
||||
}
|
||||
if (offset >= consumed.length) return { isError: false, consumed, upstreamDone };
|
||||
if (nl === -1) nl = consumed.length;
|
||||
|
||||
consumed += decoder.decode(value, { stream: true });
|
||||
const nl = consumed.indexOf("\n");
|
||||
if (nl === -1) continue; // need a full line first
|
||||
|
||||
const line = consumed.slice(0, nl).replace(/\r$/, "").trim();
|
||||
const line = consumed.slice(offset, nl).replace(/\r$/, "").trim();
|
||||
offset = nl + 1;
|
||||
if (!line.startsWith("data:")) continue;
|
||||
|
||||
const data = line.slice(5).trimStart();
|
||||
if (data === "[DONE]") return { isBilling: false, consumed };
|
||||
if (data === "[DONE]") return { isError: false, consumed, upstreamDone };
|
||||
|
||||
let envelope;
|
||||
try { envelope = JSON.parse(data); } catch { return { isBilling: false, consumed }; }
|
||||
try { envelope = JSON.parse(data); } catch { return { isError: false, consumed, upstreamDone }; }
|
||||
|
||||
const statusVal = typeof envelope.statusCodeValue === "number" ? envelope.statusCodeValue : 200;
|
||||
const inner = typeof envelope.body === "string" ? envelope.body : "";
|
||||
// statusCodeValue is documented numeric, but accept numeric strings defensively.
|
||||
const raw = Number(envelope?.statusCodeValue);
|
||||
const statusVal = Number.isNaN(raw) ? 200 : raw;
|
||||
const inner = typeof envelope?.body === "string"
|
||||
? envelope.body
|
||||
: envelope?.body != null ? JSON.stringify(envelope.body) : "";
|
||||
|
||||
if (statusVal !== 200 && isBillingBlock(inner)) {
|
||||
return { isBilling: true, statusVal, message: inner || `qoder billing block (${statusVal})` };
|
||||
if (statusVal !== 200) {
|
||||
return { isError: true, isBilling: isBillingBlock(inner), statusVal, message: inner || `upstream status ${statusVal}` };
|
||||
}
|
||||
return { isBilling: false, consumed };
|
||||
return { isError: false, consumed, upstreamDone };
|
||||
}
|
||||
}
|
||||
|
||||
@@ -388,8 +406,8 @@ async function peekFirstQoderFrame(reader, decoder) {
|
||||
* Each upstream line looks like:
|
||||
* data: {"statusCodeValue":200,"body":"{\"choices\":[{\"delta\":{...}}]}"}
|
||||
* The inner body is an OpenAI streaming chunk (or "[DONE]"). We unwrap it
|
||||
* and re-emit as `data: <inner>\n\n`. Errors become a synthetic OpenAI error
|
||||
* chunk + [DONE].
|
||||
* and re-emit as `data: <inner>\n\n`. First-frame errors become HTTP errors;
|
||||
* errors after streaming starts retain the synthetic chunk + [DONE] path.
|
||||
*
|
||||
* Critical: Qoder's SSE often keeps the socket open after the terminal
|
||||
* [DONE]/error frame (agent keepalive). Non-streaming clients drain via
|
||||
@@ -401,24 +419,28 @@ async function peekFirstQoderFrame(reader, decoder) {
|
||||
* usage from the finish chunk, so we coalesce those two frames (see
|
||||
* createQoderSseCoalescer) before forwarding.
|
||||
*
|
||||
* NEW: Peek first frame to detect billing blocks (code 112/10605/pricingUrl).
|
||||
* If detected, return 403 response so chatCore marks connection unavailable
|
||||
* and triggers combo fallback instead of leaking error text into chat.
|
||||
* Peek the first frame for errors before committing to HTTP 200. Preserve
|
||||
* upstream error statuses so chatCore can handle failures instead of recording
|
||||
* error text as a successful completion. Billing blocks retain the existing
|
||||
* 403 mapping for quota/account fallback.
|
||||
*/
|
||||
async function wrapQoderSSE(response, model) {
|
||||
async function wrapQoderSSE(response, model, log = null) {
|
||||
if (!response.ok || !response.body) return response;
|
||||
|
||||
const decoder = new TextDecoder();
|
||||
const reader = response.body.getReader();
|
||||
|
||||
// Peek first frame to detect billing block
|
||||
// Detect errors before returning a successful streaming response.
|
||||
const peek = await peekFirstQoderFrame(reader, decoder);
|
||||
if (peek?.isBilling) {
|
||||
// Billing block detected — return 403 so chatCore fails this connection
|
||||
if (peek.isError) {
|
||||
await reader.cancel().catch(() => {});
|
||||
const status = peek.isBilling
|
||||
? HTTP_STATUS.FORBIDDEN
|
||||
: Number.isInteger(peek.statusVal) && peek.statusVal >= HTTP_STATUS.BAD_REQUEST && peek.statusVal <= 599
|
||||
? peek.statusVal : HTTP_STATUS.BAD_GATEWAY;
|
||||
return new Response(
|
||||
JSON.stringify({ error: { message: peek.message, code: peek.statusVal } }),
|
||||
{ status: 403, headers: { "Content-Type": "application/json" } }
|
||||
{ status, headers: { "Content-Type": "application/json" } }
|
||||
);
|
||||
}
|
||||
|
||||
@@ -449,11 +471,35 @@ async function wrapQoderSSE(response, model) {
|
||||
|
||||
let envelope;
|
||||
try { envelope = JSON.parse(data); } catch { return; }
|
||||
const statusVal = typeof envelope.statusCodeValue === "number" ? envelope.statusCodeValue : 200;
|
||||
const statusVal = Number(envelope.statusCodeValue) || 200;
|
||||
const inner = typeof envelope.body === "string"
|
||||
? envelope.body
|
||||
: envelope.body != null ? JSON.stringify(envelope.body) : "";
|
||||
if (statusVal !== 200) {
|
||||
// Always visible: error envelopes are rare and worth one stderr line at
|
||||
// any log level (response bodies carry no credentials).
|
||||
try {
|
||||
console.error(`[QODER] error envelope status=${statusVal} statusType=${typeof envelope.statusCodeValue} bodyType=${typeof envelope.body} body=${truncate(inner, 300)}`);
|
||||
} catch { /* logging must not break the stream */ }
|
||||
if (isBillingBlock(inner)) {
|
||||
// Billing/quota envelope at any stream position (peek only covers the
|
||||
// first frame): emit a structured error chunk, not fake assistant text.
|
||||
// parseSSEToOpenAIResponse understands chunk.error and turns it into a
|
||||
// non-200 result so chat.js locks the model and falls back. Streaming
|
||||
// clients receive a real SSE error instead of "[qoder error ...]" text.
|
||||
const errObj = JSON.stringify({
|
||||
error: {
|
||||
message: inner || `qoder billing block (${statusVal})`,
|
||||
code: "qoder_billing_block",
|
||||
status: 403,
|
||||
type: "quota_error",
|
||||
},
|
||||
});
|
||||
controller.enqueue(encoder.encode(`data: ${errObj}\n\n`));
|
||||
controller.enqueue(encoder.encode(SSE_DONE));
|
||||
doneEmitted = true;
|
||||
return;
|
||||
}
|
||||
const msg = inner || `upstream status ${statusVal}`;
|
||||
const errChunk = JSON.stringify({
|
||||
id: `qoder-error-${Date.now()}`,
|
||||
@@ -552,12 +598,13 @@ async function wrapQoderSSE(response, model) {
|
||||
}
|
||||
|
||||
export class QoderExecutor extends BaseExecutor {
|
||||
constructor() {
|
||||
super("qoder", PROVIDERS.qoder);
|
||||
constructor(provider = "qoder") {
|
||||
super(provider, PROVIDERS[provider]);
|
||||
this.region = provider === "qoder-cn" ? "cn" : "intl";
|
||||
}
|
||||
|
||||
buildUrl(credentials) {
|
||||
return `${qoderInferenceBase(credentials)}/algo${QODER_CHAT_SIG_PATH}?FetchKeys=llm_model_result&AgentId=agent_common&Encode=1`;
|
||||
return `${qoderInferenceBase(credentials, this.region)}/algo${QODER_CHAT_SIG_PATH}?FetchKeys=llm_model_result&AgentId=agent_common&Encode=1`;
|
||||
}
|
||||
|
||||
// Override execute entirely — Qoder needs:
|
||||
@@ -572,7 +619,7 @@ export class QoderExecutor extends BaseExecutor {
|
||||
const rawToken = credentials?.apiKey || credentials?.accessToken;
|
||||
if (isQoderPat(rawToken)) {
|
||||
try {
|
||||
credentials = await resolveQoderCredentials(credentials, proxyOptions, signal);
|
||||
credentials = await resolveQoderCredentials(credentials, proxyOptions, signal, this.region);
|
||||
} catch (err) {
|
||||
log?.error?.("QODER", `PAT exchange failed: ${err.message}`);
|
||||
const fakeResp = new Response(
|
||||
@@ -607,7 +654,7 @@ export class QoderExecutor extends BaseExecutor {
|
||||
let qoderKey;
|
||||
let payload;
|
||||
try {
|
||||
({ qoderKey, payload } = await buildQoderRequestBody({ model, body, credentials, log, proxyOptions, signal }));
|
||||
({ qoderKey, payload } = await buildQoderRequestBody({ model, body, credentials, log, proxyOptions, signal, region: this.region }));
|
||||
} catch (err) {
|
||||
const fakeResp = new Response(
|
||||
JSON.stringify({ error: { message: err.message } }),
|
||||
@@ -666,8 +713,15 @@ export class QoderExecutor extends BaseExecutor {
|
||||
response = await proxyAwareFetch(
|
||||
url,
|
||||
{ method: "POST", headers, body: encodedBodyBuf, signal: mergedSignal },
|
||||
proxyOptions,
|
||||
// A failed proxy request may already have reached Qoder. Replaying
|
||||
// the same COSY signature directly reuses its requestId and returns
|
||||
// 403/code 103. Let the caller retry through execute() with fresh signing.
|
||||
{ ...proxyOptions, strictProxy: true },
|
||||
);
|
||||
} catch (err) {
|
||||
// strictProxy wraps transport errors; retain caller cancellation semantics.
|
||||
if (mergedSignal.aborted) throw mergedSignal.reason;
|
||||
throw err;
|
||||
} finally {
|
||||
clearTimeout(connectTimer);
|
||||
}
|
||||
@@ -677,7 +731,7 @@ export class QoderExecutor extends BaseExecutor {
|
||||
return { response, url, headers, transformedBody: payload };
|
||||
}
|
||||
|
||||
const wrapped = await wrapQoderSSE(response, `qoder/${qoderKey}`);
|
||||
const wrapped = await wrapQoderSSE(response, `${this.provider}/${qoderKey}`, log);
|
||||
return { response: wrapped, url, headers, transformedBody: payload };
|
||||
}
|
||||
|
||||
|
||||
@@ -1,10 +1,15 @@
|
||||
import { DefaultExecutor } from "./default.js";
|
||||
import { getMimoAccountCookie, invalidateMimoAccountCookieCache, MIMO_API_BASE, MIMO_API_UA } from "../shared/mimoAccount.js";
|
||||
import { getMimoAccountCookie, invalidateMimoAccountCookieCache, resolveMimoServerBase, MIMO_API_UA } from "../shared/mimoAccount.js";
|
||||
|
||||
// Desktop-exclusive Preview models. These are served by the account service's
|
||||
// /api/route proxy, authorized by the Xiaomi account session (NOT the sk- key).
|
||||
// See shared/mimoAccount.js for the session handshake.
|
||||
const PREVIEW_MODELS = new Set(["mimo-x-pro-preview", "mimo-x-flash-preview"]);
|
||||
// Dual-route v2.6 models.
|
||||
// v2.6 models dynamically route to the account service when desktop session credentials
|
||||
// (mimoPassToken or account cookie) are present to consume weekly quota, falling back to
|
||||
// the cloud API (sk- key) otherwise.
|
||||
const ACCOUNT_MODELS = new Set([
|
||||
"mimo-v2.6-pro",
|
||||
"mimo-v2.6-flash",
|
||||
"mimo-v2.6-pro-ultraspeed",
|
||||
]);
|
||||
|
||||
// Session cookie resolved in execute() (async) and read back by buildHeaders()
|
||||
// (sync — BaseExecutor.execute does not await it). Carried on the per-request
|
||||
@@ -23,15 +28,24 @@ export class XiaomiMimoExecutor extends DefaultExecutor {
|
||||
super("xiaomi-mimo");
|
||||
}
|
||||
|
||||
static isPreviewModel(model) {
|
||||
return PREVIEW_MODELS.has(bareModel(model));
|
||||
static isAccountRoute(model, credentials) {
|
||||
const bare = bareModel(model);
|
||||
if (!ACCOUNT_MODELS.has(bare)) return false;
|
||||
return Boolean(
|
||||
credentials?.[COOKIE_KEY] ||
|
||||
credentials?.providerSpecificData?.mimoPassToken
|
||||
);
|
||||
}
|
||||
|
||||
isAccountRoute(model, credentials) {
|
||||
return XiaomiMimoExecutor.isAccountRoute(model, credentials);
|
||||
}
|
||||
|
||||
buildUrl(model, stream, urlIndex = 0, credentials = null) {
|
||||
// Preview models live on the account-service route, which is not one of the
|
||||
// Account route models live on the account-service route, which is not one of the
|
||||
// declared transports — resolve it before the default runtimeTransport path.
|
||||
if (XiaomiMimoExecutor.isPreviewModel(model)) {
|
||||
return `${MIMO_API_BASE}/api/route/chat/completions`;
|
||||
if (this.isAccountRoute(model, credentials)) {
|
||||
return `${resolveMimoServerBase(credentials?.providerSpecificData)}/api/route/chat/completions`;
|
||||
}
|
||||
// Cloud API models keep default handling, so a Claude-format client reaches
|
||||
// the /anthropic/v1/messages transport.
|
||||
@@ -39,8 +53,8 @@ export class XiaomiMimoExecutor extends DefaultExecutor {
|
||||
}
|
||||
|
||||
buildHeaders(credentials, stream = true, url, model) {
|
||||
if (XiaomiMimoExecutor.isPreviewModel(model) && credentials?.[COOKIE_KEY]) {
|
||||
// Preview models authenticate with the account-session cookie, not the key.
|
||||
if (this.isAccountRoute(model, credentials) && credentials?.[COOKIE_KEY]) {
|
||||
// Account route models authenticate with the account-session cookie, not the key.
|
||||
return {
|
||||
"Content-Type": "application/json",
|
||||
Accept: stream ? "text/event-stream" : "application/json",
|
||||
@@ -52,17 +66,22 @@ export class XiaomiMimoExecutor extends DefaultExecutor {
|
||||
}
|
||||
|
||||
transformRequest(model, body, stream, credentials) {
|
||||
// super runs stripUnsupportedParams, which flattens Preview content-part
|
||||
// super runs stripUnsupportedParams, which flattens content-part
|
||||
// arrays (see the xiaomi-mimo rule in translator/concerns/paramSupport.js).
|
||||
const out = super.transformRequest(model, body, stream, credentials);
|
||||
|
||||
// Preview models: thinking/params get defaults only — never override what the
|
||||
// caller set explicitly. (body.model is already `xiaomi/<id>` via upstreamModelId.)
|
||||
if (XiaomiMimoExecutor.isPreviewModel(model)) {
|
||||
if (out.thinking == null) out.thinking = { type: "enabled" };
|
||||
// Account route models: bridge reasoning_effort to official output_config.effort
|
||||
// (matches MiMo Desktop app.asar behavior).
|
||||
if (this.isAccountRoute(model, credentials)) {
|
||||
const rawEffort = out.reasoning_effort || body?.reasoning_effort || body?.output_config?.effort;
|
||||
if (rawEffort) {
|
||||
delete out.reasoning_effort;
|
||||
const norm = String(rawEffort).toLowerCase() === "xhigh" ? "high" : String(rawEffort).toLowerCase();
|
||||
out.output_config = { ...(out.output_config || {}), effort: norm };
|
||||
}
|
||||
|
||||
if (out.temperature == null) out.temperature = 1.0;
|
||||
if (out.top_p == null) out.top_p = 0.95;
|
||||
if (!out.max_tokens) out.max_tokens = 4096;
|
||||
}
|
||||
|
||||
return out;
|
||||
@@ -70,13 +89,11 @@ export class XiaomiMimoExecutor extends DefaultExecutor {
|
||||
|
||||
async execute(args) {
|
||||
const { model, credentials, proxyOptions = null } = args;
|
||||
if (!XiaomiMimoExecutor.isPreviewModel(model)) return super.execute(args);
|
||||
if (!this.isAccountRoute(model, credentials)) return super.execute(args);
|
||||
|
||||
const cookie = await getMimoAccountCookie(credentials?.providerSpecificData, proxyOptions);
|
||||
if (!cookie) {
|
||||
throw new Error(
|
||||
"Xiaomi MiMo account session unavailable. Sign in to MiMo Desktop once so its passToken is present, then retry.",
|
||||
);
|
||||
return super.execute(args);
|
||||
}
|
||||
credentials[COOKIE_KEY] = cookie;
|
||||
const result = await super.execute(args);
|
||||
@@ -94,6 +111,6 @@ export class XiaomiMimoExecutor extends DefaultExecutor {
|
||||
}
|
||||
}
|
||||
|
||||
export const __test__ = { PREVIEW_MODELS, bareModel, COOKIE_KEY };
|
||||
export const __test__ = { ACCOUNT_MODELS, bareModel, COOKIE_KEY };
|
||||
|
||||
export default XiaomiMimoExecutor;
|
||||
|
||||
Reference in New Issue
Block a user