From c933eefc2711ab94d3bb0508973f827532c7bd07 Mon Sep 17 00:00:00 2001 From: Amir Seify Date: Wed, 16 Sep 2026 21:31:22 +0700 Subject: [PATCH] fix(cursor): stop AgentService empty turns and silent tool hangs Cursor-hosted models (cu/composer-2.5, cu/cursor-grok-*, cu/default) returned HTTP 200 with an empty turn, or hung, whenever a client sent tools. - Fold system prompts into the current user message. custom_system_prompt (RunRequest field 8) makes AgentService return an empty turn. - Send ModelDetails (field 3); thinking variants (Composer, Grok, *-thinking) return an empty turn when only requested_model (field 9) is set. - Route tool-call history and declared tool schemas through AgentService: encode OpenAI tools into mcp_tools (field 4), decode McpArgs and emit real tool_calls with finish_reason tool_calls. - Map Composer thinking / Grok thinking_delta (field 4) into visible content instead of dropping the answer with the unsigned reasoning. - Ack request_context without echoing MCP tools (double-advertise stalls the HTTP/2 stream) and ack kv_server_message so the run proceeds. - Reject IDE builtin execs instead of failing the turn, so the model can continue with MCP tools or a text answer. - Add google.protobuf.Value / MCP encoders and a FIXED64 branch to encodeField in cursorProtobuf.js. RTK now compresses the source-format body before translation for cursor only: its translator rewrites role:tool into user XML, so the post-translate pass missed those tool results. Every other provider keeps the post-translate pass unchanged. --- CHANGELOG.md | 2 + open-sse/AGENTS.md | 2 +- open-sse/executors/cursor.js | 269 +++++++++++++++---- open-sse/handlers/chatCore.js | 24 +- open-sse/rtk/index.js | 2 +- open-sse/utils/cursorProtobuf.js | 221 ++++++++++++++- tests/unit/cursor-agent-exec-request.test.js | 81 +++++- tests/unit/cursor-agent-proto.test.js | 9 + tests/unit/rtk-cursor-pretranslate.test.js | 131 +++++++++ 9 files changed, 670 insertions(+), 71 deletions(-) create mode 100644 tests/unit/rtk-cursor-pretranslate.test.js diff --git a/CHANGELOG.md b/CHANGELOG.md index f16c8f53..2bcc3390 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,8 @@ - **i18n**: integrate Persian (fa) translation ## Fixes +- **Cursor**: stop AgentService empty turns (`OUT 0`) and silent hangs — fold system prompts instead of `custom_system_prompt`, send `ModelDetails`, read Composer/Grok `thinking_delta`, ack request-context without echoing MCP tools, and reject IDE execs so the model can continue +- **RTK**: for Cursor, compress source-format `tool_result` / `role:tool` **before** translation — its translator rewrites those shapes, so post-translate compression missed them. Other providers keep the post-translate pass unchanged - **OpenCode / OpenCode Go**: resolve 403 `FreeTierError` and 429 rate limits with canonical session format, valid User-Agent, and stable upstream session reuse; force stream and declare `forceStream` for free-tier SSE aggregation; cloak decoy tools, normalize Muse Free tool choice, and strip prior reasoning items on Responses models; route Union Alpha via Messages API - **Kiro**: preserve underscores in tool names (`mcp__server__tool`) and restore client tool names in responses; use neutral placeholder for tool-result-only turns; forward tool-result images - **Stream**: report aborts after HTTP 200 in-band (per-format error frames) instead of closing silently diff --git a/open-sse/AGENTS.md b/open-sse/AGENTS.md index 5ea0d971..7e63c051 100644 --- a/open-sse/AGENTS.md +++ b/open-sse/AGENTS.md @@ -4,7 +4,7 @@ Provider-agnostic SSE engine: one OpenAI-style request → any provider (LLM cha ## Request lifecycle (chat) -`handlers/chatCore.js` → `services/model.js` `parseModel` (resolve `provider/model`) → **pre-translate hooks** (`rtk/` tool_result compress, `rtk/headroom.js` proxy compress, `rtk/caveman.js` system inject — all fail-open) → `executors/index.js` `getExecutor(provider)` → `translator/index.js` `translateRequest` (client format → provider format) → `executor.execute()` (streams upstream) → `translateResponse` (provider chunks → client format) → SSE out. +`handlers/chatCore.js` → `services/model.js` `parseModel` (resolve `provider/model`) → **RTK for `cursor`** (`rtk/` compresses the source-format `tool_result` / `role:tool` in-place — its translator rewrites those shapes, so this one provider must run **before** translate) → `translator/index.js` `translateRequest` (client format → provider format) → **post-translate savers** (`rtk/` compress for every other provider, `rtk/headroom.js` proxy compress, `rtk/caveman.js` / `rtk/ponytail.js` system inject — all fail-open) → `executors/index.js` `getExecutor(provider)` → `executor.execute()` (streams upstream) → `translateResponse` (provider chunks → client format) → SSE out. ## Directory map diff --git a/open-sse/executors/cursor.js b/open-sse/executors/cursor.js index 0aefc623..8f20040d 100644 --- a/open-sse/executors/cursor.js +++ b/open-sse/executors/cursor.js @@ -7,13 +7,16 @@ import { wrapConnectRPCFrame, decodeMessage, parseConnectRPCFrame, - extractTextFromResponse + extractTextFromResponse, + encodeMcpTools, + decodeMcpArgs, } from "../utils/cursorProtobuf.js"; import { buildCursorHeaders } from "../utils/cursorChecksum.js"; import { estimateUsage } from "../utils/usageTracking.js"; import { SSE_DONE, SSE_HEADERS } from "../utils/sseConstants.js"; import { chatChunkSse, sseChunk } from "../utils/sse.js"; import { FORMATS } from "../translator/formats.js"; +import { ROLE, OPENAI_BLOCK } from "../translator/schema/index.js"; import { proxyAwareFetch } from "../utils/proxyFetch.js"; import zlib from "zlib"; import crypto from "crypto"; @@ -65,55 +68,74 @@ function textFromContent(content) { if (typeof content === "string") return content; if (!Array.isArray(content)) return ""; return content - .filter((part) => part?.type === "text" && typeof part.text === "string") + .filter((part) => part?.type === OPENAI_BLOCK.TEXT && typeof part.text === "string") .map((part) => part.text) .join("\n"); } -function isAgentTextRequest(body) { - // Many compatible clients always attach their built-in tool schemas, even - // for a normal text turn. Cursor's retired ChatService rejects those - // requests; AgentService can still answer the text turn, so ignore schemas - // here. A real tool-call/result conversation is kept on the legacy path - // until its AgentService tool protocol is implemented. - return Array.isArray(body?.messages) && body.messages.every((message) => { - if (message?.tool_calls?.length || message?.role === "tool") return false; - return typeof message?.content === "string" - || Array.isArray(message?.content) && message.content.every((part) => part?.type === "text"); +function isTextPart(part) { + return !part || part.type === OPENAI_BLOCK.TEXT || typeof part === "string"; +} + +export function isAgentCapableRequest(body) { + // ChatService rejects auto/composer and most thinking variants. AgentService + // can answer text turns (including declared tool schemas) and tool-call + // history. Image parts still need the legacy protobuf path. + if (!Array.isArray(body?.messages) || body.messages.length === 0) return false; + return body.messages.every((message) => { + if (Array.isArray(message?.content)) return message.content.every(isTextPart); + return message?.content == null || typeof message.content === "string"; }); } function encodeHistoryMessage(message) { const content = textFromContent(message?.content); - if (!content) return null; + const extras = []; + if (message?.role === ROLE.ASSISTANT && message.tool_calls?.length) { + for (const tc of message.tool_calls) { + extras.push(`[tool_call id=${tc.id || ""} name=${tc.function?.name || "tool"} args=${tc.function?.arguments || "{}"}]`); + } + } + if (message?.role === ROLE.TOOL) { + extras.push(`[tool_result id=${message.tool_call_id || ""}]`); + } + const textBody = [content, ...extras].filter(Boolean).join("\n"); + if (!textBody) return null; // ConversationHistoryMessage.user / .assistant -> repeated content -> text. - const text = agentString(1, content); - if (message.role === "assistant") { + const text = agentString(1, textBody); + if (message.role === ROLE.ASSISTANT) { return agentMessage(2, agentMessage(1, agentMessage(1, text))); } return agentMessage(1, agentMessage(1, agentMessage(1, text))); } -function buildAgentRunFrame(messages, model) { +export function buildAgentRunFrame(messages, model, tools = []) { + // custom_system_prompt (RunRequest field 8) makes AgentService return an + // empty turn. Fold system text into the current user message instead. const system = messages - .filter((message) => message?.role === "system") + .filter((message) => message?.role === ROLE.SYSTEM) .map((message) => textFromContent(message.content)) .filter(Boolean) .join("\n\n"); - const chatMessages = messages.filter((message) => message?.role !== "system"); - const currentIndex = [...chatMessages].map((message) => message?.role).lastIndexOf("user"); + const chatMessages = messages.filter((message) => message?.role !== ROLE.SYSTEM); + const currentIndex = [...chatMessages].map((message) => message?.role).lastIndexOf(ROLE.USER); const current = currentIndex >= 0 ? chatMessages[currentIndex] : chatMessages.at(-1); const history = chatMessages .slice(0, currentIndex >= 0 ? currentIndex : -1) .map(encodeHistoryMessage) .filter(Boolean); - const userText = textFromContent(current?.content) || "Continue."; + const rawUser = textFromContent(current?.content) || "Continue."; + const userText = system ? `${system}\n\n${rawUser}` : rawUser; // agent.v1.UserMessageAction.user_message and its optional history. + // selected_context (3) + mode=1 (4) match cursor-agent's wire format; without + // them the server may accept the RPC and stream an empty turn. const userMessage = concatBuffers( agentString(1, userText), agentString(2, crypto.randomUUID()), + agentMessage(3, new Uint8Array()), + encodeField(4, PROTOBUF_VARINT, 1), ); const conversationHistory = history.length ? concatBuffers(...history.map((entry) => agentMessage(1, entry))) @@ -124,11 +146,20 @@ function buildAgentRunFrame(messages, model) { ); const conversationAction = agentMessage(1, userAction); const requestedModel = concatBuffers(agentString(1, model), agentBool(7, true)); + // ModelDetails (field 3): thinking variants (Composer, Grok, *-thinking) + // return an empty turn when only RequestedModel (field 9) is set. + const modelDetails = concatBuffers( + agentString(1, model), + agentString(3, model), + agentString(4, model), + ); + const mcpTools = encodeMcpTools(tools); const runRequest = concatBuffers( // An empty ConversationStateStructure starts a fresh local agent session. agentMessage(1, new Uint8Array()), agentMessage(2, conversationAction), - ...(system ? [agentString(8, system)] : []), + agentMessage(3, modelDetails), + ...(mcpTools.length ? [agentMessage(4, mcpTools)] : []), agentMessage(9, requestedModel), ); @@ -157,13 +188,51 @@ function decodeAgentFrames(buffer, onFrame) { return pending; } -function createRequestContextResponse() { - // AgentService asks every run for client context. 9router has no IDE file - // context, so acknowledge with an empty RequestContext. +function execIds(execRequest) { + const id = Number(execRequest?.get(1)?.[0]?.value || 0); + const execId = extractAgentString(execRequest, 15); + return { id, execId }; +} + +function wrapExecClientMessage(execMsgId, execId, resultField, resultPayload) { + const parts = []; + if (execMsgId) parts.push(encodeField(1, PROTOBUF_VARINT, execMsgId)); + parts.push(agentString(15, execId || "")); + parts.push(encodeField(resultField, PROTOBUF_LEN, resultPayload || new Uint8Array())); + return wrapConnectRPCFrame(agentMessage(2, concatBuffers(...parts))); +} + +function createRequestContextResponse(execRequest) { + // Tools already go out on AgentRunRequest.mcp_tools. Echoing them again on + // this ack makes AgentService stall silently (0 SSE bytes until abort). + const { id, execId } = execIds(execRequest); const requestContextSuccess = agentMessage(1, new Uint8Array()); const requestContextResult = agentMessage(1, requestContextSuccess); - const execClientMessage = agentMessage(10, requestContextResult); - return wrapConnectRPCFrame(agentMessage(2, execClientMessage)); + return wrapExecClientMessage(id, execId, 10, requestContextResult); +} + +// ExecServerMessage variant → ExecClientMessage result field (same numbers). +const EXEC_RESULT_FIELD = { + 2: 2, 3: 3, 4: 4, 5: 5, 7: 7, 8: 8, 9: 9, 16: 16, 20: 20, 23: 23, +}; + +function rejectExecRequest(execRequest) { + const { id, execId } = execIds(execRequest); + const variant = [...(execRequest?.keys?.() || [])].find((field) => field !== 1 && field !== 15); + const resultField = EXEC_RESULT_FIELD[variant]; + if (!resultField) return null; + // Diagnostics has no rejected variant — empty success unblocks the stream. + if (variant === 9) return wrapExecClientMessage(id, execId, 9, new Uint8Array()); + const rejected = agentMessage(2, agentString(2, "Tool not available in this environment. Use the MCP tools provided instead.")); + return wrapExecClientMessage(id, execId, resultField, rejected); +} + +function encodeKvClientMessage(kvId, resultField, resultPayload, metadata) { + const parts = []; + if (kvId) parts.push(encodeField(1, PROTOBUF_VARINT, kvId)); + parts.push(encodeField(resultField, PROTOBUF_LEN, resultPayload || new Uint8Array())); + if (metadata && metadata.length) parts.push(encodeField(4, PROTOBUF_LEN, metadata)); + return wrapConnectRPCFrame(agentMessage(3, concatBuffers(...parts))); } const CURSOR_STREAM_DEBUG = process.env.CURSOR_STREAM_DEBUG === "1"; @@ -479,7 +548,7 @@ export class CursorExecutor extends BaseExecutor { }; } - async executeAgent({ model, body, stream, credentials, signal }) { + async executeAgent({ model, body, stream, credentials, signal, log }) { const agentEndpoint = PROVIDER_OAUTH.cursor?.agentEndpoint; if (!agentEndpoint) throw new Error("Cursor AgentService endpoint is not configured"); @@ -491,9 +560,10 @@ export class CursorExecutor extends BaseExecutor { } let session; + const tools = body.tools || []; try { session = this.openAgentHttp2Stream(url, headers, requestController.signal); - session.write(buildAgentRunFrame(body.messages || [], model)); + session.write(buildAgentRunFrame(body.messages || [], model, tools)); } catch (error) { throw new Error(`Cursor AgentService request failed: ${error.message}`); } @@ -533,8 +603,23 @@ export class CursorExecutor extends BaseExecutor { // so strict clients such as Claude Code accept the completed stream. const responseId = `chatcmpl-msg_${Date.now()}`; const created = Math.floor(Date.now() / 1000); + const composerModel = isComposerModel(model); let pending = Buffer.alloc(0); let finished = false; + let thinkingAcc = ""; + let emittedVisible = 0; + let emittedText = false; + + const flushThinkingFallback = (onEvent) => { + if (emittedText || !thinkingAcc) return; + const fallback = composerModel + ? visibleComposerContentFromThinking(thinkingAcc) + : thinkingAcc.trim(); + if (fallback) { + emittedText = true; + onEvent({ type: "text", value: fallback }); + } + }; const consume = async (onEvent) => { try { @@ -553,32 +638,87 @@ export class CursorExecutor extends BaseExecutor { const update = decodeMessage(serverMessage.get(1)[0].value); if (update.has(1)) { const textDelta = extractAgentString(decodeMessage(update.get(1)[0].value), 1); - if (textDelta) onEvent({ type: "text", value: textDelta }); + if (textDelta) { + emittedText = true; + onEvent({ type: "text", value: textDelta }); + } } - // Cursor's AgentService emits internal reasoning without the - // cryptographic signature required by Anthropic thinking blocks. - // Forwarding it makes strict Anthropic clients (Claude Code) - // discard or wait on an otherwise complete response. Keep the - // reasoning upstream-only and emit the normal answer text. + // thinking_delta (field 4). Composer (and some Grok variants) put + // the visible answer after here and never send text_delta. + if (update.has(4)) { + const thinkingDelta = extractAgentString(decodeMessage(update.get(4)[0].value), 1); + if (thinkingDelta) { + thinkingAcc += thinkingDelta; + if (composerModel) { + const visible = visibleComposerContentFromThinking(thinkingAcc); + if (visible.length > emittedVisible) { + const deltaContent = visible.slice(emittedVisible); + emittedVisible = visible.length; + emittedText = true; + onEvent({ type: "text", value: deltaContent }); + } + } + } + } + // Keep unsigned reasoning upstream-only for Anthropic clients. if (update.has(14)) { + flushThinkingFallback(onEvent); finished = true; onEvent({ type: "done" }); } } + // KvServerMessage (field 4): get/set blob. Ack so the stream proceeds. + if (serverMessage.has(4)) { + const kv = decodeMessage(serverMessage.get(4)[0].value); + const kvId = kv.get(1)?.[0]?.value || 0; + const metadata = kv.get(4)?.[0]?.value || null; + if (kv.has(2)) { + session.write(encodeKvClientMessage(kvId, 2, agentMessage(1, new Uint8Array()), metadata)); + } else if (kv.has(3)) { + session.write(encodeKvClientMessage(kvId, 3, new Uint8Array(), metadata)); + } + } + // AgentService requests IDE context before producing a response. - // Return an empty context; 9router is not coupled to an editor. if (serverMessage.has(2)) { const execRequest = decodeMessage(serverMessage.get(2)[0].value); if (execRequest.has(10)) { - session.write(createRequestContextResponse()); + log?.info?.("CURSOR", "AgentService request_context ack"); + session.write(createRequestContextResponse(execRequest)); + } else if (execRequest.has(11)) { + const mcp = decodeMcpArgs(execRequest.get(11)[0].value); + const name = mcp.toolName || mcp.name; + if (name) { + log?.info?.("CURSOR", `AgentService MCP tool_call ${name}`); + finished = true; + onEvent({ + type: "tool_call", + value: { + id: mcp.toolCallId || `call_${crypto.randomUUID()}`, + name, + arguments: JSON.stringify(mcp.args || {}), + }, + }); + onEvent({ type: "done", finishReason: "tool_calls" }); + } else { + debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`); + finished = true; + onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" }); + } } else { - // Every other ExecServerMessage variant is an editor-backed tool - // (shell, read, write, …) that 9router cannot service. Fail the - // turn rather than narrating protocol state as assistant text. - debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`); - finished = true; - onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" }); + // Auto/Composer often probe IDE builtins (shell/read/…). Reject + // them so the model can continue with MCP tools or a text answer + // instead of stalling the h2 stream. + const rejection = rejectExecRequest(execRequest); + if (rejection) { + log?.info?.("CURSOR", `AgentService rejected IDE exec fields=${[...execRequest.keys()].join(",")}`); + session.write(rejection); + } else { + debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`); + finished = true; + onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" }); + } } } }); @@ -586,7 +726,10 @@ export class CursorExecutor extends BaseExecutor { } finally { try { session.end(); } catch {} try { session.close(); } catch {} - if (!finished) onEvent({ type: "done" }); + if (!finished) { + flushThinkingFallback(onEvent); + onEvent({ type: "done" }); + } } }; @@ -594,10 +737,21 @@ export class CursorExecutor extends BaseExecutor { let content = ""; let reasoning = ""; let agentError = null; + const toolCalls = []; + let finishReason = "stop"; await consume((event) => { if (event.type === "text") content += event.value; else if (event.type === "thinking") reasoning += event.value; + else if (event.type === "tool_call") { + toolCalls.push({ + id: event.value.id, + type: "function", + function: { name: event.value.name, arguments: event.value.arguments }, + }); + finishReason = "tool_calls"; + } else if (event.type === "error") agentError = event.value; + else if (event.type === "done" && event.finishReason) finishReason = event.finishReason; }); if (agentError) { return { @@ -611,13 +765,19 @@ export class CursorExecutor extends BaseExecutor { responseFormat: FORMATS.OPENAI, }; } + const message = { + role: "assistant", + content: content || null, + ...(reasoning ? { reasoning_content: reasoning } : {}), + ...(toolCalls.length ? { tool_calls: toolCalls } : {}), + }; return { response: new Response(JSON.stringify({ id: responseId, object: "chat.completion", created, model, - choices: [{ index: 0, message: { role: "assistant", content: content || null, ...(reasoning ? { reasoning_content: reasoning } : {}) }, finish_reason: "stop" }], + choices: [{ index: 0, message, finish_reason: finishReason }], usage: estimateUsage(body, content.length, FORMATS.OPENAI), }), { headers: { "Content-Type": "application/json" } }), url, @@ -635,6 +795,18 @@ export class CursorExecutor extends BaseExecutor { controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: { content: event.value } }))); } else if (event.type === "thinking") { controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: { reasoning_content: event.value } }))); + } else if (event.type === "tool_call") { + controller.enqueue(encoder.encode(chatChunkSse({ + id: responseId, created, model, + delta: { + tool_calls: [{ + index: 0, + id: event.value.id, + type: "function", + function: { name: event.value.name, arguments: event.value.arguments }, + }], + }, + }))); } else if (event.type === "error") { // An SSE error frame, not a content delta: a protocol failure must not // be rendered to the user as the assistant's reply, and downstream @@ -643,7 +815,10 @@ export class CursorExecutor extends BaseExecutor { controller.enqueue(encoder.encode(SSE_DONE)); controller.close(); } else if (event.type === "done") { - controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: {}, finishReason: "stop" }))); + controller.enqueue(encoder.encode(chatChunkSse({ + id: responseId, created, model, delta: {}, + finishReason: event.finishReason || "stop", + }))); controller.enqueue(encoder.encode(SSE_DONE)); controller.close(); } @@ -664,9 +839,9 @@ export class CursorExecutor extends BaseExecutor { } async execute({ model, body, stream, credentials, signal, log, proxyOptions = null }) { - if (isAgentTextRequest(body)) { + if (isAgentCapableRequest(body)) { try { - return await this.executeAgent({ model, body, stream, credentials, signal }); + return await this.executeAgent({ model, body, stream, credentials, signal, log }); } catch (error) { return { response: new Response(JSON.stringify({ diff --git a/open-sse/handlers/chatCore.js b/open-sse/handlers/chatCore.js index 56a9d69f..0db37ba1 100644 --- a/open-sse/handlers/chatCore.js +++ b/open-sse/handlers/chatCore.js @@ -116,6 +116,19 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred } } + // Per-request opt-out: client can bypass all token savers via header + const tokenSaverEnabled = clientRawRequest?.headers?.[TOKEN_SAVER_HEADER]?.toLowerCase() !== "off"; + + // Cursor's translator rewrites tool_result into user text, so RTK must run on + // the source body before translation. Every other pair translates the tool + // shapes 1:1 — keep the post-translate pass there so those providers are + // untouched (and a retry never re-compresses an already-compressed body). + const preTranslateRtk = provider === "cursor" + ? compressMessages(body, tokenSaverEnabled && rtkEnabled) + : null; + const preTranslateRtkLine = formatRtkLog(preTranslateRtk); + if (preTranslateRtkLine) console.log(preTranslateRtkLine); + const clientRequestedStreaming = body.stream === true || sourceFormat === FORMATS.ANTIGRAVITY || sourceFormat === FORMATS.GEMINI || sourceFormat === FORMATS.GEMINI_CLI; const providerRequiresStreaming = PROVIDERS[provider]?.forceStream === true; let stream = providerRequiresStreaming ? true : (body.stream !== false); @@ -252,13 +265,8 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred translatedBody.tools = defaultClaudeToolType(translatedBody.tools); } - // Per-request opt-out: client can bypass all token savers via header - const tokenSaverEnabled = clientRawRequest?.headers?.[TOKEN_SAVER_HEADER]?.toLowerCase() !== "off"; - - // RTK: compress tool_result content - const rtkStats = compressMessages(translatedBody, tokenSaverEnabled && rtkEnabled); - const rtkLine = formatRtkLog(rtkStats); - if (rtkLine) console.log(rtkLine); + // RTK: compress tool_result content. Skipped when already done pre-translate. + const rtkStats = preTranslateRtk || compressMessages(translatedBody, tokenSaverEnabled && rtkEnabled); // Headroom: optional external proxy compression; fail open if proxy is absent. const headroomDiagnostics = {}; @@ -275,6 +283,8 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred // Token-saver flags accumulator for the single "⚙" log line below. const xf = []; + if (rtkStats?.hits?.length) xf.push(`RTK:${rtkStats.hits.length}`); + // Caveman: inject terse-style system prompt if (tokenSaverEnabled && cavemanEnabled && cavemanLevel) { injectCaveman(translatedBody, finalFormat, cavemanLevel); diff --git a/open-sse/rtk/index.js b/open-sse/rtk/index.js index dd1e5018..b8488baa 100644 --- a/open-sse/rtk/index.js +++ b/open-sse/rtk/index.js @@ -1,5 +1,5 @@ // RTK port: compress tool_result content in LLM request bodies -// Injected at the top of translateRequest (before any format translation) +// Applied in chatCore on the source-format body, before translateRequest. import { RAW_CAP, MIN_COMPRESS_SIZE } from "./constants.js"; import { autoDetectFilter } from "./autodetect.js"; import { safeApply } from "./applyFilter.js"; diff --git a/open-sse/utils/cursorProtobuf.js b/open-sse/utils/cursorProtobuf.js index c870921b..47b3e0ba 100644 --- a/open-sse/utils/cursorProtobuf.js +++ b/open-sse/utils/cursorProtobuf.js @@ -218,6 +218,12 @@ export function encodeField(fieldNum, wireType, value) { return concatArrays(tagBytes, lengthBytes, dataBytes); } + if (wireType === WIRE_TYPE.FIXED64) { + const buf = Buffer.alloc(8); + buf.writeDoubleLE(Number(value)); + return concatArrays(tagBytes, buf); + } + return new Uint8Array(0); } @@ -887,6 +893,211 @@ export function extractTextFromResponse(payload) { } } +// ==================== AGENT SERVICE (google.protobuf.Value + MCP) ==================== + +const PB_VALUE = { NULL: 1, NUMBER: 2, STRING: 3, BOOL: 4, STRUCT: 5, LIST: 6 }; +const PB_STRUCT_FIELDS = 1; +const PB_MAP_KEY = 1; +const PB_MAP_VALUE = 2; +const PB_LIST_VALUES = 1; + +const MTD_NAME = 1; +const MTD_DESCRIPTION = 2; +const MTD_INPUT_SCHEMA = 3; +const MTD_PROVIDER = 4; +const MTD_TOOL_NAME = 5; + +const MCP_TOOLS_TOOL = 1; + +const MCP_ARGS_NAME = 1; +const MCP_ARGS_ENTRY = 2; +const MCP_ARGS_CALL_ID = 3; +const MCP_ARGS_TOOL_NAME = 5; + +const MCR_SUCCESS = 1; +const MCR_ERROR = 2; +const MCR_TOOL_NOT_FOUND = 5; +const MCS_CONTENT = 1; +const MCS_IS_ERROR = 2; +const MCC_TEXT = 1; +const MCC_IMAGE = 2; +const MTC_TEXT = 1; +const MIC_DATA = 1; +const MIC_MIME = 2; +const MER_MESSAGE = 1; +const TNF_NAME = 1; + +function asBytes(value) { + if (!value) return Buffer.alloc(0); + return Buffer.isBuffer(value) ? value : Buffer.from(value); +} + +/** + * Encode a JS value as google.protobuf.Value (oneof body, no outer tag). + */ +export function encodeAgentValue(value) { + if (value === null || value === undefined) { + return encodeField(PB_VALUE.NULL, WIRE_TYPE.VARINT, 0); + } + if (typeof value === "boolean") { + return encodeField(PB_VALUE.BOOL, WIRE_TYPE.VARINT, value ? 1 : 0); + } + if (typeof value === "number") { + return encodeField(PB_VALUE.NUMBER, WIRE_TYPE.FIXED64, value); + } + if (typeof value === "string") { + return encodeField(PB_VALUE.STRING, WIRE_TYPE.LEN, value); + } + if (Array.isArray(value)) { + const items = value.map((item) => encodeField(PB_LIST_VALUES, WIRE_TYPE.LEN, encodeAgentValue(item))); + return encodeField(PB_VALUE.LIST, WIRE_TYPE.LEN, concatArrays(...items)); + } + if (typeof value === "object") { + const entries = Object.entries(value).map(([key, val]) => encodeField( + PB_STRUCT_FIELDS, + WIRE_TYPE.LEN, + concatArrays( + encodeField(PB_MAP_KEY, WIRE_TYPE.LEN, key), + encodeField(PB_MAP_VALUE, WIRE_TYPE.LEN, encodeAgentValue(val)), + ), + )); + return encodeField(PB_VALUE.STRUCT, WIRE_TYPE.LEN, concatArrays(...entries)); + } + return encodeField(PB_VALUE.STRING, WIRE_TYPE.LEN, String(value)); +} + +/** + * Decode google.protobuf.Value bytes back to a JS value. + */ +export function decodeAgentValue(bytes) { + const fields = decodeMessage(asBytes(bytes)); + if (fields.has(PB_VALUE.NULL)) return null; + if (fields.has(PB_VALUE.BOOL)) return fields.get(PB_VALUE.BOOL)[0].value !== 0; + if (fields.has(PB_VALUE.NUMBER)) { + return asBytes(fields.get(PB_VALUE.NUMBER)[0].value).readDoubleLE(0); + } + if (fields.has(PB_VALUE.STRING)) { + return asBytes(fields.get(PB_VALUE.STRING)[0].value).toString("utf8"); + } + if (fields.has(PB_VALUE.STRUCT)) { + const result = {}; + for (const entry of decodeMessage(asBytes(fields.get(PB_VALUE.STRUCT)[0].value)).get(PB_STRUCT_FIELDS) || []) { + const pair = decodeMessage(asBytes(entry.value)); + const key = asBytes(pair.get(PB_MAP_KEY)?.[0]?.value).toString("utf8"); + if (key) result[key] = decodeAgentValue(pair.get(PB_MAP_VALUE)?.[0]?.value); + } + return result; + } + if (fields.has(PB_VALUE.LIST)) { + return (decodeMessage(asBytes(fields.get(PB_VALUE.LIST)[0].value)).get(PB_LIST_VALUES) || []) + .map((item) => decodeAgentValue(item.value)); + } + return null; +} + +function toolNameAndSchema(tool) { + const fn = tool?.function || tool || {}; + return { + name: fn.name || tool?.name || "", + description: fn.description || tool?.description || "", + schema: fn.parameters || tool?.parameters || tool?.inputSchema || tool?.input_schema || {}, + }; +} + +/** + * Encode agent.v1.McpToolDefinition body (name, description, Value schema, provider, tool_name). + */ +export function encodeMcpToolDefinition(tool) { + const { name, description, schema } = toolNameAndSchema(tool); + return concatArrays( + encodeField(MTD_NAME, WIRE_TYPE.LEN, name), + encodeField(MTD_DESCRIPTION, WIRE_TYPE.LEN, description), + encodeField(MTD_INPUT_SCHEMA, WIRE_TYPE.LEN, encodeAgentValue(schema)), + encodeField(MTD_PROVIDER, WIRE_TYPE.LEN, "9router"), + encodeField(MTD_TOOL_NAME, WIRE_TYPE.LEN, name), + ); +} + +/** + * Encode AgentRunRequest.mcp_tools: repeated McpToolDefinition under field 1. + */ +export function encodeMcpTools(tools = []) { + if (!tools?.length) return new Uint8Array(); + return concatArrays( + ...tools.map((tool) => encodeField(MCP_TOOLS_TOOL, WIRE_TYPE.LEN, encodeMcpToolDefinition(tool))), + ); +} + +/** + * Decode agent.v1.McpArgs (name, typed args map, toolCallId, toolName). + */ +export function decodeMcpArgs(bytes) { + const msg = decodeMessage(asBytes(bytes)); + const args = {}; + for (const entry of msg.get(MCP_ARGS_ENTRY) || []) { + const pair = decodeMessage(asBytes(entry.value)); + const key = asBytes(pair.get(PB_MAP_KEY)?.[0]?.value).toString("utf8"); + if (key) args[key] = decodeAgentValue(pair.get(PB_MAP_VALUE)?.[0]?.value); + } + const read = (field) => asBytes(msg.get(field)?.[0]?.value).toString("utf8"); + return { + name: read(MCP_ARGS_NAME), + toolCallId: read(MCP_ARGS_CALL_ID), + toolName: read(MCP_ARGS_TOOL_NAME), + args, + }; +} + +function encodeMcpTextItem(text) { + return encodeField( + MCS_CONTENT, + WIRE_TYPE.LEN, + encodeField(MCC_TEXT, WIRE_TYPE.LEN, encodeField(MTC_TEXT, WIRE_TYPE.LEN, text)), + ); +} + +function encodeMcpImageItem(image) { + const data = image?.data || image || new Uint8Array(); + const mimeType = image?.mimeType || "application/octet-stream"; + return encodeField( + MCS_CONTENT, + WIRE_TYPE.LEN, + encodeField( + MCC_IMAGE, + WIRE_TYPE.LEN, + concatArrays( + encodeField(MIC_DATA, WIRE_TYPE.LEN, data), + encodeField(MIC_MIME, WIRE_TYPE.LEN, mimeType), + ), + ), + ); +} + +export function encodeMcpResultSuccess({ textItems = [], imageItems = [], isError = false } = {}) { + const success = concatArrays( + ...textItems.map(encodeMcpTextItem), + ...imageItems.map(encodeMcpImageItem), + encodeField(MCS_IS_ERROR, WIRE_TYPE.VARINT, isError ? 1 : 0), + ); + return encodeField(MCR_SUCCESS, WIRE_TYPE.LEN, success); +} + +export function encodeMcpResultError(message) { + return encodeField( + MCR_ERROR, + WIRE_TYPE.LEN, + encodeField(MER_MESSAGE, WIRE_TYPE.LEN, String(message || "")), + ); +} + +export function encodeMcpResultToolNotFound(name) { + return encodeField( + MCR_TOOL_NOT_FOUND, + WIRE_TYPE.LEN, + encodeField(TNF_NAME, WIRE_TYPE.LEN, String(name || "")), + ); +} + // ==================== EXPORTS ==================== export default { @@ -900,5 +1111,13 @@ export default { decodeField, decodeMessage, parseConnectRPCFrame, - extractTextFromResponse + extractTextFromResponse, + encodeAgentValue, + decodeAgentValue, + encodeMcpToolDefinition, + encodeMcpTools, + decodeMcpArgs, + encodeMcpResultSuccess, + encodeMcpResultError, + encodeMcpResultToolNotFound, }; diff --git a/tests/unit/cursor-agent-exec-request.test.js b/tests/unit/cursor-agent-exec-request.test.js index 347e159f..18389c18 100644 --- a/tests/unit/cursor-agent-exec-request.test.js +++ b/tests/unit/cursor-agent-exec-request.test.js @@ -18,6 +18,18 @@ function textFrame(text) { return Buffer.from(wrapConnectRPCFrame(encodeField(1, LEN, update))); } +// InteractionUpdate.thinking_delta (field 4) + turn_ended (field 14). +function thinkingFrame(text) { + const thinkingPart = Buffer.from(encodeField(1, LEN, text)); + const update = Buffer.from(encodeField(4, LEN, thinkingPart)); + return Buffer.from(wrapConnectRPCFrame(encodeField(1, LEN, update))); +} + +function turnEndedFrame() { + const update = Buffer.from(encodeField(14, LEN, new Uint8Array())); + return Buffer.from(wrapConnectRPCFrame(encodeField(1, LEN, update))); +} + function stubAgentSession(executor, frames) { const written = []; const queue = [...frames]; @@ -48,12 +60,12 @@ function parseSSE(text) { .map((data) => JSON.parse(data)); } -async function runAgent({ frames, stream }) { +async function runAgent({ frames, stream, model = "gpt-5.2", tools }) { const executor = new CursorExecutor(); const written = stubAgentSession(executor, frames); const result = await executor.executeAgent({ - model: "gpt-5.2", - body: { messages: [{ role: "user", content: "hi" }] }, + model, + body: { messages: [{ role: "user", content: "hi" }], ...(tools ? { tools } : {}) }, stream, credentials, }); @@ -73,32 +85,45 @@ describe("CursorExecutor AgentService exec_request handling", () => { expect(content).toBe("hello"); }); + it("does not echo client tools on the request_context ack", async () => { + const { written, result } = await runAgent({ + tools: [{ function: { name: "read_file", parameters: { type: "object" } } }], + frames: [execRequestFrame(10), textFrame("hello")], + stream: true, + }); + + expect(written.length).toBe(2); + expect(written[1].toString("utf8")).not.toContain("read_file"); + const content = parseSSE(await result.response.text()) + .map((e) => e.choices?.[0]?.delta?.content || "") + .join(""); + expect(content).toBe("hello"); + }); + it("does not render an unsupported exec request as assistant content", async () => { - const { result } = await runAgent({ - frames: [textFrame("partial answer"), execRequestFrame(2)], + const { result, written } = await runAgent({ + frames: [textFrame("partial answer"), execRequestFrame(2), textFrame(" more")], stream: true, }); const body = await result.response.text(); - expect(body).not.toContain("unsupported IDE tool\\n"); + expect(body).not.toContain("unsupported IDE tool"); const events = parseSSE(body); const content = events.map((e) => e.choices?.[0]?.delta?.content || "").join(""); - expect(content).toBe("partial answer"); - - const errorEvent = events.find((e) => e.error); - expect(errorEvent?.error?.message).toContain("unsupported IDE tool"); - expect(events.some((e) => e.choices?.[0]?.finish_reason === "stop")).toBe(false); + expect(content).toBe("partial answer more"); + expect(events.some((e) => e.error)).toBe(false); + expect(written.length).toBe(2); // run frame + IDE rejection }); - it("drops frames batched behind an unsupported exec request in the same read", async () => { + it("still emits later text after rejecting an IDE exec in the same read", async () => { const { result } = await runAgent({ frames: [Buffer.concat([execRequestFrame(2), textFrame("late")])], stream: true, }); const body = await result.response.text(); - expect(body).toContain("unsupported IDE tool"); - expect(body).not.toContain("late"); + expect(body).not.toContain("unsupported IDE tool"); + expect(body).toContain("late"); }); it("returns a non-200 error body for an unsupported exec request when not streaming", async () => { @@ -111,4 +136,32 @@ describe("CursorExecutor AgentService exec_request handling", () => { const payload = await result.response.json(); expect(payload.error.message).toContain("unsupported IDE tool"); }); + + it("streams Composer visible content from thinking_delta after ", async () => { + const { result } = await runAgent({ + model: "composer-2.5", + frames: [ + thinkingFrame("private reasoning that must not leakOK"), + turnEndedFrame(), + ], + stream: true, + }); + + const events = parseSSE(await result.response.text()); + const content = events.map((e) => e.choices?.[0]?.delta?.content || "").join(""); + expect(content).toBe("OK"); + expect(JSON.stringify(events)).not.toContain("private reasoning"); + }); + + it("flushes Grok thinking as visible content when the turn has no text_delta", async () => { + const { result } = await runAgent({ + model: "grok-4.5", + frames: [thinkingFrame("hello from grok"), turnEndedFrame()], + stream: true, + }); + + const events = parseSSE(await result.response.text()); + const content = events.map((e) => e.choices?.[0]?.delta?.content || "").join(""); + expect(content).toBe("hello from grok"); + }); }); diff --git a/tests/unit/cursor-agent-proto.test.js b/tests/unit/cursor-agent-proto.test.js index 2d571aba..b112cc58 100644 --- a/tests/unit/cursor-agent-proto.test.js +++ b/tests/unit/cursor-agent-proto.test.js @@ -246,6 +246,15 @@ describe("Cursor AgentService executor helpers (cursor.js)", () => { const run = decodeMessage(clientMsg.get(1)[0].value); expect(run.has(2)).toBe(true); // action expect(run.has(9)).toBe(true); // requested_model + // custom_system_prompt (field 8) makes AgentService return an empty turn. + expect(run.has(8)).toBe(false); + expect(run.has(3)).toBe(true); // ModelDetails — required for thinking variants + const action = decodeMessage(run.get(2)[0].value); + const userAction = decodeMessage(action.get(1)[0].value); + const userMessage = decodeMessage(userAction.get(1)[0].value); + const userText = Buffer.from(userMessage.get(1)[0].value).toString("utf8"); + expect(userText).toContain("be brief"); + expect(userText).toContain("hi"); }); it("encodes mcp_tools (field 4) when tools are provided", () => { diff --git a/tests/unit/rtk-cursor-pretranslate.test.js b/tests/unit/rtk-cursor-pretranslate.test.js new file mode 100644 index 00000000..601fcfec --- /dev/null +++ b/tests/unit/rtk-cursor-pretranslate.test.js @@ -0,0 +1,131 @@ +import { describe, it, expect, vi, beforeEach } from "vitest"; + +const { executeMock } = vi.hoisted(() => ({ + executeMock: vi.fn(), +})); + +vi.mock("../../open-sse/executors/index.js", () => ({ + getExecutor: () => ({ + noAuth: true, + execute: executeMock, + }), +})); + +vi.mock("../../open-sse/utils/requestLogger.js", () => ({ + createRequestLogger: async () => ({ + logClientRawRequest: vi.fn(), + logRawRequest: vi.fn(), + logTargetRequest: vi.fn(), + logProviderResponse: vi.fn(), + logConvertedResponse: vi.fn(), + logError: vi.fn(), + }), +})); + +vi.mock("../../open-sse/utils/stream.js", () => ({ + COLORS: { red: "", reset: "" }, + createPassthroughStreamWithLogger: vi.fn(() => new TransformStream()), +})); + +vi.mock("@/lib/usageDb.js", () => ({ + trackPendingRequest: vi.fn(), + appendRequestLog: vi.fn(async () => {}), + saveRequestDetail: vi.fn(async () => {}), +})); + +const { handleChatCore } = await import("../../open-sse/handlers/chatCore.js"); + +function makeLongDiff() { + const lines = ["diff --git a/foo.js b/foo.js", "index abc..def 100644", "--- a/foo.js", "+++ b/foo.js", "@@ -1,3 +1,200 @@"]; + for (let i = 0; i < 200; i++) lines.push(`+added line ${i} UNIQUE_PADDING_${i} ${"x".repeat(20)}`); + return lines.join("\n"); +} + +describe("token savers on Cursor (pre-translate RTK)", () => { + beforeEach(() => { + vi.clearAllMocks(); + global.fetch = vi.fn(async (url, init) => { + if (String(url).includes("/v1/compress")) { + const payload = JSON.parse(init.body); + return new Response(JSON.stringify({ + messages: payload.messages, + tokens_before: 8000, + tokens_after: 2500, + tokens_saved: 5500, + }), { status: 200, headers: { "content-type": "application/json" } }); + } + throw new Error(`unexpected fetch: ${url}`); + }); + executeMock.mockResolvedValue({ + response: new Response(JSON.stringify({ + id: "chatcmpl-test", + object: "chat.completion", + choices: [{ message: { role: "assistant", content: "ok" }, finish_reason: "stop", index: 0 }], + }), { status: 200, headers: { "content-type": "application/json" } }), + url: "https://api2.cursor.sh/agent", + headers: {}, + transformedBody: null, + }); + }); + + it("compresses role:tool git diffs before openai→cursor rewrite, then injects Headroom/Caveman/Ponytail", async () => { + const diff = makeLongDiff(); + const log = { debug: vi.fn(), info: vi.fn(), warn: vi.fn(), line: vi.fn() }; + + await handleChatCore({ + body: { + model: "cu/default", + stream: false, + messages: [ + { role: "system", content: "hi" }, + { role: "user", content: "run git diff" }, + { + role: "assistant", + content: null, + tool_calls: [{ id: "call_1", type: "function", function: { name: "Bash", arguments: JSON.stringify({ command: "git diff" }) } }], + }, + { role: "tool", tool_call_id: "call_1", content: diff }, + { role: "user", content: "summarize" }, + ], + }, + modelInfo: { provider: "cursor", model: "default" }, + credentials: { apiKey: "test-key", providerSpecificData: {} }, + log, + connectionId: "test-conn", + rtkEnabled: true, + headroomEnabled: true, + headroomUrl: "http://localhost:8787", + cavemanEnabled: true, + cavemanLevel: "full", + ponytailEnabled: true, + ponytailLevel: "full", + clientRawRequest: { + endpoint: "/v1/chat/completions", + body: { model: "cu/default" }, + headers: { accept: "application/json" }, + }, + }); + + expect(executeMock).toHaveBeenCalled(); + const dispatched = executeMock.mock.calls[0][0].body; + const blob = JSON.stringify(dispatched.messages); + + expect(dispatched.messages.some((m) => m.role === "tool")).toBe(false); + expect(blob).toContain(""); + expect(blob).toContain("lines truncated"); + expect(blob).not.toContain("UNIQUE_PADDING_150"); + expect(blob).toContain("lazy senior developer"); + expect(blob).toMatch(/Respond like a caveman|drop filler|ACTIVE EVERY RESPONSE/i); + + expect(global.fetch).toHaveBeenCalledWith( + "http://localhost:8787/v1/compress", + expect.any(Object) + ); + + const xf = log.line.mock.calls.find((c) => c[1] === "⚙"); + expect(xf, "expected ⚙ saver log").toBeTruthy(); + expect(xf[2]).toContain("RTK:"); + expect(xf[2]).toContain("CAVEMAN:full"); + expect(xf[2]).toContain("PONYTAIL:full"); + }); +});