Merge branch 'master' into gitea/new_feature
This commit is contained in:
@@ -136,6 +136,10 @@ export class AntigravityExecutor extends BaseExecutor {
|
||||
transformRequest(model, body, stream, credentials) {
|
||||
const projectId = credentials?.projectId || this.generateProjectId();
|
||||
|
||||
// OpenAI clients may include stream_options even for non-streaming calls.
|
||||
// Google generateContent rejects that combination before processing the request.
|
||||
if (stream !== true) delete body.stream_options;
|
||||
|
||||
// ─── Image generation: completely different request structure ───
|
||||
if (isImageModel(model)) {
|
||||
const imageConfig = parseImageConfig(model);
|
||||
@@ -264,7 +268,7 @@ export class AntigravityExecutor extends BaseExecutor {
|
||||
return {
|
||||
...body,
|
||||
project: projectId,
|
||||
model: model,
|
||||
model: body.model || model,
|
||||
userAgent: "antigravity",
|
||||
requestType: "agent",
|
||||
requestId: buildIdeRequestId({ body, request: transformedRequest, credentials, model, requestType: "agent" }),
|
||||
|
||||
@@ -127,7 +127,7 @@ export class BaseExecutor {
|
||||
for (let urlIndex = 0; urlIndex < fallbackCount; urlIndex++) {
|
||||
const url = this.buildUrl(model, stream, urlIndex, credentials);
|
||||
const transformedBody = this.transformRequest(model, body, stream, credentials);
|
||||
const headers = this.buildHeaders(credentials, stream);
|
||||
const headers = this.buildHeaders(credentials, stream, url);
|
||||
|
||||
if (!retryAttemptsByUrl[urlIndex]) retryAttemptsByUrl[urlIndex] = 0;
|
||||
|
||||
|
||||
30
open-sse/executors/codebuddy-intl.js
Normal file
30
open-sse/executors/codebuddy-intl.js
Normal file
@@ -0,0 +1,30 @@
|
||||
import { DefaultExecutor } from "./default.js";
|
||||
|
||||
/**
|
||||
* CodeBuddyIntlExecutor — talks to https://www.codebuddy.ai/v2/chat/completions
|
||||
*
|
||||
* Same OpenAI-compatible-but-stream-only gateway behavior as codebuddy-cn:
|
||||
* non-stream requests are rejected, and reasoning is surfaced only when the
|
||||
* request carries the IDE's OpenAI-style reasoning params. Force stream and
|
||||
* mirror reasoning_summary exactly like CodeBuddyExecutor.
|
||||
*/
|
||||
export class CodeBuddyIntlExecutor extends DefaultExecutor {
|
||||
constructor() {
|
||||
super("codebuddy-intl");
|
||||
}
|
||||
|
||||
transformRequest(model, body, stream, credentials) {
|
||||
const transformed = super.transformRequest(model, body, stream, credentials);
|
||||
transformed.stream = true;
|
||||
|
||||
const eff = transformed.reasoning_effort;
|
||||
if (eff === "none" || eff === "off") {
|
||||
delete transformed.reasoning_effort;
|
||||
} else if (eff) {
|
||||
transformed.reasoning_summary = "auto";
|
||||
}
|
||||
return transformed;
|
||||
}
|
||||
}
|
||||
|
||||
export default CodeBuddyIntlExecutor;
|
||||
@@ -1,18 +1,22 @@
|
||||
import { BaseExecutor } from "./base.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
import { PROVIDERS, PROVIDER_OAUTH } from "../config/providers.js";
|
||||
import { HTTP_STATUS } from "../config/runtimeConfig.js";
|
||||
import {
|
||||
generateCursorBody,
|
||||
encodeField,
|
||||
wrapConnectRPCFrame,
|
||||
decodeMessage,
|
||||
parseConnectRPCFrame,
|
||||
extractTextFromResponse
|
||||
} 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 { chatChunkSse, sseChunk } from "../utils/sse.js";
|
||||
import { FORMATS } from "../translator/formats.js";
|
||||
import { proxyAwareFetch } from "../utils/proxyFetch.js";
|
||||
import zlib from "zlib";
|
||||
import crypto from "crypto";
|
||||
|
||||
// Detect cloud environment
|
||||
const isCloudEnv = () => {
|
||||
@@ -38,6 +42,130 @@ const COMPRESS_FLAG = {
|
||||
GZIP_TRAILER: 0x03
|
||||
};
|
||||
|
||||
const AGENT_RUN_PATH = "/agent.v1.AgentService/Run";
|
||||
const PROTOBUF_LEN = 2;
|
||||
const PROTOBUF_VARINT = 0;
|
||||
|
||||
function concatBuffers(...parts) {
|
||||
const length = parts.reduce((total, part) => total + part.length, 0);
|
||||
const result = new Uint8Array(length);
|
||||
let offset = 0;
|
||||
for (const part of parts) {
|
||||
result.set(part, offset);
|
||||
offset += part.length;
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
const agentString = (field, value) => encodeField(field, PROTOBUF_LEN, value);
|
||||
const agentMessage = (field, value) => encodeField(field, PROTOBUF_LEN, value);
|
||||
const agentBool = (field, value) => encodeField(field, PROTOBUF_VARINT, value ? 1 : 0);
|
||||
|
||||
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")
|
||||
.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 encodeHistoryMessage(message) {
|
||||
const content = textFromContent(message?.content);
|
||||
if (!content) return null;
|
||||
|
||||
// ConversationHistoryMessage.user / .assistant -> repeated content -> text.
|
||||
const text = agentString(1, content);
|
||||
if (message.role === "assistant") {
|
||||
return agentMessage(2, agentMessage(1, agentMessage(1, text)));
|
||||
}
|
||||
return agentMessage(1, agentMessage(1, agentMessage(1, text)));
|
||||
}
|
||||
|
||||
function buildAgentRunFrame(messages, model) {
|
||||
const system = messages
|
||||
.filter((message) => message?.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 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.";
|
||||
|
||||
// agent.v1.UserMessageAction.user_message and its optional history.
|
||||
const userMessage = concatBuffers(
|
||||
agentString(1, userText),
|
||||
agentString(2, crypto.randomUUID()),
|
||||
);
|
||||
const conversationHistory = history.length
|
||||
? concatBuffers(...history.map((entry) => agentMessage(1, entry)))
|
||||
: null;
|
||||
const userAction = concatBuffers(
|
||||
agentMessage(1, userMessage),
|
||||
...(conversationHistory ? [agentMessage(7, conversationHistory)] : []),
|
||||
);
|
||||
const conversationAction = agentMessage(1, userAction);
|
||||
const requestedModel = concatBuffers(agentString(1, model), agentBool(7, true));
|
||||
const runRequest = concatBuffers(
|
||||
// An empty ConversationStateStructure starts a fresh local agent session.
|
||||
agentMessage(1, new Uint8Array()),
|
||||
agentMessage(2, conversationAction),
|
||||
...(system ? [agentString(8, system)] : []),
|
||||
agentMessage(9, requestedModel),
|
||||
);
|
||||
|
||||
// agent.v1.AgentClientMessage.run_request.
|
||||
return wrapConnectRPCFrame(agentMessage(1, runRequest));
|
||||
}
|
||||
|
||||
function extractAgentString(message, field) {
|
||||
const value = message?.get(field)?.[0]?.value;
|
||||
return value ? Buffer.from(value).toString("utf8") : "";
|
||||
}
|
||||
|
||||
function decodeAgentFrames(buffer, onFrame) {
|
||||
let pending = Buffer.from(buffer || []);
|
||||
while (pending.length >= 5) {
|
||||
const flags = pending[0];
|
||||
const length = pending.readUInt32BE(1);
|
||||
if (pending.length < 5 + length) break;
|
||||
let payload = pending.subarray(5, 5 + length);
|
||||
pending = pending.subarray(5 + length);
|
||||
if (flags & COMPRESS_FLAG.GZIP) {
|
||||
payload = zlib.gunzipSync(payload);
|
||||
}
|
||||
if (!(flags & COMPRESS_FLAG.TRAILER)) onFrame(payload);
|
||||
}
|
||||
return pending;
|
||||
}
|
||||
|
||||
function createRequestContextResponse() {
|
||||
// AgentService asks every run for client context. 9router has no IDE file
|
||||
// context, so acknowledge with an empty RequestContext.
|
||||
const requestContextSuccess = agentMessage(1, new Uint8Array());
|
||||
const requestContextResult = agentMessage(1, requestContextSuccess);
|
||||
const execClientMessage = agentMessage(10, requestContextResult);
|
||||
return wrapConnectRPCFrame(agentMessage(2, execClientMessage));
|
||||
}
|
||||
|
||||
const CURSOR_STREAM_DEBUG = process.env.CURSOR_STREAM_DEBUG === "1";
|
||||
const debugLog = (...args) => {
|
||||
if (CURSOR_STREAM_DEBUG) console.log(...args);
|
||||
@@ -253,7 +381,304 @@ export class CursorExecutor extends BaseExecutor {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* AgentService (agent.api5.cursor.sh) is HTTP/2-only. Node's fetch/undici speaks
|
||||
* HTTP/1.1 and fails with HTTPParserError on the h2 preface — use http2 duplex.
|
||||
*/
|
||||
openAgentHttp2Stream(url, headers, signal) {
|
||||
if (!http2) {
|
||||
throw new Error("HTTP/2 is required for Cursor AgentService (endpoint is h2-only)");
|
||||
}
|
||||
|
||||
const urlObj = new URL(url);
|
||||
const client = http2.connect(`https://${urlObj.host}`);
|
||||
const chunkQueue = [];
|
||||
let waiting = null;
|
||||
let ended = false;
|
||||
let streamError = null;
|
||||
let req = null;
|
||||
|
||||
const wake = (result) => {
|
||||
if (!waiting) return;
|
||||
const resolve = waiting;
|
||||
waiting = null;
|
||||
resolve(result);
|
||||
};
|
||||
|
||||
const fail = (error) => {
|
||||
if (streamError) return;
|
||||
streamError = error;
|
||||
ended = true;
|
||||
wake(null);
|
||||
};
|
||||
|
||||
const close = () => {
|
||||
try { req?.destroy(); } catch {}
|
||||
try { client.close(); } catch {}
|
||||
};
|
||||
|
||||
client.on("error", fail);
|
||||
|
||||
req = client.request({
|
||||
":method": "POST",
|
||||
":path": urlObj.pathname,
|
||||
":authority": urlObj.host,
|
||||
":scheme": "https",
|
||||
...headers,
|
||||
});
|
||||
|
||||
req.on("error", fail);
|
||||
req.on("data", (chunk) => {
|
||||
if (waiting) wake({ value: chunk, done: false });
|
||||
else chunkQueue.push(chunk);
|
||||
});
|
||||
req.on("end", () => {
|
||||
ended = true;
|
||||
wake({ value: undefined, done: true });
|
||||
});
|
||||
|
||||
if (signal) {
|
||||
const onAbort = () => {
|
||||
fail(new Error("Request aborted"));
|
||||
close();
|
||||
};
|
||||
if (signal.aborted) onAbort();
|
||||
else signal.addEventListener("abort", onAbort, { once: true });
|
||||
}
|
||||
|
||||
const responseHeaders = new Promise((resolve, reject) => {
|
||||
const onEarlyError = (error) => reject(error);
|
||||
client.once("error", onEarlyError);
|
||||
req.once("error", onEarlyError);
|
||||
req.once("response", (hdrs) => {
|
||||
client.off("error", onEarlyError);
|
||||
req.off("error", onEarlyError);
|
||||
resolve(hdrs);
|
||||
});
|
||||
});
|
||||
|
||||
return {
|
||||
responseHeaders,
|
||||
write(frame) {
|
||||
if (req && !req.destroyed) req.write(Buffer.from(frame));
|
||||
},
|
||||
end() {
|
||||
try { if (req && !req.destroyed) req.end(); } catch {}
|
||||
},
|
||||
close,
|
||||
async read() {
|
||||
if (chunkQueue.length) return { value: chunkQueue.shift(), done: false };
|
||||
if (ended) {
|
||||
if (streamError) throw streamError;
|
||||
return { value: undefined, done: true };
|
||||
}
|
||||
const result = await new Promise((resolve) => { waiting = resolve; });
|
||||
if (streamError) throw streamError;
|
||||
return result || { value: undefined, done: true };
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
async executeAgent({ model, body, stream, credentials, signal }) {
|
||||
const agentEndpoint = PROVIDER_OAUTH.cursor?.agentEndpoint;
|
||||
if (!agentEndpoint) throw new Error("Cursor AgentService endpoint is not configured");
|
||||
|
||||
const url = `${agentEndpoint}${AGENT_RUN_PATH}`;
|
||||
const headers = this.buildHeaders(credentials);
|
||||
const requestController = new AbortController();
|
||||
if (signal?.addEventListener) {
|
||||
signal.addEventListener("abort", () => requestController.abort(signal.reason), { once: true });
|
||||
}
|
||||
|
||||
let session;
|
||||
try {
|
||||
session = this.openAgentHttp2Stream(url, headers, requestController.signal);
|
||||
session.write(buildAgentRunFrame(body.messages || [], model));
|
||||
} catch (error) {
|
||||
throw new Error(`Cursor AgentService request failed: ${error.message}`);
|
||||
}
|
||||
|
||||
let responseHeaders;
|
||||
try {
|
||||
responseHeaders = await session.responseHeaders;
|
||||
} catch (error) {
|
||||
session.close();
|
||||
throw new Error(`Cursor AgentService request failed: ${error.message}`);
|
||||
}
|
||||
|
||||
const status = Number(responseHeaders[":status"] || 0);
|
||||
if (status !== 200) {
|
||||
let errorText = "";
|
||||
try {
|
||||
while (true) {
|
||||
const { done, value } = await session.read();
|
||||
if (done) break;
|
||||
errorText += Buffer.from(value).toString("utf8");
|
||||
}
|
||||
} catch {}
|
||||
session.close();
|
||||
return {
|
||||
response: new Response(JSON.stringify({
|
||||
error: { message: `Cursor AgentService ${status}: ${errorText || "request failed"}`, type: "api_error" },
|
||||
}), { status: status || HTTP_STATUS.SERVER_ERROR, headers: { "Content-Type": "application/json" } }),
|
||||
url,
|
||||
headers,
|
||||
transformedBody: body,
|
||||
responseFormat: FORMATS.OPENAI,
|
||||
};
|
||||
}
|
||||
|
||||
// The Claude SSE translator derives Anthropic's message ID by stripping
|
||||
// `chatcmpl-`. Keep the remaining ID in Anthropic's required `msg_` form
|
||||
// so strict clients such as Claude Code accept the completed stream.
|
||||
const responseId = `chatcmpl-msg_${Date.now()}`;
|
||||
const created = Math.floor(Date.now() / 1000);
|
||||
let pending = Buffer.alloc(0);
|
||||
let finished = false;
|
||||
|
||||
const consume = async (onEvent) => {
|
||||
try {
|
||||
while (!finished) {
|
||||
const { done, value } = await session.read();
|
||||
if (done) break;
|
||||
pending = Buffer.concat([pending, Buffer.from(value)]);
|
||||
pending = decodeAgentFrames(pending, (payload) => {
|
||||
// A single read can carry several frames; once the turn is over the
|
||||
// rest of the batch must not reach the already-closed controller.
|
||||
if (finished) return;
|
||||
const serverMessage = decodeMessage(payload);
|
||||
|
||||
// agent.v1.AgentServerMessage.interaction_update
|
||||
if (serverMessage.has(1)) {
|
||||
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 });
|
||||
}
|
||||
// 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.
|
||||
if (update.has(14)) {
|
||||
finished = true;
|
||||
onEvent({ type: "done" });
|
||||
}
|
||||
}
|
||||
|
||||
// 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());
|
||||
} 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" });
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
} finally {
|
||||
try { session.end(); } catch {}
|
||||
try { session.close(); } catch {}
|
||||
if (!finished) onEvent({ type: "done" });
|
||||
}
|
||||
};
|
||||
|
||||
if (stream === false) {
|
||||
let content = "";
|
||||
let reasoning = "";
|
||||
let agentError = null;
|
||||
await consume((event) => {
|
||||
if (event.type === "text") content += event.value;
|
||||
else if (event.type === "thinking") reasoning += event.value;
|
||||
else if (event.type === "error") agentError = event.value;
|
||||
});
|
||||
if (agentError) {
|
||||
return {
|
||||
response: new Response(JSON.stringify({ error: { message: agentError, type: "api_error" } }), {
|
||||
status: HTTP_STATUS.BAD_REQUEST,
|
||||
headers: { "Content-Type": "application/json" },
|
||||
}),
|
||||
url,
|
||||
headers,
|
||||
transformedBody: body,
|
||||
responseFormat: FORMATS.OPENAI,
|
||||
};
|
||||
}
|
||||
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" }],
|
||||
usage: estimateUsage(body, content.length, FORMATS.OPENAI),
|
||||
}), { headers: { "Content-Type": "application/json" } }),
|
||||
url,
|
||||
headers,
|
||||
transformedBody: body,
|
||||
responseFormat: FORMATS.OPENAI,
|
||||
};
|
||||
}
|
||||
|
||||
const encoder = new TextEncoder();
|
||||
const responseStream = new ReadableStream({
|
||||
start(controller) {
|
||||
consume((event) => {
|
||||
if (event.type === "text") {
|
||||
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 === "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
|
||||
// usage tracking must not record the turn as a success.
|
||||
controller.enqueue(encoder.encode(sseChunk({ error: { message: event.value, type: "api_error" } })));
|
||||
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(SSE_DONE));
|
||||
controller.close();
|
||||
}
|
||||
}).catch((error) => controller.error(error));
|
||||
},
|
||||
cancel() {
|
||||
requestController.abort();
|
||||
},
|
||||
});
|
||||
|
||||
return {
|
||||
response: new Response(responseStream, { headers: SSE_HEADERS }),
|
||||
url,
|
||||
headers,
|
||||
transformedBody: body,
|
||||
responseFormat: FORMATS.OPENAI,
|
||||
};
|
||||
}
|
||||
|
||||
async execute({ model, body, stream, credentials, signal, log, proxyOptions = null }) {
|
||||
if (isAgentTextRequest(body)) {
|
||||
try {
|
||||
return await this.executeAgent({ model, body, stream, credentials, signal });
|
||||
} catch (error) {
|
||||
return {
|
||||
response: new Response(JSON.stringify({
|
||||
error: { message: error.message, type: "connection_error", code: "" },
|
||||
}), { status: HTTP_STATUS.SERVER_ERROR, headers: { "Content-Type": "application/json" } }),
|
||||
url: `${PROVIDER_OAUTH.cursor?.agentEndpoint || ""}${AGENT_RUN_PATH}`,
|
||||
headers: {},
|
||||
transformedBody: body,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
const url = this.buildUrl();
|
||||
const headers = this.buildHeaders(credentials);
|
||||
const transformedBody = this.transformRequest(model, body, stream, credentials);
|
||||
|
||||
@@ -38,7 +38,8 @@ function applyAuth(headers, desc, credentials) {
|
||||
|
||||
// Provider-specific header quirks kept as small hooks (not pure auth).
|
||||
const HEADER_HOOKS = {
|
||||
kimiHeaders: (h) => Object.assign(h, buildKimiHeaders()),
|
||||
// Stable device_id from OAuth connection (CLIProxyAPI KimiTokenStorage.DeviceID)
|
||||
kimiHeaders: (h, c) => Object.assign(h, buildKimiHeaders(c?.providerSpecificData?.deviceId)),
|
||||
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) => {
|
||||
@@ -227,7 +228,8 @@ export class DefaultExecutor extends BaseExecutor {
|
||||
kiro: () => this.refreshKiro(credentials.refreshToken, proxyOptions),
|
||||
cline: () => this.refreshCline(credentials.refreshToken, proxyOptions),
|
||||
clinepass: () => this.refreshCline(credentials.refreshToken, proxyOptions),
|
||||
"kimi-coding": () => this.refreshKimiCoding(credentials.refreshToken, proxyOptions),
|
||||
kimi: () => this.refreshKimi(credentials, proxyOptions),
|
||||
"kimi-coding": () => this.refreshKimi(credentials, proxyOptions),
|
||||
kilocode: () => this.refreshKilocode(credentials.refreshToken, proxyOptions)
|
||||
};
|
||||
|
||||
@@ -307,16 +309,20 @@ export class DefaultExecutor extends BaseExecutor {
|
||||
return { accessToken, refreshToken: data?.refreshToken || refreshToken, expiresIn };
|
||||
}
|
||||
|
||||
async refreshKimiCoding(refreshToken, proxyOptions = null) {
|
||||
const kimiHeaders = buildKimiHeaders();
|
||||
const response = await proxyAwareFetch(PROVIDERS["kimi-coding"].refreshUrl, {
|
||||
// CLIProxyAPI DeviceFlowClient.RefreshToken — form body + X-Msh-* headers + stable device_id
|
||||
async refreshKimi(credentials, proxyOptions = null) {
|
||||
const refreshToken = credentials.refreshToken;
|
||||
const cfg = PROVIDERS.kimi || PROVIDERS["kimi-coding"];
|
||||
if (!cfg?.refreshUrl || !cfg?.clientId) return null;
|
||||
const kimiHeaders = buildKimiHeaders(credentials?.providerSpecificData?.deviceId);
|
||||
const response = await proxyAwareFetch(cfg.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: PROVIDERS["kimi-coding"].clientId })
|
||||
body: new URLSearchParams({ grant_type: "refresh_token", refresh_token: refreshToken, client_id: cfg.clientId })
|
||||
}, proxyOptions);
|
||||
if (!response.ok) return null;
|
||||
const tokens = await response.json();
|
||||
|
||||
847
open-sse/executors/devin-cli.js
Normal file
847
open-sse/executors/devin-cli.js
Normal file
@@ -0,0 +1,847 @@
|
||||
/**
|
||||
* DevinCliExecutor — routes completions through the official Devin CLI binary
|
||||
* via the Agent Client Protocol (ACP) JSON-RPC 2.0 over stdio.
|
||||
*
|
||||
* Protocol flow:
|
||||
* 1. Spawn `devin acp` (default agent = full built-in tools: fs/shell/search).
|
||||
* Set CLI_DEVIN_AGENT_TYPE=summarizer for a tool-less, text-only mode.
|
||||
* 2. Send: initialize → session/new (with model + cwd + mcpServers) → session/prompt.
|
||||
* 3. Receive: session/update notifications (agent_message_chunk = reply text,
|
||||
* tool_call/tool_call_update = built-in tool invocations, surfaced as text).
|
||||
* When devin calls a client-tool from the exposed MCP ("Calling mcp_X from
|
||||
* clientTools"), it is bridged to an OpenAI tool_use and the turn ends.
|
||||
* 4. Emit deltas as OpenAI-compatible SSE chunks.
|
||||
* 5. Kill subprocess on _cognition.ai/agent_stopped or error.
|
||||
*
|
||||
* Auth: noAuth — the subprocess inherits the parent env and uses credentials
|
||||
* stored by `devin auth login` (~/.local/share/devin/credentials.toml).
|
||||
*
|
||||
* Binary discovery: CLI_DEVIN_BIN env → PATH lookup → platform installer paths.
|
||||
*/
|
||||
|
||||
import { spawn } from "node:child_process";
|
||||
import path from "node:path";
|
||||
import os from "node:os";
|
||||
import fs from "node:fs";
|
||||
import { BaseExecutor } from "./base.js";
|
||||
|
||||
// ─── Binary discovery ────────────────────────────────────────────────────────
|
||||
|
||||
function resolveDevinBin() {
|
||||
// 1. Explicit override
|
||||
const envBin = process.env.CLI_DEVIN_BIN?.trim();
|
||||
if (envBin) return envBin;
|
||||
|
||||
const isWin = process.platform === "win32";
|
||||
const home = os.homedir();
|
||||
|
||||
// 2. Known installer / package-manager locations. spawn uses shell:false on
|
||||
// macOS/Linux, so process.env.PATH alone may miss ~/.local/bin, Homebrew,
|
||||
// Scoop, etc. when the server runs detached (tray/daemon/launchd) without
|
||||
// a login shell — probe these explicitly before falling back to PATH.
|
||||
const candidates = isWin
|
||||
? [
|
||||
// Official installer: %LOCALAPPDATA%\devin\cli\bin\devin.exe
|
||||
path.join(process.env.LOCALAPPDATA || path.join(home, "AppData", "Local"), "devin", "cli", "bin", "devin.exe"),
|
||||
path.join(home, ".local", "bin", "devin.exe"),
|
||||
path.join(home, "scoop", "shims", "devin.exe"),
|
||||
path.join(process.env.LOCALAPPDATA || path.join(home, "AppData", "Local"), "Programs", "devin", "devin.exe"),
|
||||
]
|
||||
: [
|
||||
path.join(home, ".local", "share", "devin", "bin", "devin"),
|
||||
path.join(home, ".devin", "bin", "devin"),
|
||||
path.join(home, ".local", "bin", "devin"), // pipx / user install
|
||||
"/opt/homebrew/bin/devin", // Homebrew (Apple Silicon)
|
||||
"/usr/local/bin/devin", // Homebrew (Intel) / manual
|
||||
"/usr/bin/devin",
|
||||
];
|
||||
for (const candidate of candidates) {
|
||||
if (fs.existsSync(candidate)) return candidate;
|
||||
}
|
||||
|
||||
// 3. Fallback — rely on process.env.PATH
|
||||
return isWin ? "devin.exe" : "devin";
|
||||
}
|
||||
|
||||
// ─── ACP JSON-RPC helper ────────────────────────────────────────────────────
|
||||
|
||||
function rpc(method, params, id) {
|
||||
const msg = { jsonrpc: "2.0", method, params };
|
||||
if (id !== undefined) msg.id = id;
|
||||
return JSON.stringify(msg) + "\n";
|
||||
}
|
||||
|
||||
// ─── Client-tools → MCP bridge ───────────────────────────────────────────────
|
||||
// devin only invokes built-in + MCP tools, not OpenAI function-calling schemas.
|
||||
// body.tools are exposed as a stdio MCP server "clientTools" so devin can call
|
||||
// them. When devin calls one, we emit OpenAI tool_use and end the turn; the
|
||||
// client executes and returns tool_result on the next request. That next request
|
||||
// re-spawns with the full history (including tool_calls + tool results) and
|
||||
// seeds the MCP server with those results so a re-call gets the real data.
|
||||
// Tool schemas via DEVIN_MCP_TOOLS; prior results via DEVIN_MCP_RESULTS.
|
||||
|
||||
const CLIENT_TOOLS_MCP_SCRIPT = `
|
||||
import readline from "node:readline";
|
||||
const TOOLS = JSON.parse(process.env.DEVIN_MCP_TOOLS || "[]");
|
||||
const RESULTS = JSON.parse(process.env.DEVIN_MCP_RESULTS || "{}");
|
||||
const rl = readline.createInterface({ input: process.stdin });
|
||||
function send(o){ process.stdout.write(JSON.stringify(o) + "\\n"); }
|
||||
rl.on("line", (line) => {
|
||||
let m; try { m = JSON.parse(line); } catch { return; }
|
||||
if (m.method === "initialize") {
|
||||
send({ jsonrpc: "2.0", id: m.id, result: { protocolVersion: "2024-11-05", capabilities: { tools: {} }, serverInfo: { name: "clientTools", version: "1.0" } } });
|
||||
} else if (m.method === "tools/list") {
|
||||
send({ jsonrpc: "2.0", id: m.id, result: { tools: TOOLS } });
|
||||
} else if (m.method === "tools/call") {
|
||||
const name = m.params?.name || "";
|
||||
const seeded = RESULTS[name];
|
||||
const text = seeded !== undefined
|
||||
? String(seeded)
|
||||
: "(awaiting client tool_result)";
|
||||
process.stderr.write("[client-tools] tool_call name=" + name + " seeded=" + (seeded !== undefined) + "\\n");
|
||||
send({ jsonrpc: "2.0", id: m.id, result: { content: [{ type: "text", text }] } });
|
||||
}
|
||||
});
|
||||
`.trimStart();
|
||||
|
||||
function ensureClientToolsScript() {
|
||||
const scriptPath = path.join(os.tmpdir(), "9router-devin-client-tools.mjs");
|
||||
// Always rewrite so script upgrades land without a process restart.
|
||||
fs.writeFileSync(scriptPath, CLIENT_TOOLS_MCP_SCRIPT);
|
||||
return scriptPath;
|
||||
}
|
||||
|
||||
// Map OpenAI tools ([{type:"function",function:{name,description,parameters}}])
|
||||
// to MCP tool declarations ([{name,description,inputSchema}]).
|
||||
// devin only discovers MCP tools whose name carries the `mcp_` prefix, so we
|
||||
// add it here and strip it back when bridging the call to the client.
|
||||
const MCP_TOOL_PREFIX = "mcp_";
|
||||
function toMcpToolName(name) {
|
||||
return name.startsWith(MCP_TOOL_PREFIX) ? name : MCP_TOOL_PREFIX + name;
|
||||
}
|
||||
function fromMcpToolName(name) {
|
||||
return name.startsWith(MCP_TOOL_PREFIX) ? name.slice(MCP_TOOL_PREFIX.length) : name;
|
||||
}
|
||||
|
||||
function buildClientToolsMcp(tools, resultMap) {
|
||||
const mcpTools = [];
|
||||
for (const t of tools) {
|
||||
if (!t) continue;
|
||||
const f = t.function || t;
|
||||
if (!f?.name) continue;
|
||||
mcpTools.push({
|
||||
name: toMcpToolName(f.name),
|
||||
description: f.description || "",
|
||||
inputSchema: f.parameters || f.input_schema || { type: "object", properties: {} },
|
||||
});
|
||||
}
|
||||
if (!mcpTools.length) return null;
|
||||
const env = { DEVIN_MCP_TOOLS: JSON.stringify(mcpTools) };
|
||||
if (resultMap && Object.keys(resultMap).length) {
|
||||
env.DEVIN_MCP_RESULTS = JSON.stringify(resultMap);
|
||||
}
|
||||
return {
|
||||
command: process.execPath,
|
||||
args: [ensureClientToolsScript()],
|
||||
env,
|
||||
};
|
||||
}
|
||||
|
||||
// Extract tool_result content keyed by MCP tool name (mcp_<original>).
|
||||
// Walks messages: assistant.tool_calls id→name, role=tool tool_call_id→content.
|
||||
function extractClientToolResults(messages) {
|
||||
const idToMcpName = new Map();
|
||||
const results = {};
|
||||
for (const m of messages) {
|
||||
if (m?.role === "assistant" && Array.isArray(m.tool_calls)) {
|
||||
for (const tc of m.tool_calls) {
|
||||
const name = tc?.function?.name || tc?.name;
|
||||
if (tc?.id && name) idToMcpName.set(tc.id, toMcpToolName(name));
|
||||
}
|
||||
}
|
||||
// Claude-style tool_use blocks in content
|
||||
if (m?.role === "assistant" && Array.isArray(m.content)) {
|
||||
for (const b of m.content) {
|
||||
if (b?.type === "tool_use" && b.id && b.name) {
|
||||
idToMcpName.set(b.id, toMcpToolName(b.name));
|
||||
}
|
||||
}
|
||||
}
|
||||
if (m?.role === "tool" && m.tool_call_id) {
|
||||
const mcpName = idToMcpName.get(m.tool_call_id);
|
||||
if (mcpName) {
|
||||
results[mcpName] =
|
||||
typeof m.content === "string" ? m.content : JSON.stringify(m.content ?? "");
|
||||
}
|
||||
}
|
||||
// Claude-style tool_result blocks in user content
|
||||
if (m?.role === "user" && Array.isArray(m.content)) {
|
||||
for (const b of m.content) {
|
||||
if (b?.type === "tool_result" && b.tool_use_id) {
|
||||
const mcpName = idToMcpName.get(b.tool_use_id);
|
||||
if (mcpName) {
|
||||
const c = b.content;
|
||||
results[mcpName] =
|
||||
typeof c === "string" ? c : JSON.stringify(c ?? "");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return results;
|
||||
}
|
||||
|
||||
// Resolve workspace cwd from client request (Codex/CLI env context, body fields).
|
||||
// Prefer an absolute existing path so agent file tools hit the user's project
|
||||
// instead of os.tmpdir() (which made relative create/delete inconsistent).
|
||||
function resolveWorkspaceCwd(body) {
|
||||
const candidates = [];
|
||||
const push = (v) => {
|
||||
if (typeof v === "string" && v.trim()) candidates.push(v.trim());
|
||||
};
|
||||
push(body?.cwd);
|
||||
push(body?.working_directory);
|
||||
push(body?.workdir);
|
||||
push(body?.workspace);
|
||||
push(body?.metadata?.cwd);
|
||||
push(body?.metadata?.working_directory);
|
||||
|
||||
const scanText = (text) => {
|
||||
if (typeof text !== "string") return;
|
||||
for (const m of text.matchAll(/<cwd>\s*([^<]+?)\s*<\/cwd>/gi)) push(m[1]);
|
||||
};
|
||||
const scanMessages = (msgs) => {
|
||||
if (!Array.isArray(msgs)) return;
|
||||
for (const msg of msgs) {
|
||||
if (!msg) continue;
|
||||
if (typeof msg.content === "string") scanText(msg.content);
|
||||
else if (Array.isArray(msg.content)) {
|
||||
for (const p of msg.content) {
|
||||
if (typeof p === "string") scanText(p);
|
||||
else if (p && typeof p === "object") {
|
||||
scanText(p.text);
|
||||
scanText(p.input_text);
|
||||
scanText(p.content);
|
||||
}
|
||||
}
|
||||
}
|
||||
// Responses API input items
|
||||
if (typeof msg === "string") scanText(msg);
|
||||
if (msg.type === "message" && Array.isArray(msg.content)) {
|
||||
for (const p of msg.content) scanText(p?.text || p?.input_text);
|
||||
}
|
||||
}
|
||||
};
|
||||
scanMessages(body?.messages);
|
||||
scanMessages(body?.input);
|
||||
|
||||
for (const c of candidates) {
|
||||
try {
|
||||
if (path.isAbsolute(c) && fs.existsSync(c) && fs.statSync(c).isDirectory()) {
|
||||
return c;
|
||||
}
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
}
|
||||
return os.tmpdir();
|
||||
}
|
||||
|
||||
// ─── Multi-turn message → single prompt builder ─────────────────────────────
|
||||
|
||||
function buildPromptText(messages) {
|
||||
// Inline the whole conversation so the model has full context, including
|
||||
// prior tool_calls / tool_results so it can continue after a client round-trip.
|
||||
const lines = [];
|
||||
for (const m of messages) {
|
||||
const role = String(m.role || "user");
|
||||
let text = "";
|
||||
if (typeof m.content === "string") {
|
||||
text = m.content;
|
||||
} else if (Array.isArray(m.content)) {
|
||||
for (const p of m.content) {
|
||||
if (!p || typeof p !== "object") continue;
|
||||
if (p.type === "text") text += String(p.text || "");
|
||||
else if (p.type === "tool_use") {
|
||||
text += `\n[Tool call ${p.name} id=${p.id}]\n${JSON.stringify(p.input ?? {})}\n`;
|
||||
} else if (p.type === "tool_result") {
|
||||
const c =
|
||||
typeof p.content === "string" ? p.content : JSON.stringify(p.content ?? "");
|
||||
text += `\n[Tool result id=${p.tool_use_id}]\n${c}\n`;
|
||||
}
|
||||
}
|
||||
}
|
||||
// OpenAI tool_calls on assistant messages
|
||||
if (role === "assistant" && Array.isArray(m.tool_calls) && m.tool_calls.length) {
|
||||
const parts = m.tool_calls.map((tc) => {
|
||||
const name = tc.function?.name || tc.name || "tool";
|
||||
const args = tc.function?.arguments ?? tc.arguments ?? {};
|
||||
const argStr = typeof args === "string" ? args : JSON.stringify(args);
|
||||
return `[Tool call ${name} id=${tc.id}]\n${argStr}`;
|
||||
});
|
||||
text = [text, ...parts].filter(Boolean).join("\n\n");
|
||||
}
|
||||
// OpenAI role=tool messages
|
||||
if (role === "tool") {
|
||||
const c = typeof m.content === "string" ? m.content : JSON.stringify(m.content ?? "");
|
||||
text = `[Tool result id=${m.tool_call_id || ""}]\n${c}`;
|
||||
}
|
||||
if (!text.trim()) continue;
|
||||
if (role === "system") {
|
||||
lines.push(`[System]\n${text}`);
|
||||
} else if (role === "assistant") {
|
||||
lines.push(`[Assistant]\n${text}`);
|
||||
} else if (role === "tool") {
|
||||
lines.push(`[Tool]\n${text}`);
|
||||
} else {
|
||||
lines.push(`[User]\n${text}`);
|
||||
}
|
||||
}
|
||||
return lines.join("\n\n") || "(empty)";
|
||||
}
|
||||
|
||||
// ─── DevinCliExecutor ─────────────────────────────────────────────────────────
|
||||
|
||||
export class DevinCliExecutor extends BaseExecutor {
|
||||
constructor() {
|
||||
super("devin-cli", { id: "devin-cli", baseUrl: "devin://acp/stdio" });
|
||||
}
|
||||
|
||||
buildUrl() {
|
||||
return "devin://acp/stdio";
|
||||
}
|
||||
|
||||
buildHeaders() {
|
||||
return {};
|
||||
}
|
||||
|
||||
transformRequest() {
|
||||
return null;
|
||||
}
|
||||
|
||||
async execute({ model, body, credentials, signal, log }) {
|
||||
const b = body ?? {};
|
||||
const messages = Array.isArray(b.messages)
|
||||
? b.messages
|
||||
: Array.isArray(b.input)
|
||||
? b.input
|
||||
: [];
|
||||
const promptText = buildPromptText(messages);
|
||||
const workspaceCwd = resolveWorkspaceCwd(b);
|
||||
const devinBin = resolveDevinBin();
|
||||
|
||||
log?.info?.(
|
||||
"DEVIN",
|
||||
`devin acp → model=${model}, bin=${devinBin}, cwd=${workspaceCwd}`
|
||||
);
|
||||
|
||||
// Optional MCP servers via DEVIN_MCP_SERVERS (JSON object, devin config format):
|
||||
// {"echo":{"command":"/abs/node","args":["/srv/echo.js"],"env":{"K":"V"}}}
|
||||
// Plus body.tools (OpenAI schema) → exposed as a "clientTools" MCP
|
||||
// server so devin can invoke client-defined tools (bridged back in Phase 2).
|
||||
// When any are present, a throwaway XDG_CONFIG_HOME holds devin/config.json so
|
||||
// the agent auto-connects them (session/new mcpServers alone doesn't spawn
|
||||
// them — see ACP mcp/connect, still unstable). Cleaned up on finish.
|
||||
// NOTE: this replaces the user's global devin MCP config for the subprocess.
|
||||
let mcpConfigDir = null;
|
||||
const mcpServers = {};
|
||||
const mcpJson = process.env.DEVIN_MCP_SERVERS?.trim();
|
||||
if (mcpJson) {
|
||||
try {
|
||||
Object.assign(mcpServers, JSON.parse(mcpJson));
|
||||
} catch (e) {
|
||||
log?.info?.("DEVIN", `DEVIN_MCP_SERVERS parse failed: ${e.message}`);
|
||||
}
|
||||
}
|
||||
const clientTools = Array.isArray(b.tools) ? b.tools.filter(Boolean) : [];
|
||||
const clientToolResults = extractClientToolResults(messages);
|
||||
const clientToolsMcp = buildClientToolsMcp(clientTools, clientToolResults);
|
||||
const hasClientTools = !!clientToolsMcp;
|
||||
if (clientToolsMcp) {
|
||||
mcpServers["clientTools"] = clientToolsMcp;
|
||||
const seeded = Object.keys(clientToolResults).length;
|
||||
log?.info?.(
|
||||
"DEVIN",
|
||||
`exposing ${clientTools.length} client tool(s) as MCP` +
|
||||
(seeded ? ` (seeded ${seeded} result(s))` : "")
|
||||
);
|
||||
}
|
||||
if (Object.keys(mcpServers).length) {
|
||||
try {
|
||||
mcpConfigDir = fs.mkdtempSync(path.join(os.tmpdir(), "devin-mcp-"));
|
||||
const cfgDev = path.join(mcpConfigDir, "devin");
|
||||
fs.mkdirSync(cfgDev, { recursive: true });
|
||||
fs.writeFileSync(
|
||||
path.join(cfgDev, "config.json"),
|
||||
JSON.stringify({ mcpServers })
|
||||
);
|
||||
log?.info?.("DEVIN", `mcp config written → ${mcpConfigDir}`);
|
||||
} catch (e) {
|
||||
log?.info?.("DEVIN", `mcp config write failed: ${e.message}`);
|
||||
mcpConfigDir = null;
|
||||
}
|
||||
}
|
||||
const cleanupMcp = () => {
|
||||
if (!mcpConfigDir) return;
|
||||
try {
|
||||
fs.rmSync(mcpConfigDir, { recursive: true, force: true });
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
mcpConfigDir = null;
|
||||
};
|
||||
|
||||
const sseStream = new ReadableStream({
|
||||
start(controller) {
|
||||
const enc = new TextEncoder();
|
||||
const emit = (data) => controller.enqueue(enc.encode(data));
|
||||
|
||||
// Inherit the parent environment so devin resolves stored CLI credentials
|
||||
// (~/.local/share/devin/credentials.toml from `devin auth login`). Do NOT
|
||||
// inject WINDSURF_API_KEY: this provider is noAuth, and a bogus/leaked key
|
||||
// overrides stored creds and makes devin return -32000 "invalid api key".
|
||||
const env = { ...process.env };
|
||||
// Auto-approve tool execution so the agent doesn't block waiting for a
|
||||
// session/request_permission response we never send (default mode would
|
||||
// hang the stream on the first shell/exec tool call). Override via env.
|
||||
// WARNING: bypass lets the agent run shell/modify FS unattended — local only.
|
||||
env.DEVIN_PERMISSION_MODE = process.env.DEVIN_PERMISSION_MODE || "bypass";
|
||||
if (mcpConfigDir) env.XDG_CONFIG_HOME = mcpConfigDir;
|
||||
|
||||
// Agent type: default (omitted) = full agent with built-in tools
|
||||
// (fs/shell/search) so the model can actually perform tasks. Override to
|
||||
// `summarizer` (no tools, text-only) via CLI_DEVIN_AGENT_TYPE for a safer,
|
||||
// tool-less mode. WARNING: the default agent can run shell commands and
|
||||
// modify the filesystem on the host running 9router — only expose locally.
|
||||
const agentType = process.env.CLI_DEVIN_AGENT_TYPE?.trim();
|
||||
const acpArgs = ["acp"];
|
||||
if (agentType) acpArgs.push("--agent-type", agentType);
|
||||
|
||||
// Spawn in the client workspace cwd (from <cwd> env context) so built-in
|
||||
// file tools create/delete relative paths in the user's project.
|
||||
// MCP config still comes from XDG_CONFIG_HOME (throwaway), not project .devin/.
|
||||
const child = spawn(devinBin, acpArgs, {
|
||||
env,
|
||||
cwd: workspaceCwd,
|
||||
stdio: ["pipe", "pipe", "pipe"],
|
||||
// On Windows, devin.exe may need shell resolution
|
||||
shell: process.platform === "win32",
|
||||
});
|
||||
|
||||
let spawnError = null;
|
||||
let stdinClosed = false;
|
||||
|
||||
child.on("error", (err) => {
|
||||
spawnError = err;
|
||||
const msg =
|
||||
err.message.includes("ENOENT") || err.message.includes("not found")
|
||||
? `Devin CLI not found: ${devinBin}. Install via https://cli.devin.ai or set CLI_DEVIN_BIN env var.`
|
||||
: `Devin CLI spawn error: ${err.message}`;
|
||||
emit(
|
||||
`data: ${JSON.stringify({ error: { message: msg, type: "devin_cli_error", code: "spawn_failed" } })}\n\n`
|
||||
);
|
||||
emit("data: [DONE]\n\n");
|
||||
controller.close();
|
||||
});
|
||||
|
||||
if (signal) {
|
||||
signal.addEventListener("abort", () => {
|
||||
if (!child.killed) child.kill("SIGTERM");
|
||||
});
|
||||
}
|
||||
|
||||
// ── JSON-RPC state machine ──────────────────────────────────────────
|
||||
let idCounter = 1;
|
||||
let sessionId = null;
|
||||
let initDone = false;
|
||||
let sessionCreated = false;
|
||||
let promptSent = false;
|
||||
const responseId = `chatcmpl-devin-${Date.now()}`;
|
||||
const created = Math.floor(Date.now() / 1000);
|
||||
let roleEmitted = false;
|
||||
let totalText = "";
|
||||
let finished = false;
|
||||
|
||||
const sendRpc = (method, params) => {
|
||||
if (stdinClosed || child.stdin.destroyed) return;
|
||||
const id = idCounter++;
|
||||
try {
|
||||
child.stdin.write(rpc(method, params, id));
|
||||
} catch {
|
||||
/* ignore write errors after close */
|
||||
}
|
||||
return id;
|
||||
};
|
||||
|
||||
// Emit a content delta as an OpenAI-compatible SSE chunk (handles the
|
||||
// leading role chunk once).
|
||||
const emitDelta = (delta) => {
|
||||
if (!roleEmitted) {
|
||||
emit(
|
||||
`data: ${JSON.stringify({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, delta: { role: "assistant", content: "" }, finish_reason: null }],
|
||||
})}\n\n`
|
||||
);
|
||||
roleEmitted = true;
|
||||
}
|
||||
totalText += delta;
|
||||
emit(
|
||||
`data: ${JSON.stringify({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, delta: { content: delta }, finish_reason: null }],
|
||||
})}\n\n`
|
||||
);
|
||||
};
|
||||
|
||||
// Emit an OpenAI tool_call delta (function calling). Ends the turn with
|
||||
// finish_reason "tool_calls" so the client executes and returns tool_result.
|
||||
let toolUseEmitted = false;
|
||||
// ACP tool_call is upsert-by-id: the first event has title, a later update
|
||||
// may only carry rawInput (title omitted). Track pending client-tool calls.
|
||||
const pendingClientTools = new Map(); // toolCallId → original tool name
|
||||
const emitToolUse = (toolName, args, toolCallId) => {
|
||||
const argsStr = typeof args === "string" ? args : JSON.stringify(args ?? {});
|
||||
if (!roleEmitted) {
|
||||
emit(
|
||||
`data: ${JSON.stringify({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, delta: { role: "assistant", content: null }, finish_reason: null }],
|
||||
})}\n\n`
|
||||
);
|
||||
roleEmitted = true;
|
||||
}
|
||||
emit(
|
||||
`data: ${JSON.stringify({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [
|
||||
{
|
||||
index: 0,
|
||||
delta: {
|
||||
tool_calls: [
|
||||
{
|
||||
index: 0,
|
||||
id: toolCallId,
|
||||
type: "function",
|
||||
function: { name: toolName, arguments: argsStr },
|
||||
},
|
||||
],
|
||||
},
|
||||
finish_reason: null,
|
||||
},
|
||||
],
|
||||
})}\n\n`
|
||||
);
|
||||
};
|
||||
|
||||
const finish = (error, finishReason = "stop") => {
|
||||
if (finished) return;
|
||||
finished = true;
|
||||
|
||||
if (error) {
|
||||
emit(
|
||||
`data: ${JSON.stringify({ error: { message: error, type: "devin_cli_error" } })}\n\n`
|
||||
);
|
||||
} else {
|
||||
// Emit finish chunk
|
||||
emit(
|
||||
`data: ${JSON.stringify({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, delta: {}, finish_reason: finishReason }],
|
||||
usage: {
|
||||
prompt_tokens: Math.ceil(promptText.length / 4),
|
||||
completion_tokens: Math.ceil(totalText.length / 4),
|
||||
total_tokens: Math.ceil((promptText.length + totalText.length) / 4),
|
||||
estimated: true,
|
||||
},
|
||||
})}\n\n`
|
||||
);
|
||||
}
|
||||
emit("data: [DONE]\n\n");
|
||||
|
||||
// Gracefully close stdin → devin will exit
|
||||
try {
|
||||
if (!stdinClosed) {
|
||||
stdinClosed = true;
|
||||
child.stdin.end();
|
||||
}
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
|
||||
// Give it 2s to exit cleanly, then SIGKILL
|
||||
const killTimer = setTimeout(() => {
|
||||
if (!child.killed) child.kill("SIGKILL");
|
||||
}, 2000);
|
||||
killTimer.unref?.();
|
||||
|
||||
controller.close();
|
||||
cleanupMcp();
|
||||
};
|
||||
|
||||
// ── stdout reader (NDJSON) ──────────────────────────────────────────
|
||||
let buffer = "";
|
||||
|
||||
child.stdout.on("data", (chunk) => {
|
||||
buffer += chunk.toString("utf8");
|
||||
let nl;
|
||||
// Each ACP message is a newline-terminated JSON line
|
||||
while ((nl = buffer.indexOf("\n")) !== -1) {
|
||||
const line = buffer.slice(0, nl).trim();
|
||||
buffer = buffer.slice(nl + 1);
|
||||
if (!line) continue;
|
||||
|
||||
let msg;
|
||||
try {
|
||||
msg = JSON.parse(line);
|
||||
} catch {
|
||||
continue; // ignore non-JSON lines (banner text, etc.)
|
||||
}
|
||||
|
||||
// ── Initialize response ───────────────────────────────────────
|
||||
if (!initDone && msg.result !== undefined && !msg.method) {
|
||||
initDone = true;
|
||||
// Create session with the client workspace cwd so agent file tools
|
||||
// resolve relative paths against the project (not /tmp).
|
||||
// `mcpServers` is required by devin 3000.2.x (must be a sequence);
|
||||
// omitting it returns -32602 "Invalid params: missing field mcpServers".
|
||||
sendRpc("session/new", {
|
||||
cwd: workspaceCwd,
|
||||
mcpServers: [],
|
||||
model: model || undefined,
|
||||
});
|
||||
continue;
|
||||
}
|
||||
|
||||
// ── session/new response → get sessionId ──────────────────────
|
||||
if (initDone && !sessionCreated && msg.result !== undefined && !msg.method) {
|
||||
const res = msg.result || {};
|
||||
sessionId = res.sessionId || null;
|
||||
if (!sessionId) {
|
||||
finish("Devin ACP: session/new returned no sessionId");
|
||||
return;
|
||||
}
|
||||
sessionCreated = true;
|
||||
// Send the prompt. devin 3000.2.x expects `prompt` (a sequence),
|
||||
// not `content` — using `content` returns -32602 "missing field prompt".
|
||||
promptSent = true;
|
||||
sendRpc("session/prompt", {
|
||||
sessionId,
|
||||
prompt: [{ type: "text", text: promptText }],
|
||||
});
|
||||
continue;
|
||||
}
|
||||
|
||||
// ── session/prompt response (ack / final result) ────────────
|
||||
if (sessionCreated && promptSent && msg.result !== undefined && !msg.method) {
|
||||
// Devin 3000.2.x only resolves session/prompt with the final result
|
||||
// (stopReason) after streaming completes. Streaming notifications are
|
||||
// handled below; nothing to do here unless we never streamed.
|
||||
if (!roleEmitted) {
|
||||
const res = msg.result || undefined;
|
||||
const content = extractResultText(res);
|
||||
if (content) {
|
||||
totalText = content;
|
||||
emitDelta(content);
|
||||
}
|
||||
const stopReason = (res && res.stopReason) || "";
|
||||
if (stopReason && stopReason !== "cancelled") {
|
||||
finish();
|
||||
return;
|
||||
}
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
// ── Permission requests → auto-approve the first allow option ──
|
||||
// Devi asks before running shell/exec tools; as a headless proxy we
|
||||
// grant once. (DEVIN_PERMISSION_MODE=bypass usually prevents these,
|
||||
// but some tool kinds still prompt, so handle them here too.)
|
||||
if (msg.method === "session/request_permission" && msg.id !== undefined) {
|
||||
const options = msg.params?.options || [];
|
||||
const allow =
|
||||
options.find((o) => /allow/i.test(String(o.kind || ""))) || options[0];
|
||||
if (allow) {
|
||||
child.stdin.write(
|
||||
JSON.stringify({
|
||||
jsonrpc: "2.0",
|
||||
id: msg.id,
|
||||
result: { outcome: { outcome: "selected", optionId: allow.optionId } },
|
||||
}) + "\n"
|
||||
);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
// ── Agent stopped notification (devin 3000.2.x stop signal) ───
|
||||
if (msg.method === "_cognition.ai/agent_stopped" || msg.method === "$/agent_stopped") {
|
||||
const cause = msg.params?.cause;
|
||||
if (cause === "error") {
|
||||
// devin uses errorMessage on this notification (not message/error).
|
||||
const errText =
|
||||
msg.params?.errorMessage ||
|
||||
msg.params?.message ||
|
||||
msg.params?.error ||
|
||||
"Devin agent error";
|
||||
finish(String(errText));
|
||||
} else {
|
||||
finish();
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
// ── Streaming notifications (session/update) ──────────────────
|
||||
if (msg.method === "session/update" || msg.method === "$/update") {
|
||||
const params = msg.params;
|
||||
if (!params) continue;
|
||||
|
||||
// devin 3000.2.x nests the payload under params.update.sessionUpdate;
|
||||
// older devin used a flat params.type.
|
||||
const update = params.update || {};
|
||||
const type = update.sessionUpdate || params.type;
|
||||
const contentField = update.content !== undefined ? update.content : params.content;
|
||||
const deltaText =
|
||||
typeof contentField === "string"
|
||||
? contentField
|
||||
: contentField?.text ?? params.delta ?? params.text ?? "";
|
||||
|
||||
// ── Client-tool bridge: devin calling a tool from our exposed MCP ──
|
||||
// ACP title shape: "Calling mcp_<name> from clientTools".
|
||||
// tool_call is upsert-by-id: title may only appear on the first event,
|
||||
// rawInput on a later tool_call_update. Track pending ids so we don't
|
||||
// require both fields on the same notification.
|
||||
if (
|
||||
hasClientTools &&
|
||||
!toolUseEmitted &&
|
||||
(type === "tool_call" || type === "tool_call_update")
|
||||
) {
|
||||
const tcId = update.toolCallId;
|
||||
if (typeof update.title === "string" && update.title.startsWith("Calling mcp_") && /from clientTools\b/.test(update.title)) {
|
||||
const nameMatch = update.title.match(/^Calling (mcp_\S+)\b/);
|
||||
const mcpName = nameMatch ? nameMatch[1] : "";
|
||||
const origName = fromMcpToolName(mcpName);
|
||||
if (tcId && origName) pendingClientTools.set(tcId, origName);
|
||||
}
|
||||
const origName = tcId ? pendingClientTools.get(tcId) : null;
|
||||
if (origName && update.rawInput) {
|
||||
toolUseEmitted = true;
|
||||
pendingClientTools.delete(tcId);
|
||||
emitToolUse(origName, update.rawInput, tcId || `call_${Date.now()}`);
|
||||
finish(null, "tool_calls");
|
||||
return;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
if (type === "agent_message_chunk" || type === "message_delta" || type === "text_delta" || type === "content_delta") {
|
||||
if (deltaText) emitDelta(deltaText);
|
||||
} else if (type === "agent_thought_chunk") {
|
||||
// Internal reasoning — not surfaced to the client.
|
||||
} else if (type === "message_stop" || type === "stop" || type === "done") {
|
||||
finish();
|
||||
return;
|
||||
} else if (type === "error") {
|
||||
finish(String(params.message || params.error || "Devin ACP error"));
|
||||
return;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
// ── Error responses ───────────────────────────────────────────
|
||||
if (msg.error) {
|
||||
finish(`Devin ACP error ${msg.error.code}: ${msg.error.message}`);
|
||||
return;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
child.stderr.on("data", (chunk) => {
|
||||
log?.debug?.("DEVIN", `stderr: ${chunk.toString("utf8").slice(0, 200)}`);
|
||||
});
|
||||
|
||||
child.on("close", (code) => {
|
||||
if (!finished) {
|
||||
if (code !== 0 && !spawnError) {
|
||||
finish(roleEmitted ? undefined : `Devin CLI exited with code ${code}`);
|
||||
} else {
|
||||
finish();
|
||||
}
|
||||
} else {
|
||||
cleanupMcp();
|
||||
}
|
||||
});
|
||||
|
||||
// ── Send initialize ───────────────────────────────────────────────
|
||||
sendRpc("initialize", {
|
||||
protocolVersion: "0.3",
|
||||
clientInfo: { name: "9router", version: "1.0" },
|
||||
capabilities: {},
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
return {
|
||||
response: new Response(sseStream, {
|
||||
status: 200,
|
||||
headers: {
|
||||
"Content-Type": "text/event-stream",
|
||||
"Cache-Control": "no-cache",
|
||||
Connection: "keep-alive",
|
||||
},
|
||||
}),
|
||||
url: "devin://acp/stdio",
|
||||
headers: {},
|
||||
transformedBody: {
|
||||
model,
|
||||
cwd: workspaceCwd,
|
||||
clientTools: clientTools.map((t) => t?.function?.name || t?.name).filter(Boolean),
|
||||
clientToolResults: Object.keys(clientToolResults),
|
||||
mcpServers: Object.keys(mcpServers),
|
||||
promptLength: Array.isArray(body?.messages)
|
||||
? body.messages.length
|
||||
: Array.isArray(body?.input)
|
||||
? body.input.length
|
||||
: 0,
|
||||
},
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
// ─── Helpers ─────────────────────────────────────────────────────────────────
|
||||
|
||||
// Extract text from a final ACP session/prompt result object across common shapes.
|
||||
function extractResultText(result) {
|
||||
// { message: { content: "..." } }
|
||||
// { messages: [{ content: "..." }] }
|
||||
// { content: "..." }
|
||||
// { text: "..." }
|
||||
if (typeof result.content === "string") return result.content;
|
||||
if (typeof result.text === "string") return result.text;
|
||||
const msg = result.message;
|
||||
if (msg && typeof msg.content === "string") return msg.content;
|
||||
const msgs = result.messages;
|
||||
if (Array.isArray(msgs)) {
|
||||
return msgs
|
||||
.filter((m) => m.role === "assistant")
|
||||
.map((m) => String(m.content || ""))
|
||||
.join("\n");
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
export default DevinCliExecutor;
|
||||
@@ -20,7 +20,12 @@ import { CommandCodeExecutor } from "./commandcode.js";
|
||||
import { XiaomiTokenplanExecutor } from "./xiaomi-tokenplan.js";
|
||||
import { MimoFreeExecutor } from "./mimo-free.js";
|
||||
import { CodeBuddyExecutor } from "./codebuddy-cn.js";
|
||||
import { CodeBuddyIntlExecutor } from "./codebuddy-intl.js";
|
||||
import TraeExecutor from "./trae.js";
|
||||
import ZedExecutor from "./zed.js";
|
||||
import WindsurfExecutor from "./windsurf.js";
|
||||
import { DefaultExecutor } from "./default.js";
|
||||
import { DevinCliExecutor } from "./devin-cli.js";
|
||||
|
||||
const executors = {
|
||||
antigravity: new AntigravityExecutor(),
|
||||
@@ -50,6 +55,11 @@ const executors = {
|
||||
"mimo-free": new MimoFreeExecutor(),
|
||||
mmf: new MimoFreeExecutor(), // Alias for mimo-free
|
||||
"codebuddy-cn": new CodeBuddyExecutor(),
|
||||
"codebuddy-intl": new CodeBuddyIntlExecutor(),
|
||||
trae: new TraeExecutor(),
|
||||
zed: new ZedExecutor(),
|
||||
windsurf: new WindsurfExecutor(),
|
||||
"devin-cli": new DevinCliExecutor(),
|
||||
};
|
||||
|
||||
const defaultCache = new Map();
|
||||
@@ -88,3 +98,8 @@ export { CommandCodeExecutor } from "./commandcode.js";
|
||||
export { XiaomiTokenplanExecutor } from "./xiaomi-tokenplan.js";
|
||||
export { MimoFreeExecutor } from "./mimo-free.js";
|
||||
export { CodeBuddyExecutor } from "./codebuddy-cn.js";
|
||||
export { CodeBuddyIntlExecutor } from "./codebuddy-intl.js";
|
||||
export { default as TraeExecutor } from "./trae.js";
|
||||
export { default as ZedExecutor } from "./zed.js";
|
||||
export { default as WindsurfExecutor } from "./windsurf.js";
|
||||
export { DevinCliExecutor } from "./devin-cli.js";
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -33,7 +33,11 @@ import { FETCH_CONNECT_TIMEOUT_MS } from "../config/runtimeConfig.js";
|
||||
import { resolveProviderTimeoutMs } from "../services/providerTimeout.js";
|
||||
import {
|
||||
QODER_CHAT_URL_ENCODED,
|
||||
QODER_JOB_TOKEN_EXCHANGE_URL,
|
||||
QODER_USERINFO_URL,
|
||||
QODER_MODEL_MAP,
|
||||
QODER_IDE_VERSION,
|
||||
QODER_CLIENT_TYPE,
|
||||
} from "../shared/qoder/constants.js";
|
||||
import { getQoderModelConfig, resolveQoderModels } from "../services/qoderModels.js";
|
||||
|
||||
@@ -221,8 +225,13 @@ async function buildQoderRequestBody({ model, body, credentials, log, proxyOptio
|
||||
* 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 `data: [DONE]\n\n` plus
|
||||
* a synthetic OpenAI error chunk.
|
||||
* and re-emit as `data: <inner>\n\n`. Errors become a synthetic OpenAI error
|
||||
* chunk + [DONE].
|
||||
*
|
||||
* Critical: Qoder's SSE often keeps the socket open after the terminal
|
||||
* [DONE]/error frame (agent keepalive). Non-streaming clients drain via
|
||||
* response.text() which hangs until the socket closes — so on terminal
|
||||
* events we cancel the upstream reader and close our stream immediately.
|
||||
*/
|
||||
function wrapQoderSSE(response, model) {
|
||||
if (!response.ok || !response.body) return response;
|
||||
@@ -231,15 +240,14 @@ function wrapQoderSSE(response, model) {
|
||||
const encoder = new TextEncoder();
|
||||
let buffer = "";
|
||||
let doneEmitted = false;
|
||||
const reader = response.body.getReader();
|
||||
|
||||
// Process one already-extracted SSE line (no trailing newline). Returns
|
||||
// false when the line indicated end-of-stream so the caller can stop
|
||||
// forwarding any remaining chunks after [DONE].
|
||||
// Process one already-extracted SSE line (no trailing newline).
|
||||
const processLine = (line, controller) => {
|
||||
const trimmed = line.replace(/\r$/, "").trim();
|
||||
if (!trimmed) return;
|
||||
if (!trimmed.startsWith("data:")) return;
|
||||
if (doneEmitted) return; // never forward chunks past stream end
|
||||
if (doneEmitted) return;
|
||||
|
||||
const data = trimmed.slice(5).trimStart();
|
||||
if (data === "[DONE]") {
|
||||
@@ -272,47 +280,60 @@ function wrapQoderSSE(response, model) {
|
||||
doneEmitted = true;
|
||||
return;
|
||||
}
|
||||
// Inner is an OpenAI-shaped chunk. Strip any embedded newlines so the
|
||||
// SSE frame stays a single event (a literal "\n" inside `inner` would
|
||||
// otherwise split the frame across multiple data: lines and downstream
|
||||
// parsers would reassemble them as separate events).
|
||||
// Strip embedded newlines so the SSE frame stays a single event.
|
||||
const sanitized = inner.replace(/\r?\n/g, "");
|
||||
controller.enqueue(encoder.encode(`data: ${sanitized}\n\n`));
|
||||
};
|
||||
|
||||
const transform = new TransformStream({
|
||||
transform(chunk, controller) {
|
||||
buffer += decoder.decode(chunk, { stream: true });
|
||||
let nl;
|
||||
while ((nl = buffer.indexOf("\n")) !== -1) {
|
||||
const line = buffer.slice(0, nl);
|
||||
buffer = buffer.slice(nl + 1);
|
||||
processLine(line, controller);
|
||||
const stream = new ReadableStream({
|
||||
// Use start()+loop (not pull): a pull that buffers a partial line without
|
||||
// enqueueing would never be re-invoked, hanging consumers like .text().
|
||||
async start(controller) {
|
||||
try {
|
||||
while (!doneEmitted) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) {
|
||||
buffer += decoder.decode();
|
||||
if (buffer.length > 0) {
|
||||
processLine(buffer, controller);
|
||||
buffer = "";
|
||||
}
|
||||
break;
|
||||
}
|
||||
|
||||
buffer += decoder.decode(value, { stream: true });
|
||||
let nl;
|
||||
while ((nl = buffer.indexOf("\n")) !== -1) {
|
||||
const line = buffer.slice(0, nl);
|
||||
buffer = buffer.slice(nl + 1);
|
||||
processLine(line, controller);
|
||||
if (doneEmitted) {
|
||||
// Terminal frame received — drop upstream keepalive and end.
|
||||
await reader.cancel().catch(() => {});
|
||||
controller.close();
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
// fall through to terminal [DONE] + close
|
||||
} finally {
|
||||
if (!doneEmitted) {
|
||||
try {
|
||||
controller.enqueue(encoder.encode(SSE_DONE));
|
||||
doneEmitted = true;
|
||||
} catch { /* already closed */ }
|
||||
}
|
||||
try { controller.close(); } catch { /* already closed */ }
|
||||
await reader.cancel().catch(() => {});
|
||||
}
|
||||
},
|
||||
flush(controller) {
|
||||
// Finalize the decoder so any pending multi-byte sequence is
|
||||
// released into `buffer` instead of being silently dropped.
|
||||
buffer += decoder.decode();
|
||||
// Drain any trailing line that arrived without a terminating newline
|
||||
// (e.g. upstream closed the socket immediately after the last write,
|
||||
// or a CDN stripped the final CRLF). Without this, the chunk that
|
||||
// carries finish_reason is silently lost.
|
||||
if (buffer.length > 0) {
|
||||
processLine(buffer, controller);
|
||||
buffer = "";
|
||||
}
|
||||
if (!doneEmitted) {
|
||||
controller.enqueue(encoder.encode(SSE_DONE));
|
||||
doneEmitted = true;
|
||||
}
|
||||
cancel() {
|
||||
return reader.cancel().catch(() => {});
|
||||
},
|
||||
});
|
||||
|
||||
const transformed = response.body.pipeThrough(transform);
|
||||
// Build a Response with passable headers; the streaming handler reads
|
||||
// `.body` as a ReadableStream regardless of Content-Type.
|
||||
return new Response(transformed, {
|
||||
return new Response(stream, {
|
||||
status: response.status,
|
||||
statusText: response.statusText,
|
||||
headers: {
|
||||
@@ -322,6 +343,92 @@ function wrapQoderSSE(response, model) {
|
||||
});
|
||||
}
|
||||
|
||||
// ── PAT (Personal Access Token) → job-token exchange ───────────────────────
|
||||
// PATs (pt-...) cannot sign COSY requests directly. Exchange them for a
|
||||
// short-lived job token (jt-...) via /api/v1/jobToken/exchange (plain JSON,
|
||||
// not COSY-signed), then resolve the userId from userinfo. Mirrors the
|
||||
// official qodercli flow. Cached per-PAT until near-expiry.
|
||||
const PAT_PREFIX = "pt-";
|
||||
const PAT_REFRESH_BUFFER_MS = 5 * 60 * 1000;
|
||||
const patJobCache = new Map();
|
||||
|
||||
export function isQoderPat(token) {
|
||||
return typeof token === "string" && token.startsWith(PAT_PREFIX);
|
||||
}
|
||||
|
||||
async function exchangeJobToken(pat, proxyOptions = null, signal = null) {
|
||||
const res = await proxyAwareFetch(
|
||||
QODER_JOB_TOKEN_EXCHANGE_URL,
|
||||
{
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/json",
|
||||
"User-Agent": "qodercli/1.0.0",
|
||||
"Cosy-Version": QODER_IDE_VERSION,
|
||||
"Cosy-ClientType": QODER_CLIENT_TYPE,
|
||||
},
|
||||
body: JSON.stringify({ personal_token: pat }),
|
||||
signal,
|
||||
},
|
||||
proxyOptions,
|
||||
);
|
||||
if (!res.ok) {
|
||||
const text = await res.text().catch(() => "");
|
||||
throw new Error(`qoder PAT exchange failed: ${res.status} ${text.slice(0, 200)}`);
|
||||
}
|
||||
const data = await res.json();
|
||||
if (!data.token) throw new Error("qoder PAT exchange returned no job token");
|
||||
|
||||
let expiresAt = Date.now() + 24 * 60 * 60 * 1000;
|
||||
if (data.expires_at) {
|
||||
const parsed = Date.parse(data.expires_at);
|
||||
if (!Number.isNaN(parsed)) expiresAt = parsed;
|
||||
} else if (typeof data.expires_in === "number" && data.expires_in > 0) {
|
||||
expiresAt = Date.now() + data.expires_in;
|
||||
}
|
||||
return { jobToken: data.token, jobRefreshToken: data.refresh_token || "", expiresAt };
|
||||
}
|
||||
|
||||
async function fetchUserIdForJobToken(jobToken, proxyOptions = null, signal = null) {
|
||||
try {
|
||||
const res = await proxyAwareFetch(
|
||||
QODER_USERINFO_URL,
|
||||
{
|
||||
method: "GET",
|
||||
headers: {
|
||||
Authorization: `Bearer ${jobToken}`,
|
||||
Accept: "application/json",
|
||||
"User-Agent": "qodercli/1.0.0",
|
||||
},
|
||||
signal,
|
||||
},
|
||||
proxyOptions,
|
||||
);
|
||||
if (!res.ok) return "";
|
||||
const info = await res.json().catch(() => ({}));
|
||||
return info.id || info.userId || info.user_id || "";
|
||||
} catch {
|
||||
return "";
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Exchange a PAT for a job token + userId, caching until near-expiry so repeat
|
||||
* chat requests don't re-exchange. Returns { accessToken, userId }.
|
||||
*/
|
||||
async function resolvePatCredential(pat, proxyOptions = null, signal = null) {
|
||||
const cached = patJobCache.get(pat);
|
||||
if (cached && cached.expiresAt - Date.now() > PAT_REFRESH_BUFFER_MS) {
|
||||
return cached;
|
||||
}
|
||||
const { jobToken, expiresAt } = await exchangeJobToken(pat, proxyOptions, signal);
|
||||
const userId = await fetchUserIdForJobToken(jobToken, proxyOptions, signal);
|
||||
const entry = { accessToken: jobToken, userId, expiresAt };
|
||||
patJobCache.set(pat, entry);
|
||||
return entry;
|
||||
}
|
||||
|
||||
export class QoderExecutor extends BaseExecutor {
|
||||
constructor() {
|
||||
super("qoder", PROVIDERS.qoder);
|
||||
@@ -339,6 +446,34 @@ export class QoderExecutor extends BaseExecutor {
|
||||
async execute({ model, body, stream, credentials, signal, log, proxyOptions = null }) {
|
||||
const url = this.buildUrl();
|
||||
|
||||
// PAT (pt-...) → exchange for short-lived job token + resolve userId so
|
||||
// downstream COSY signing + catalog fetch work. Device tokens (dt-...) and
|
||||
// job tokens (jt-...) skip this and are used directly.
|
||||
const rawToken = credentials?.apiKey || credentials?.accessToken;
|
||||
if (isQoderPat(rawToken)) {
|
||||
try {
|
||||
const resolved = await resolvePatCredential(rawToken, proxyOptions, signal);
|
||||
credentials = {
|
||||
...credentials,
|
||||
accessToken: resolved.accessToken,
|
||||
apiKey: undefined,
|
||||
providerSpecificData: {
|
||||
authMethod: "pat",
|
||||
...(credentials?.providerSpecificData || {}),
|
||||
userId: resolved.userId || credentials?.providerSpecificData?.userId || "",
|
||||
machineId: credentials?.providerSpecificData?.machineId || "",
|
||||
},
|
||||
};
|
||||
} catch (err) {
|
||||
log?.error?.("QODER", `PAT exchange failed: ${err.message}`);
|
||||
const fakeResp = new Response(
|
||||
JSON.stringify({ error: { message: `qoder PAT exchange failed: ${err.message}` } }),
|
||||
{ status: 401, headers: { "Content-Type": "application/json" } },
|
||||
);
|
||||
return { response: fakeResp, url, headers: {}, transformedBody: body };
|
||||
}
|
||||
}
|
||||
|
||||
const psd = credentials?.providerSpecificData || {};
|
||||
if (!psd.userId) {
|
||||
// No user id → no way to sign. Surface a 401 so the dashboard nudges
|
||||
@@ -456,4 +591,6 @@ export const __test__ = {
|
||||
normalizeMessages,
|
||||
wrapQoderSSE,
|
||||
buildQoderRequestBody,
|
||||
isQoderPat,
|
||||
resolvePatCredential,
|
||||
};
|
||||
|
||||
339
open-sse/executors/trae.js
Normal file
339
open-sse/executors/trae.js
Normal file
@@ -0,0 +1,339 @@
|
||||
import { BaseExecutor } from "./base.js";
|
||||
import { proxyAwareFetch } from "../utils/proxyFetch.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
|
||||
// Trae executor — SOLO remote agent API.
|
||||
//
|
||||
// Flow:
|
||||
// 1. POST {base}/chat_sessions → { code:0, data:{ chat_session_id, message_id } }
|
||||
// 2. GET {base}/chat_sessions/{id}/events?reply_to_message_id={message_id}
|
||||
// → text/event-stream. Assistant text streams in `plan_item` events under
|
||||
// the `thought` field (cumulative per plan-item id). `token_usage` carries
|
||||
// usage; `done` ends the turn; `error` carries upstream errors.
|
||||
//
|
||||
// Auth: header `Authorization: Cloud-IDE-JWT <jwt>` (RS256, ~14-day lifetime).
|
||||
// Identity fields for common_params live in credentials.providerSpecificData.
|
||||
|
||||
const STREAM_TIMEOUT_MS = parseInt(process.env.TRAE_STREAM_TIMEOUT_MS || "300000", 10);
|
||||
const TRAE_UA =
|
||||
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 " +
|
||||
"(KHTML, like Gecko) Chrome/149.0.0.0 Safari/537.36";
|
||||
|
||||
function flattenQuery(messages) {
|
||||
const parts = [];
|
||||
for (const m of messages) {
|
||||
let content = "";
|
||||
if (typeof m.content === "string") content = m.content;
|
||||
else if (Array.isArray(m.content)) {
|
||||
content = m.content
|
||||
.map((p) => {
|
||||
if (typeof p === "string") return p;
|
||||
if (p && typeof p === "object") return String(p.text ?? "");
|
||||
return "";
|
||||
})
|
||||
.join("");
|
||||
}
|
||||
if (m.role === "system") parts.push(`[System]\n${content}`);
|
||||
else if (m.role === "assistant") parts.push(`[Assistant]\n${content}`);
|
||||
else parts.push(content);
|
||||
}
|
||||
// Trae expects query as a JSON-encoded string of typed content blocks.
|
||||
return JSON.stringify([{ type: "text", data: { content: parts.join("\n\n") } }]);
|
||||
}
|
||||
|
||||
export default class TraeExecutor extends BaseExecutor {
|
||||
constructor() {
|
||||
super("trae", PROVIDERS.trae);
|
||||
}
|
||||
|
||||
base() {
|
||||
return (this.config.baseUrl || "https://core-normal.trae.ai/api/remote/v1").replace(/\/$/, "");
|
||||
}
|
||||
|
||||
buildHeaders(credentials, stream = true) {
|
||||
const token = credentials?.accessToken || "";
|
||||
const psd = credentials?.providerSpecificData || {};
|
||||
return {
|
||||
Authorization: `Cloud-IDE-JWT ${token}`,
|
||||
"Content-Type": "application/json",
|
||||
"X-Trae-Client-Type": "web",
|
||||
"X-Preferenced-Language": psd.appLanguage || "en",
|
||||
"x-user-region": psd.userRegion || "US",
|
||||
Referer: "https://solo.trae.ai/",
|
||||
"User-Agent": TRAE_UA,
|
||||
Accept: stream ? "text/event-stream" : "application/json",
|
||||
};
|
||||
}
|
||||
|
||||
// SOLO session modes: "code" (model picker) vs "work" (fast auto lane).
|
||||
resolveMode(model) {
|
||||
const m = (model || "").trim().toLowerCase();
|
||||
if (m === "work" || m === "auto-work" || m === "solo-work") {
|
||||
return { mode: "work", strategy: "auto", modelName: "" };
|
||||
}
|
||||
const auto = !m || m === "auto";
|
||||
return { mode: "code", strategy: auto ? "auto" : "manual", modelName: auto ? "" : model };
|
||||
}
|
||||
|
||||
// common_params is a JSON-encoded string embedded inside initial_message.
|
||||
commonParams(psd, mode, sessionId) {
|
||||
const cp = {
|
||||
language: "en-us",
|
||||
app_language: psd.appLanguage || "en",
|
||||
quality: "stable",
|
||||
app_version: psd.appVersion || "1.0.0.1229",
|
||||
web_id: psd.webId || "",
|
||||
user_identity: psd.userIdentity || "Free",
|
||||
is_freshman: "0",
|
||||
biz_user_id: psd.bizUserId || "",
|
||||
user_unique_id: psd.userUniqueId || "",
|
||||
scope: psd.scope || "marscode-us",
|
||||
tenant: psd.tenant || "marscode",
|
||||
region: psd.region || "US-East",
|
||||
aiRegion: psd.aiRegion || psd.region || "US-East",
|
||||
is_privacy_mode: 0,
|
||||
privacy_mode: "off",
|
||||
solo_chat_mode: mode,
|
||||
};
|
||||
if (sessionId) cp.biz_session_id = sessionId;
|
||||
return JSON.stringify(cp);
|
||||
}
|
||||
|
||||
// POST /chat_sessions — creates a session and submits the first turn.
|
||||
async createSession(headers, query, model, psd, signal) {
|
||||
const { mode, strategy, modelName } = this.resolveMode(model);
|
||||
const body = {
|
||||
mode,
|
||||
environment_id: "default",
|
||||
initial_message: {
|
||||
chat_session_id: "",
|
||||
content: [],
|
||||
query,
|
||||
model_name: modelName,
|
||||
agent_type: "solo_agent_remote",
|
||||
model_selection_strategy: strategy,
|
||||
common_params: this.commonParams(psd, mode),
|
||||
},
|
||||
env: "remote",
|
||||
auto_create_project: false,
|
||||
origin: "web",
|
||||
};
|
||||
const res = await proxyAwareFetch(`${this.base()}/chat_sessions`, {
|
||||
method: "POST",
|
||||
headers,
|
||||
body: JSON.stringify(body),
|
||||
signal,
|
||||
}, null);
|
||||
const text = await res.text();
|
||||
if (!res.ok) throw new Error(`[${res.status}] ${text}`);
|
||||
const json = JSON.parse(text);
|
||||
if (json?.code !== 0) throw new Error(`Trae create_session: ${JSON.stringify(json)}`);
|
||||
return { sessionId: json.data.chat_session_id, messageId: json.data.message_id };
|
||||
}
|
||||
|
||||
// GET /events SSE → invoke onEvent(eventType, dataObj) per frame.
|
||||
// Resolves when `done`/`error` arrives, the stream ends, or timeout fires.
|
||||
async streamEvents(headers, sessionId, replyTo, onEvent, signal) {
|
||||
const url = `${this.base()}/chat_sessions/${sessionId}/events?reply_to_message_id=${encodeURIComponent(replyTo)}`;
|
||||
const ctrl = new AbortController();
|
||||
if (signal?.aborted) ctrl.abort();
|
||||
const timer = setTimeout(() => ctrl.abort(new Error("trae stream timeout")), STREAM_TIMEOUT_MS);
|
||||
const onAbort = () => ctrl.abort();
|
||||
if (signal) signal.addEventListener("abort", onAbort, { once: true });
|
||||
try {
|
||||
const res = await proxyAwareFetch(url, { method: "GET", headers, signal: ctrl.signal }, null);
|
||||
if (!res.ok || !res.body) throw new Error(`[${res.status}] events stream failed`);
|
||||
const reader = res.body.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
let buf = "";
|
||||
let ev = null;
|
||||
for (;;) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
buf += decoder.decode(value, { stream: true });
|
||||
let nl;
|
||||
while ((nl = buf.indexOf("\n")) >= 0) {
|
||||
const line = buf.slice(0, nl).replace(/\r$/, "");
|
||||
buf = buf.slice(nl + 1);
|
||||
if (line.startsWith("event:")) ev = line.slice(6).trim();
|
||||
else if (line.startsWith("data:")) {
|
||||
const payload = line.slice(5).trim();
|
||||
let data;
|
||||
try { data = JSON.parse(payload); } catch { data = { _raw: payload }; }
|
||||
if (onEvent(ev, data)) {
|
||||
await reader.cancel().catch(() => {});
|
||||
return;
|
||||
}
|
||||
} else if (line === "") ev = null;
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
if (signal) signal.removeEventListener("abort", onAbort);
|
||||
}
|
||||
}
|
||||
|
||||
async execute({ model, body, stream, credentials, signal }) {
|
||||
const headers = this.buildHeaders(credentials, stream !== false);
|
||||
const psd = credentials?.providerSpecificData || {};
|
||||
const query = flattenQuery(body?.messages || []);
|
||||
const responseId = `chatcmpl-trae-${Date.now()}`;
|
||||
const created = Math.floor(Date.now() / 1000);
|
||||
|
||||
const errResponse = (status, message) => new Response(
|
||||
JSON.stringify({ error: { message, type: "api_error", code: "" } }),
|
||||
{ status, headers: { "Content-Type": "application/json" } }
|
||||
);
|
||||
|
||||
let session;
|
||||
try {
|
||||
session = await this.createSession(headers, query, model, psd, signal);
|
||||
} catch (err) {
|
||||
return { response: errResponse(502, err?.message ? String(err.message) : String(err)), url: this.base(), headers, transformedBody: body };
|
||||
}
|
||||
|
||||
// Shared per-turn state: plan_item thoughts (cumulative, longest wins).
|
||||
const order = [];
|
||||
const thoughts = {};
|
||||
let sent = 0;
|
||||
let usage = null;
|
||||
let errorEvent = null;
|
||||
const renderNewText = (data) => {
|
||||
const pid = data.id;
|
||||
if (!pid) return "";
|
||||
if (!(pid in thoughts)) order.push(pid);
|
||||
const t = data.thought || "";
|
||||
if (t.length >= (thoughts[pid] || "").length) thoughts[pid] = t;
|
||||
const full = order.map((i) => thoughts[i]).join("");
|
||||
const piece = full.slice(sent);
|
||||
sent = full.length;
|
||||
return piece;
|
||||
};
|
||||
|
||||
if (stream !== false) {
|
||||
const enc = new TextEncoder();
|
||||
const sse = new ReadableStream({
|
||||
start: async (controller) => {
|
||||
const emit = (obj) => controller.enqueue(enc.encode(`data: ${JSON.stringify(obj)}\n\n`));
|
||||
emit({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, delta: { role: "assistant" }, finish_reason: null }],
|
||||
});
|
||||
try {
|
||||
await this.streamEvents(headers, session.sessionId, session.messageId, (ev, data) => {
|
||||
if (ev === "error") { errorEvent = data; return true; }
|
||||
if (ev === "token_usage") usage = data;
|
||||
if (ev === "plan_item") {
|
||||
const piece = renderNewText(data);
|
||||
if (piece) {
|
||||
emit({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, delta: { content: piece }, finish_reason: null }],
|
||||
});
|
||||
}
|
||||
}
|
||||
return ev === "done";
|
||||
}, signal);
|
||||
if (errorEvent) {
|
||||
emit({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [],
|
||||
error: { message: `trae ${errorEvent.code || ""}: ${errorEvent.message || ""}`, type: "api_error" },
|
||||
});
|
||||
} else {
|
||||
emit({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
|
||||
});
|
||||
if (usage) {
|
||||
emit({
|
||||
id: responseId,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [],
|
||||
usage: {
|
||||
prompt_tokens: usage.prompt_tokens || 0,
|
||||
completion_tokens: usage.completion_tokens || 0,
|
||||
total_tokens: usage.total_tokens || 0,
|
||||
},
|
||||
});
|
||||
}
|
||||
}
|
||||
controller.enqueue(enc.encode("data: [DONE]\n\n"));
|
||||
controller.close();
|
||||
} catch (err) {
|
||||
controller.error(err);
|
||||
}
|
||||
},
|
||||
});
|
||||
return {
|
||||
response: new Response(sse, {
|
||||
status: 200,
|
||||
headers: {
|
||||
"Content-Type": "text/event-stream",
|
||||
"Cache-Control": "no-cache",
|
||||
"Connection": "keep-alive",
|
||||
},
|
||||
}),
|
||||
url: this.base(),
|
||||
headers,
|
||||
transformedBody: body,
|
||||
};
|
||||
}
|
||||
|
||||
// Non-streaming: drive to completion, return chat.completion JSON.
|
||||
try {
|
||||
await this.streamEvents(headers, session.sessionId, session.messageId, (ev, data) => {
|
||||
if (ev === "error") { errorEvent = data; return true; }
|
||||
if (ev === "token_usage") usage = data;
|
||||
if (ev === "plan_item") renderNewText(data);
|
||||
return ev === "done";
|
||||
}, signal);
|
||||
} catch (err) {
|
||||
return { response: errResponse(502, err?.message ? String(err.message) : String(err)), url: this.base(), headers, transformedBody: body };
|
||||
}
|
||||
if (errorEvent) {
|
||||
return { response: errResponse(502, `trae ${errorEvent.code || ""}: ${errorEvent.message || ""}`), url: this.base(), headers, transformedBody: body };
|
||||
}
|
||||
const content = order.map((i) => thoughts[i]).join("");
|
||||
const out = {
|
||||
id: responseId,
|
||||
object: "chat.completion",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, message: { role: "assistant", content }, finish_reason: "stop" }],
|
||||
};
|
||||
if (usage) {
|
||||
out.usage = {
|
||||
prompt_tokens: usage.prompt_tokens || 0,
|
||||
completion_tokens: usage.completion_tokens || 0,
|
||||
total_tokens: usage.total_tokens || 0,
|
||||
};
|
||||
}
|
||||
return {
|
||||
response: new Response(JSON.stringify(out), { status: 200, headers: { "Content-Type": "application/json" } }),
|
||||
url: this.base(),
|
||||
headers,
|
||||
transformedBody: body,
|
||||
};
|
||||
}
|
||||
|
||||
// Refresh hook placeholder — Cloud-IDE-JWT is long-lived (~14d); refresh via
|
||||
// ExchangeToken (refresh→access) is wired in services/tokenRefresh/providers.js.
|
||||
async refreshCredentials() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
588
open-sse/executors/windsurf.js
Normal file
588
open-sse/executors/windsurf.js
Normal file
@@ -0,0 +1,588 @@
|
||||
import { BaseExecutor } from "./base.js";
|
||||
import { proxyAwareFetch } from "../utils/proxyFetch.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
import { randomUUID } from "node:crypto";
|
||||
|
||||
// WindsurfExecutor — Codeium gRPC-web chat.
|
||||
//
|
||||
// Wire protocol: gRPC-web over HTTPS (Content-Type: application/grpc-web+proto).
|
||||
// Service: exa.language_server_pb.LanguageServerService
|
||||
// Method: GetChatMessage (unary request → streamed CompletionChunk frames)
|
||||
//
|
||||
// Auth: credentials.accessToken = Codeium apiKey (sk-ws-... or Firebase-derived)
|
||||
// — placed in Metadata.api_key protobuf field of every request + Bearer header.
|
||||
|
||||
const WS_BASE_URL = "https://server.codeium.com";
|
||||
const WS_SERVICE = "exa.language_server_pb.LanguageServerService";
|
||||
const WS_METHOD_CHAT = "GetChatMessage";
|
||||
const WS_CHAT_URL = `${WS_BASE_URL}/${WS_SERVICE}/${WS_METHOD_CHAT}`;
|
||||
|
||||
const WS_IDE_NAME = "windsurf";
|
||||
const WS_IDE_VERSION = "3.14.0";
|
||||
const WS_EXT_VERSION = "3.14.0";
|
||||
const WS_LOCALE = "en-US";
|
||||
|
||||
// ─── Model alias map (catalog name → Windsurf wire name) ─────────────────────
|
||||
const MODEL_ALIAS_MAP = {
|
||||
// ── Cognition SWE ───────────────────────────────────────────────────────
|
||||
"swe-1.6-fast": "swe-1-6-fast",
|
||||
"swe-1.6": "swe-1-6",
|
||||
"swe-1.5-fast": "swe-1-5-fast",
|
||||
"swe-1.5": "swe-1-5",
|
||||
// ── Claude Opus 4.7 — effort-tiered ─────────────────────────────────────
|
||||
"claude-opus-4.7-max": "claude-opus-4-7-max",
|
||||
"claude-opus-4.7-xhigh": "claude-opus-4-7-xhigh",
|
||||
"claude-opus-4.7-high": "claude-opus-4-7-high",
|
||||
"claude-opus-4.7-medium": "claude-opus-4-7-medium",
|
||||
"claude-opus-4.7-low": "claude-opus-4-7-low",
|
||||
"claude-opus-4.7-review": "opus-4-7-review",
|
||||
// ── Claude Opus/Sonnet 4.6 ──────────────────────────────────────────────
|
||||
"claude-sonnet-4.6-thinking-1m": "claude-sonnet-4-6-thinking-1m",
|
||||
"claude-sonnet-4.6-1m": "claude-sonnet-4-6-1m",
|
||||
"claude-sonnet-4.6-thinking": "claude-sonnet-4-6-thinking",
|
||||
"claude-sonnet-4.6": "claude-sonnet-4-6",
|
||||
"claude-opus-4.6-thinking": "claude-opus-4-6-thinking",
|
||||
"claude-opus-4.6": "claude-opus-4-6",
|
||||
// ── Claude 4.5 ──────────────────────────────────────────────────────────
|
||||
"claude-opus-4.5-thinking": "MODEL_CLAUDE_4_5_OPUS_THINKING",
|
||||
"claude-opus-4.5": "MODEL_CLAUDE_4_5_OPUS",
|
||||
"claude-sonnet-4.5-thinking": "MODEL_PRIVATE_3",
|
||||
"claude-sonnet-4.5": "MODEL_PRIVATE_2",
|
||||
"claude-haiku-4.5": "MODEL_PRIVATE_11",
|
||||
// ── GPT-5.5 ─────────────────────────────────────────────────────────────
|
||||
"gpt-5.5-xhigh-fast": "gpt-5-5-xhigh-priority",
|
||||
"gpt-5.5-high-fast": "gpt-5-5-high-priority",
|
||||
"gpt-5.5-medium-fast": "gpt-5-5-medium-priority",
|
||||
"gpt-5.5-low-fast": "gpt-5-5-low-priority",
|
||||
"gpt-5.5-none-fast": "gpt-5-5-none-priority",
|
||||
"gpt-5.5-xhigh": "gpt-5-5-xhigh",
|
||||
"gpt-5.5-high": "gpt-5-5-high",
|
||||
"gpt-5.5-medium": "gpt-5-5-medium",
|
||||
"gpt-5.5-low": "gpt-5-5-low",
|
||||
"gpt-5.5-none": "gpt-5-5-none",
|
||||
"gpt-5.5-review": "gpt-5-5-review",
|
||||
"gpt-5.5": "gpt-5-5-medium",
|
||||
// ── GPT-5.4 ─────────────────────────────────────────────────────────────
|
||||
"gpt-5.4-xhigh-fast": "gpt-5-4-xhigh-priority",
|
||||
"gpt-5.4-high-fast": "gpt-5-4-high-priority",
|
||||
"gpt-5.4-medium-fast": "gpt-5-4-medium-priority",
|
||||
"gpt-5.4-low-fast": "gpt-5-4-low-priority",
|
||||
"gpt-5.4-none-fast": "gpt-5-4-none-priority",
|
||||
"gpt-5.4-xhigh": "gpt-5-4-xhigh",
|
||||
"gpt-5.4-high": "gpt-5-4-high",
|
||||
"gpt-5.4-medium": "gpt-5-4-medium",
|
||||
"gpt-5.4-low": "gpt-5-4-low",
|
||||
"gpt-5.4-none": "gpt-5-4-none",
|
||||
"gpt-5.4-mini-xhigh": "gpt-5-4-mini-xhigh",
|
||||
"gpt-5.4-mini-high": "gpt-5-4-mini-high",
|
||||
"gpt-5.4-mini-medium": "gpt-5-4-mini-medium",
|
||||
"gpt-5.4-mini-low": "gpt-5-4-mini-low",
|
||||
"gpt-5.4": "gpt-5-4-medium",
|
||||
// ── GPT-5.3-Codex ───────────────────────────────────────────────────────
|
||||
"gpt-5.3-codex-xhigh-fast": "gpt-5-3-codex-xhigh-priority",
|
||||
"gpt-5.3-codex-high-fast": "gpt-5-3-codex-high-priority",
|
||||
"gpt-5.3-codex-medium-fast": "gpt-5-3-codex-medium-priority",
|
||||
"gpt-5.3-codex-low-fast": "gpt-5-3-codex-low-priority",
|
||||
"gpt-5.3-codex-xhigh": "gpt-5-3-codex-xhigh",
|
||||
"gpt-5.3-codex-high": "gpt-5-3-codex-high",
|
||||
"gpt-5.3-codex-medium": "gpt-5-3-codex-medium",
|
||||
"gpt-5.3-codex-low": "gpt-5-3-codex-low",
|
||||
"gpt-5.3-codex": "gpt-5-3-codex-medium",
|
||||
// ── GPT-5.2 ─────────────────────────────────────────────────────────────
|
||||
"gpt-5.2-xhigh": "MODEL_GPT_5_2_XHIGH",
|
||||
"gpt-5.2-high": "MODEL_GPT_5_2_HIGH",
|
||||
"gpt-5.2-medium": "MODEL_GPT_5_2_MEDIUM",
|
||||
"gpt-5.2-low": "MODEL_GPT_5_2_LOW",
|
||||
"gpt-5.2-none": "MODEL_GPT_5_2_NONE",
|
||||
"gpt-5.2": "MODEL_GPT_5_2_MEDIUM",
|
||||
// ── GPT-5 ───────────────────────────────────────────────────────────────
|
||||
"gpt-5": "gpt-5",
|
||||
// ── GPT-4.1 / 4o ────────────────────────────────────────────────────────
|
||||
"gpt-4.1": "MODEL_CHAT_GPT_4_1_2025_04_14",
|
||||
"gpt-4.1-mini": "gpt-4.1-mini",
|
||||
"gpt-4o": "MODEL_CHAT_GPT_4O_2024_08_06",
|
||||
// ── Gemini ──────────────────────────────────────────────────────────────
|
||||
"gemini-3.1-pro-high": "gemini-3-1-pro-high",
|
||||
"gemini-3.1-pro-low": "gemini-3-1-pro-low",
|
||||
"gemini-3.1-pro": "gemini-3-1-pro-high",
|
||||
"gemini-3.0-flash-high": "MODEL_GOOGLE_GEMINI_3_0_FLASH_HIGH",
|
||||
"gemini-3.0-flash-medium": "MODEL_GOOGLE_GEMINI_3_0_FLASH_MEDIUM",
|
||||
"gemini-3.0-flash-low": "MODEL_GOOGLE_GEMINI_3_0_FLASH_LOW",
|
||||
"gemini-3.0-flash-minimal": "MODEL_GOOGLE_GEMINI_3_0_FLASH_MINIMAL",
|
||||
"gemini-3.0-flash": "MODEL_GOOGLE_GEMINI_3_0_FLASH_HIGH",
|
||||
"gemini-2.5-pro": "MODEL_GOOGLE_GEMINI_2_5_PRO",
|
||||
// ── Others ──────────────────────────────────────────────────────────────
|
||||
"deepseek-v4": "deepseek-v4",
|
||||
"kimi-k2.6": "kimi-k2-6",
|
||||
"kimi-k2.5": "kimi-k2-5",
|
||||
"glm-5.1": "glm-5-1",
|
||||
};
|
||||
|
||||
export function resolveWsModelId(model) {
|
||||
return MODEL_ALIAS_MAP[model] ?? model;
|
||||
}
|
||||
|
||||
// ─── Minimal protobuf encoder ────────────────────────────────────────────────
|
||||
// Wire types: 0 = varint, 2 = length-delimited.
|
||||
|
||||
function encodeVarint(value) {
|
||||
const bytes = [];
|
||||
let v = value >>> 0;
|
||||
while (v > 0x7f) {
|
||||
bytes.push((v & 0x7f) | 0x80);
|
||||
v >>>= 7;
|
||||
}
|
||||
bytes.push(v & 0x7f);
|
||||
return new Uint8Array(bytes);
|
||||
}
|
||||
|
||||
function concatBytes(arrays) {
|
||||
const total = arrays.reduce((n, a) => n + a.length, 0);
|
||||
const out = new Uint8Array(total);
|
||||
let off = 0;
|
||||
for (const a of arrays) {
|
||||
out.set(a, off);
|
||||
off += a.length;
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
const TEXT_ENC = new TextEncoder();
|
||||
const TEXT_DEC = new TextDecoder();
|
||||
|
||||
function encodeField(fieldNum, payload) {
|
||||
const tag = encodeVarint((fieldNum << 3) | 2);
|
||||
const len = encodeVarint(payload.length);
|
||||
return concatBytes([tag, len, payload]);
|
||||
}
|
||||
|
||||
function encodeString(fieldNum, value) {
|
||||
return encodeField(fieldNum, TEXT_ENC.encode(value));
|
||||
}
|
||||
|
||||
function encodeMessage(fieldNum, msg) {
|
||||
return encodeField(fieldNum, msg);
|
||||
}
|
||||
|
||||
// ─── Protobuf message builders ───────────────────────────────────────────────
|
||||
|
||||
function buildMetadata(apiKey, sessionId) {
|
||||
return concatBytes([
|
||||
encodeString(1, apiKey),
|
||||
encodeString(2, WS_IDE_NAME),
|
||||
encodeString(3, WS_IDE_VERSION),
|
||||
encodeString(4, WS_EXT_VERSION),
|
||||
encodeString(5, sessionId),
|
||||
encodeString(6, WS_LOCALE),
|
||||
]);
|
||||
}
|
||||
|
||||
function buildModelOrAlias(model) {
|
||||
return encodeString(1, model);
|
||||
}
|
||||
|
||||
function buildChatMessage(msg) {
|
||||
const parts = [encodeString(1, msg.role), encodeString(2, msg.content)];
|
||||
if (msg.toolCallId) parts.push(encodeString(3, msg.toolCallId));
|
||||
return concatBytes(parts);
|
||||
}
|
||||
|
||||
export function buildGetChatMessageRequest(apiKey, model, messages) {
|
||||
const sessionId = randomUUID();
|
||||
const cascadeId = randomUUID();
|
||||
|
||||
const parts = [
|
||||
encodeMessage(1, buildMetadata(apiKey, sessionId)), // metadata
|
||||
encodeString(2, cascadeId), // cascade_id
|
||||
encodeMessage(3, buildModelOrAlias(model)), // model_or_alias
|
||||
];
|
||||
|
||||
for (const msg of messages) {
|
||||
parts.push(encodeMessage(4, buildChatMessage(msg))); // repeated messages
|
||||
}
|
||||
|
||||
return concatBytes(parts);
|
||||
}
|
||||
|
||||
// ─── gRPC-web framing ────────────────────────────────────────────────────────
|
||||
|
||||
export function grpcWebFrame(payload) {
|
||||
const frame = new Uint8Array(5 + payload.length);
|
||||
frame[0] = 0x00; // no compression
|
||||
const view = new DataView(frame.buffer);
|
||||
view.setUint32(1, payload.length, false); // big-endian length
|
||||
frame.set(payload, 5);
|
||||
return frame;
|
||||
}
|
||||
|
||||
// ─── Protobuf response decoder ───────────────────────────────────────────────
|
||||
// CompletionChunk (oneof):
|
||||
// field 1 → ContentChunk { field 1: string text }
|
||||
// field 2 → ToolCallChunk (skipped)
|
||||
// field 3 → DoneChunk { field 1: UsageStats{ field1: prompt, field2: completion } }
|
||||
// field 4 → ErrorChunk { field 1: string message }
|
||||
|
||||
function readVarint(buf, offset) {
|
||||
let result = 0;
|
||||
let shift = 0;
|
||||
while (offset < buf.length) {
|
||||
const b = buf[offset++];
|
||||
result |= (b & 0x7f) << shift;
|
||||
if ((b & 0x80) === 0) break;
|
||||
shift += 7;
|
||||
}
|
||||
return [result >>> 0, offset];
|
||||
}
|
||||
|
||||
function decodeStringField(buf, targetField) {
|
||||
let offset = 0;
|
||||
while (offset < buf.length) {
|
||||
let tag;
|
||||
[tag, offset] = readVarint(buf, offset);
|
||||
const fieldNum = tag >>> 3;
|
||||
const wireType = tag & 0x07;
|
||||
if (wireType === 2) {
|
||||
let len;
|
||||
[len, offset] = readVarint(buf, offset);
|
||||
const payload = buf.slice(offset, offset + len);
|
||||
offset += len;
|
||||
if (fieldNum === targetField) return TEXT_DEC.decode(payload);
|
||||
} else if (wireType === 0) {
|
||||
let v;
|
||||
[v, offset] = readVarint(buf, offset);
|
||||
} else if (wireType === 1) {
|
||||
offset += 8;
|
||||
} else if (wireType === 5) {
|
||||
offset += 4;
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function decodeDoneChunk(buf) {
|
||||
// DoneChunk: field 1 = UsageStats (nested)
|
||||
// UsageStats: field 1 = prompt_tokens (varint), field 2 = completion_tokens (varint)
|
||||
let offset = 0;
|
||||
let usageBytes = null;
|
||||
while (offset < buf.length) {
|
||||
let tag;
|
||||
[tag, offset] = readVarint(buf, offset);
|
||||
const fieldNum = tag >>> 3;
|
||||
const wireType = tag & 0x07;
|
||||
if (wireType === 2) {
|
||||
let len;
|
||||
[len, offset] = readVarint(buf, offset);
|
||||
if (fieldNum === 1) usageBytes = buf.slice(offset, offset + len);
|
||||
offset += len;
|
||||
} else if (wireType === 0) {
|
||||
let v;
|
||||
[v, offset] = readVarint(buf, offset);
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (!usageBytes) return [0, 0];
|
||||
let promptTokens = 0;
|
||||
let completionTokens = 0;
|
||||
offset = 0;
|
||||
while (offset < usageBytes.length) {
|
||||
let tag;
|
||||
[tag, offset] = readVarint(usageBytes, offset);
|
||||
const fieldNum = tag >>> 3;
|
||||
const wireType = tag & 0x07;
|
||||
if (wireType === 0) {
|
||||
let v;
|
||||
[v, offset] = readVarint(usageBytes, offset);
|
||||
if (fieldNum === 1) promptTokens = v;
|
||||
else if (fieldNum === 2) completionTokens = v;
|
||||
} else if (wireType === 2) {
|
||||
let len;
|
||||
[len, offset] = readVarint(usageBytes, offset);
|
||||
offset += len;
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
return [promptTokens, completionTokens];
|
||||
}
|
||||
|
||||
export function decodeCompletionChunk(buf) {
|
||||
let offset = 0;
|
||||
while (offset < buf.length) {
|
||||
let tag;
|
||||
[tag, offset] = readVarint(buf, offset);
|
||||
const fieldNum = tag >>> 3;
|
||||
const wireType = tag & 0x07;
|
||||
|
||||
if (wireType === 2) {
|
||||
let len;
|
||||
[len, offset] = readVarint(buf, offset);
|
||||
const payload = buf.slice(offset, offset + len);
|
||||
offset += len;
|
||||
|
||||
if (fieldNum === 1) {
|
||||
const text = decodeStringField(payload, 1);
|
||||
if (text !== null) return { kind: "content", text };
|
||||
} else if (fieldNum === 3) {
|
||||
const usage = decodeDoneChunk(payload);
|
||||
return { kind: "done", promptTokens: usage[0], completionTokens: usage[1] };
|
||||
} else if (fieldNum === 4) {
|
||||
const msg = decodeStringField(payload, 1);
|
||||
return { kind: "error", message: msg ?? "unknown windsurf error" };
|
||||
}
|
||||
// field 2 = ToolCallChunk — not yet handled; skip
|
||||
} else if (wireType === 0) {
|
||||
let v;
|
||||
[v, offset] = readVarint(buf, offset);
|
||||
} else if (wireType === 1) {
|
||||
offset += 8;
|
||||
} else if (wireType === 5) {
|
||||
offset += 4;
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
return { kind: "unknown" };
|
||||
}
|
||||
|
||||
// ─── OpenAI messages → Windsurf wire ─────────────────────────────────────────
|
||||
|
||||
function openAIMessagesToWs(messages) {
|
||||
const out = [];
|
||||
for (const m of messages) {
|
||||
const role = String(m.role || "user");
|
||||
let content = "";
|
||||
if (typeof m.content === "string") {
|
||||
content = m.content;
|
||||
} else if (Array.isArray(m.content)) {
|
||||
for (const part of m.content) {
|
||||
if (part && typeof part === "object" && part.type === "text") {
|
||||
content += String(part.text || "");
|
||||
}
|
||||
}
|
||||
}
|
||||
out.push({ role, content, toolCallId: m.tool_call_id });
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
// ─── WindsurfExecutor ────────────────────────────────────────────────────────
|
||||
|
||||
export class WindsurfExecutor extends BaseExecutor {
|
||||
constructor() {
|
||||
super("windsurf", PROVIDERS.windsurf || { id: "windsurf", baseUrl: WS_CHAT_URL });
|
||||
}
|
||||
|
||||
buildUrl() {
|
||||
return WS_CHAT_URL;
|
||||
}
|
||||
|
||||
buildHeaders(credentials, stream = true) {
|
||||
const token = credentials?.accessToken || credentials?.apiKey || "";
|
||||
return {
|
||||
"Content-Type": "application/grpc-web+proto",
|
||||
Accept: "application/grpc-web+proto",
|
||||
// Codeium apiKey also goes in Metadata.api_key (protobuf field) — see request body.
|
||||
...(token ? { Authorization: `Bearer ${token}` } : {}),
|
||||
"User-Agent": `windsurf/${WS_IDE_VERSION}`,
|
||||
"X-Grpc-Web": "1",
|
||||
};
|
||||
}
|
||||
|
||||
// Request body is built manually in execute() — requires model + messages.
|
||||
transformRequest() {
|
||||
return null;
|
||||
}
|
||||
|
||||
async execute({ model, body, stream, credentials, signal, log, upstreamExtraHeaders, proxyOptions = null }) {
|
||||
const apiKey = credentials?.accessToken || credentials?.apiKey || "";
|
||||
const wsModel = resolveWsModelId(model);
|
||||
|
||||
const b = body ?? {};
|
||||
const rawMessages = Array.isArray(b.messages) ? b.messages : [];
|
||||
let wsMessages = openAIMessagesToWs(rawMessages);
|
||||
if (wsMessages.length === 0) {
|
||||
wsMessages.push({ role: "user", content: "" });
|
||||
}
|
||||
|
||||
const protoPayload = buildGetChatMessageRequest(apiKey, wsModel, wsMessages);
|
||||
const framedPayload = grpcWebFrame(protoPayload);
|
||||
|
||||
const url = this.buildUrl();
|
||||
const headers = this.buildHeaders(credentials);
|
||||
if (upstreamExtraHeaders) Object.assign(headers, upstreamExtraHeaders);
|
||||
|
||||
log?.debug?.("WS", `Windsurf → ${wsModel} (${wsMessages.length} messages)`);
|
||||
|
||||
const upstream = await proxyAwareFetch(url, {
|
||||
method: "POST",
|
||||
headers,
|
||||
body: framedPayload,
|
||||
signal,
|
||||
}, proxyOptions);
|
||||
|
||||
if (!upstream.ok && upstream.status !== 200) {
|
||||
return { response: upstream, url, headers, transformedBody: protoPayload };
|
||||
}
|
||||
|
||||
const sseResponse = this.transformToSSE(upstream, model);
|
||||
return { response: sseResponse, url, headers, transformedBody: protoPayload };
|
||||
}
|
||||
|
||||
// Convert a gRPC-web binary response into an OpenAI-compatible SSE stream.
|
||||
transformToSSE(upstream, model) {
|
||||
const responseId = `chatcmpl-ws-${Date.now()}`;
|
||||
const created = Math.floor(Date.now() / 1000);
|
||||
const executor = this;
|
||||
|
||||
const sseStream = new ReadableStream({
|
||||
async start(controller) {
|
||||
const enc = new TextEncoder();
|
||||
let roleEmitted = false;
|
||||
let totalText = "";
|
||||
let promptTokens = 0;
|
||||
let completionTokens = 0;
|
||||
let hadError = null;
|
||||
|
||||
const emit = (data) => controller.enqueue(enc.encode(data));
|
||||
|
||||
try {
|
||||
let pending = new Uint8Array(0);
|
||||
const reader = upstream.body?.getReader();
|
||||
|
||||
const handleFrame = (flag, payload) => {
|
||||
if (flag === 0x80) {
|
||||
// Trailer frame — contains grpc-status, grpc-message
|
||||
const trailer = TEXT_DEC.decode(payload);
|
||||
const statusMatch = /grpc-status:\s*(\d+)/i.exec(trailer);
|
||||
if (statusMatch && statusMatch[1] !== "0") {
|
||||
const msgMatch = /grpc-message:\s*(.+)/i.exec(trailer);
|
||||
hadError = msgMatch
|
||||
? decodeURIComponent(msgMatch[1].trim())
|
||||
: `gRPC status ${statusMatch[1]}`;
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (flag !== 0x00) return; // skip unknown flags
|
||||
|
||||
const chunk = executor.constructor.decodeCompletionChunk
|
||||
? executor.constructor.decodeCompletionChunk(payload)
|
||||
: decodeCompletionChunk(payload);
|
||||
|
||||
if (chunk.kind === "content" && chunk.text) {
|
||||
totalText += chunk.text;
|
||||
if (!roleEmitted) {
|
||||
emit(`data: ${JSON.stringify({
|
||||
id: responseId, object: "chat.completion.chunk", created, model,
|
||||
choices: [{ index: 0, delta: { role: "assistant", content: "" }, finish_reason: null }],
|
||||
})}\n\n`);
|
||||
roleEmitted = true;
|
||||
}
|
||||
emit(`data: ${JSON.stringify({
|
||||
id: responseId, object: "chat.completion.chunk", created, model,
|
||||
choices: [{ index: 0, delta: { content: chunk.text }, finish_reason: null }],
|
||||
})}\n\n`);
|
||||
} else if (chunk.kind === "done") {
|
||||
promptTokens = chunk.promptTokens;
|
||||
completionTokens = chunk.completionTokens;
|
||||
} else if (chunk.kind === "error") {
|
||||
hadError = chunk.message;
|
||||
}
|
||||
};
|
||||
|
||||
const drainFrames = () => {
|
||||
let offset = 0;
|
||||
while (offset + 5 <= pending.length) {
|
||||
const flag = pending[offset];
|
||||
const len =
|
||||
(pending[offset + 1] << 24) |
|
||||
(pending[offset + 2] << 16) |
|
||||
(pending[offset + 3] << 8) |
|
||||
pending[offset + 4];
|
||||
if (len < 0 || offset + 5 + len > pending.length) break;
|
||||
handleFrame(flag, pending.slice(offset + 5, offset + 5 + len));
|
||||
offset += 5 + len;
|
||||
}
|
||||
if (offset > 0) pending = pending.slice(offset);
|
||||
};
|
||||
|
||||
if (reader) {
|
||||
try {
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
if (!value) continue;
|
||||
pending = pending.length === 0 ? value : concatBytes([pending, value]);
|
||||
drainFrames();
|
||||
}
|
||||
} finally {
|
||||
reader.releaseLock();
|
||||
}
|
||||
}
|
||||
drainFrames();
|
||||
|
||||
if (hadError) {
|
||||
emit(`data: ${JSON.stringify({
|
||||
error: { message: hadError, type: "windsurf_error", code: "upstream_error" },
|
||||
})}\n\n`);
|
||||
emit("data: [DONE]\n\n");
|
||||
controller.close();
|
||||
return;
|
||||
}
|
||||
|
||||
// Unary fallback: nothing streamed but text decoded → emit as one chunk.
|
||||
if (!roleEmitted && totalText) {
|
||||
emit(`data: ${JSON.stringify({
|
||||
id: responseId, object: "chat.completion.chunk", created, model,
|
||||
choices: [{ index: 0, delta: { role: "assistant", content: "" }, finish_reason: null }],
|
||||
})}\n\n`);
|
||||
emit(`data: ${JSON.stringify({
|
||||
id: responseId, object: "chat.completion.chunk", created, model,
|
||||
choices: [{ index: 0, delta: { content: totalText }, finish_reason: null }],
|
||||
})}\n\n`);
|
||||
}
|
||||
|
||||
const finishPayload = {
|
||||
id: responseId, object: "chat.completion.chunk", created, model,
|
||||
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
|
||||
};
|
||||
if (promptTokens > 0 || completionTokens > 0) {
|
||||
finishPayload.usage = {
|
||||
prompt_tokens: promptTokens,
|
||||
completion_tokens: completionTokens,
|
||||
total_tokens: promptTokens + completionTokens,
|
||||
};
|
||||
}
|
||||
emit(`data: ${JSON.stringify(finishPayload)}\n\n`);
|
||||
emit("data: [DONE]\n\n");
|
||||
} catch (err) {
|
||||
const msg = err?.message ? String(err.message) : String(err);
|
||||
emit(`data: ${JSON.stringify({
|
||||
error: { message: `Windsurf stream error: ${msg}`, type: "windsurf_error" },
|
||||
})}\n\n`);
|
||||
emit("data: [DONE]\n\n");
|
||||
}
|
||||
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
return new Response(sseStream, {
|
||||
status: 200,
|
||||
headers: {
|
||||
"Content-Type": "text/event-stream",
|
||||
"Cache-Control": "no-cache",
|
||||
Connection: "keep-alive",
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
// apiKey is long-lived (Firebase-derived or Devin ide_token); refresh handled out-of-band.
|
||||
async refreshCredentials() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
export default WindsurfExecutor;
|
||||
304
open-sse/executors/zed.js
Normal file
304
open-sse/executors/zed.js
Normal file
@@ -0,0 +1,304 @@
|
||||
// ZedHostedExecutor — routes requests to Zed's hosted LLM aggregator
|
||||
// (cloud.zed.dev/completions), a multi-format proxy fronting
|
||||
// Anthropic/OpenAI/Google/xAI depending on the requested model.
|
||||
//
|
||||
// Wire protocol: POST /completions with an NDJSON/SSE-ish body-per-line
|
||||
// response stream (`{"event": <provider-shaped-chunk>}` / `{"status": ...}` /
|
||||
// `[DONE]`), authenticated with a short-lived LLM bearer token exchanged from
|
||||
// the RSA-decrypted access_token (see open-sse/shared/zedAuth.js). The
|
||||
// provider-shaped chunk is Claude/Gemini/OpenAI-Responses/xAI(OpenAI-shaped)
|
||||
// depending on which upstream Zed fronts for the model — translated back to
|
||||
// OpenAI Chat Completions by reusing the existing translators.
|
||||
//
|
||||
// Overrides execute() entirely (does NOT use DefaultExecutor's pipeline) because the Zed wire
|
||||
// shape (thread envelope, LLM-token exchange, NDJSON status frames) doesn't
|
||||
// fit the generic transformRequest/buildUrl contract.
|
||||
|
||||
import { BaseExecutor } from "./base.js";
|
||||
import { FORMATS } from "../translator/formats.js";
|
||||
import { initState } from "../translator/index.js";
|
||||
import { openaiToClaudeRequest } from "../translator/request/openai-to-claude.js";
|
||||
import { openaiToGeminiRequest } from "../translator/request/openai-to-gemini.js";
|
||||
import { openaiToOpenAIResponsesRequest } from "../translator/request/openai-responses.js";
|
||||
import { claudeToOpenAIResponse } from "../translator/response/claude-to-openai.js";
|
||||
import { geminiToOpenAIResponse } from "../translator/response/gemini-to-openai.js";
|
||||
import { openaiResponsesToOpenAIResponse } from "../translator/response/openai-responses.js";
|
||||
import {
|
||||
ZED_HEADERS,
|
||||
resolveZedModels,
|
||||
zedLlmFetch,
|
||||
} from "../shared/zedAuth.js";
|
||||
|
||||
const ZED_PROVIDER = {
|
||||
anthropic: "Anthropic",
|
||||
openai: "OpenAi",
|
||||
google: "Google",
|
||||
xai: "XAi",
|
||||
};
|
||||
|
||||
function normalizeZedProvider(value, model) {
|
||||
const raw = String(value || "").toLowerCase();
|
||||
if (raw === "anthropic") return ZED_PROVIDER.anthropic;
|
||||
if (raw === "openai" || raw === "open_ai") return ZED_PROVIDER.openai;
|
||||
if (raw === "google" || raw === "gemini") return ZED_PROVIDER.google;
|
||||
if (raw === "xai" || raw === "x_ai" || raw === "x-ai") return ZED_PROVIDER.xai;
|
||||
|
||||
const m = String(model || "").toLowerCase();
|
||||
if (m.includes("claude")) return ZED_PROVIDER.anthropic;
|
||||
if (m.includes("gemini")) return ZED_PROVIDER.google;
|
||||
if (m.includes("grok") || m.includes("xai")) return ZED_PROVIDER.xai;
|
||||
return ZED_PROVIDER.openai;
|
||||
}
|
||||
|
||||
function buildProviderRequest(provider, model, body, stream, credentials) {
|
||||
if (provider === ZED_PROVIDER.anthropic) {
|
||||
return openaiToClaudeRequest(model, body, true);
|
||||
}
|
||||
if (provider === ZED_PROVIDER.google) {
|
||||
return openaiToGeminiRequest(model, body, true);
|
||||
}
|
||||
if (provider === ZED_PROVIDER.openai) {
|
||||
return openaiToOpenAIResponsesRequest(model, body, true, credentials);
|
||||
}
|
||||
// xAI is OpenAI-shaped — forward as-is.
|
||||
return { ...(body || {}), model, stream: stream !== false };
|
||||
}
|
||||
|
||||
function initProviderState(provider, model) {
|
||||
if (provider === ZED_PROVIDER.anthropic) return initState(FORMATS.CLAUDE);
|
||||
if (provider === ZED_PROVIDER.google) return initState(FORMATS.GEMINI);
|
||||
if (provider === ZED_PROVIDER.openai) return initState(FORMATS.OPENAI_RESPONSES);
|
||||
const state = initState(FORMATS.OPENAI);
|
||||
state.model = model;
|
||||
return state;
|
||||
}
|
||||
|
||||
function convertProviderEvent(provider, event, state) {
|
||||
if (provider === ZED_PROVIDER.anthropic) return claudeToOpenAIResponse(event, state);
|
||||
if (provider === ZED_PROVIDER.google) return geminiToOpenAIResponse(event, state);
|
||||
if (provider === ZED_PROVIDER.openai) return openaiResponsesToOpenAIResponse(event, state);
|
||||
return event;
|
||||
}
|
||||
|
||||
function createErrorChunk(model, message) {
|
||||
return {
|
||||
id: `chatcmpl-zed-error-${Date.now()}`,
|
||||
object: "chat.completion.chunk",
|
||||
created: Math.floor(Date.now() / 1000),
|
||||
model,
|
||||
choices: [
|
||||
{ index: 0, delta: { content: `[Zed error] ${message}` }, finish_reason: "stop" },
|
||||
],
|
||||
};
|
||||
}
|
||||
|
||||
function enqueueSseObject(controller, encoder, chunk) {
|
||||
if (!chunk) return;
|
||||
const items = Array.isArray(chunk) ? chunk : [chunk];
|
||||
for (const item of items) {
|
||||
if (!item) continue;
|
||||
controller.enqueue(encoder.encode(`data: ${JSON.stringify(item)}\n\n`));
|
||||
}
|
||||
}
|
||||
|
||||
function unwrapZedLine(line) {
|
||||
let text = line.replace(/\r$/, "").trim();
|
||||
if (!text) return null;
|
||||
if (text.startsWith("data:")) text = text.slice(5).trimStart();
|
||||
if (text === "[DONE]") return { done: true };
|
||||
try {
|
||||
const parsed = JSON.parse(text);
|
||||
if (parsed && Object.prototype.hasOwnProperty.call(parsed, "event")) {
|
||||
return { event: parsed.event };
|
||||
}
|
||||
if (parsed && Object.prototype.hasOwnProperty.call(parsed, "status")) {
|
||||
return { status: parsed.status };
|
||||
}
|
||||
return { event: parsed };
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function normalizeStatus(status) {
|
||||
if (!status) return null;
|
||||
if (typeof status === "string") return { type: status };
|
||||
if (typeof status === "object") {
|
||||
const key = Object.keys(status)[0];
|
||||
if (key && typeof status[key] === "object") return { type: key, ...status[key] };
|
||||
return status;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function wrapZedCompletionStream(response, provider, model) {
|
||||
if (!response.ok || !response.body) return response;
|
||||
|
||||
const decoder = new TextDecoder();
|
||||
const encoder = new TextEncoder();
|
||||
const state = initProviderState(provider, model);
|
||||
let buffer = "";
|
||||
let done = false;
|
||||
|
||||
const finish = (controller) => {
|
||||
if (done) return;
|
||||
const finalChunk = convertProviderEvent(provider, null, state);
|
||||
enqueueSseObject(controller, encoder, finalChunk);
|
||||
controller.enqueue(encoder.encode("data: [DONE]\n\n"));
|
||||
done = true;
|
||||
};
|
||||
|
||||
const processLine = (line, controller) => {
|
||||
if (done) return;
|
||||
const payload = unwrapZedLine(line);
|
||||
if (!payload) return;
|
||||
if (payload.done) {
|
||||
finish(controller);
|
||||
return;
|
||||
}
|
||||
if (payload.status) {
|
||||
const status = normalizeStatus(payload.status);
|
||||
if (status?.type === "failed" || status?.failed) {
|
||||
const failed = status.failed || status;
|
||||
const message = String(failed.message || failed.error || failed.code || "request failed");
|
||||
enqueueSseObject(controller, encoder, createErrorChunk(model, message));
|
||||
finish(controller);
|
||||
} else if (status?.type === "stream_ended" || status === "stream_ended") {
|
||||
finish(controller);
|
||||
}
|
||||
return;
|
||||
}
|
||||
const converted = convertProviderEvent(provider, payload.event, state);
|
||||
enqueueSseObject(controller, encoder, converted);
|
||||
};
|
||||
|
||||
const transformed = response.body.pipeThrough(
|
||||
new TransformStream({
|
||||
transform(chunk, controller) {
|
||||
buffer += decoder.decode(chunk, { stream: true });
|
||||
let nl;
|
||||
while ((nl = buffer.indexOf("\n")) !== -1) {
|
||||
const line = buffer.slice(0, nl);
|
||||
buffer = buffer.slice(nl + 1);
|
||||
processLine(line, controller);
|
||||
}
|
||||
},
|
||||
flush(controller) {
|
||||
buffer += decoder.decode();
|
||||
if (buffer) {
|
||||
processLine(buffer, controller);
|
||||
buffer = "";
|
||||
}
|
||||
finish(controller);
|
||||
},
|
||||
}),
|
||||
);
|
||||
|
||||
return new Response(transformed, {
|
||||
status: response.status,
|
||||
statusText: response.statusText,
|
||||
headers: {
|
||||
"Content-Type": "text/event-stream",
|
||||
"Cache-Control": "no-cache",
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
class ZedExecutor extends BaseExecutor {
|
||||
constructor() {
|
||||
super("zed");
|
||||
}
|
||||
|
||||
async resolveModel(model, credentials, signal, log) {
|
||||
try {
|
||||
const catalog = await resolveZedModels(credentials, { config: this.config, signal });
|
||||
let raw = catalog?.rawById?.get(model) ?? null;
|
||||
if (!raw) {
|
||||
const refreshed = await resolveZedModels(credentials, {
|
||||
config: this.config,
|
||||
signal,
|
||||
forceRefresh: true,
|
||||
});
|
||||
raw = refreshed?.rawById?.get(model) ?? null;
|
||||
}
|
||||
return { raw, provider: normalizeZedProvider(raw?.provider, model) };
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
log?.warn?.("ZED", `model catalog unavailable, inferring provider for ${model}: ${message}`);
|
||||
return { raw: null, provider: normalizeZedProvider(null, model) };
|
||||
}
|
||||
}
|
||||
|
||||
async execute({ model, body, stream, credentials, signal, log, proxyOptions = null }) {
|
||||
const { provider } = await this.resolveModel(model, credentials, signal, log);
|
||||
const providerRequest = buildProviderRequest(provider, model, body, stream, credentials);
|
||||
const bodyRecord = body || {};
|
||||
const payload = {
|
||||
thread_id: bodyRecord.thread_id || credentials?._clientSessionId,
|
||||
prompt_id: bodyRecord.prompt_id,
|
||||
provider,
|
||||
model,
|
||||
provider_request: providerRequest,
|
||||
};
|
||||
|
||||
const response = await zedLlmFetch(credentials, "/completions", {
|
||||
config: this.config,
|
||||
signal,
|
||||
fetchOptions: {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
Accept: "application/x-ndjson, text/event-stream, */*",
|
||||
"User-Agent": "9router/zed",
|
||||
"x-zed-version": this.config?.appVersion?.toString() || "0.200.0",
|
||||
[ZED_HEADERS.clientSupportsStatus]: "true",
|
||||
[ZED_HEADERS.clientSupportsStreamEnded]: "true",
|
||||
},
|
||||
body: JSON.stringify(payload),
|
||||
},
|
||||
});
|
||||
|
||||
const wrapped = response.ok ? wrapZedCompletionStream(response, provider, model) : response;
|
||||
return {
|
||||
response: wrapped,
|
||||
url: `${this.config?.llmBaseUrl || "https://cloud.zed.dev"}/completions`,
|
||||
headers: { "Content-Type": "application/json", Authorization: "Bearer <zed-llm-token>" },
|
||||
transformedBody: payload,
|
||||
};
|
||||
}
|
||||
|
||||
parseError(response, bodyText) {
|
||||
let parsed = null;
|
||||
try {
|
||||
parsed = JSON.parse(bodyText || "{}");
|
||||
} catch {
|
||||
parsed = null;
|
||||
}
|
||||
|
||||
const errorObj = parsed?.error || undefined;
|
||||
const code = parsed?.code || errorObj?.code || "";
|
||||
const rawMessage =
|
||||
parsed?.message || errorObj?.message || bodyText || response.statusText;
|
||||
if (code === "trial_blocked") {
|
||||
return {
|
||||
status: response.status,
|
||||
message: `Zed trial access is blocked upstream. The account can list hosted models, but Zed is refusing completions until trial/billing access is enabled or unblocked. Zed says: ${rawMessage}`,
|
||||
};
|
||||
}
|
||||
if (code) {
|
||||
return { status: response.status, message: `Zed ${code}: ${rawMessage}` };
|
||||
}
|
||||
return { status: response.status, message: rawMessage || `Zed upstream error: ${response.status}` };
|
||||
}
|
||||
|
||||
async refreshCredentials() {
|
||||
// Zed uses a long-lived RSA-decrypted access_token — no OAuth refresh.
|
||||
return null;
|
||||
}
|
||||
|
||||
needsRefresh() {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
export default ZedExecutor;
|
||||
Reference in New Issue
Block a user