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.
This commit is contained in:
Amir Seify
2026-09-16 21:31:22 +07:00
parent 5c217d34f3
commit c933eefc27
9 changed files with 670 additions and 71 deletions

View File

@@ -7,13 +7,16 @@ import {
wrapConnectRPCFrame,
decodeMessage,
parseConnectRPCFrame,
extractTextFromResponse
extractTextFromResponse,
encodeMcpTools,
decodeMcpArgs,
} from "../utils/cursorProtobuf.js";
import { buildCursorHeaders } from "../utils/cursorChecksum.js";
import { estimateUsage } from "../utils/usageTracking.js";
import { SSE_DONE, SSE_HEADERS } from "../utils/sseConstants.js";
import { chatChunkSse, sseChunk } from "../utils/sse.js";
import { FORMATS } from "../translator/formats.js";
import { ROLE, OPENAI_BLOCK } from "../translator/schema/index.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js";
import zlib from "zlib";
import crypto from "crypto";
@@ -65,55 +68,74 @@ function textFromContent(content) {
if (typeof content === "string") return content;
if (!Array.isArray(content)) return "";
return content
.filter((part) => part?.type === "text" && typeof part.text === "string")
.filter((part) => part?.type === OPENAI_BLOCK.TEXT && typeof part.text === "string")
.map((part) => part.text)
.join("\n");
}
function isAgentTextRequest(body) {
// Many compatible clients always attach their built-in tool schemas, even
// for a normal text turn. Cursor's retired ChatService rejects those
// requests; AgentService can still answer the text turn, so ignore schemas
// here. A real tool-call/result conversation is kept on the legacy path
// until its AgentService tool protocol is implemented.
return Array.isArray(body?.messages) && body.messages.every((message) => {
if (message?.tool_calls?.length || message?.role === "tool") return false;
return typeof message?.content === "string"
|| Array.isArray(message?.content) && message.content.every((part) => part?.type === "text");
function isTextPart(part) {
return !part || part.type === OPENAI_BLOCK.TEXT || typeof part === "string";
}
export function isAgentCapableRequest(body) {
// ChatService rejects auto/composer and most thinking variants. AgentService
// can answer text turns (including declared tool schemas) and tool-call
// history. Image parts still need the legacy protobuf path.
if (!Array.isArray(body?.messages) || body.messages.length === 0) return false;
return body.messages.every((message) => {
if (Array.isArray(message?.content)) return message.content.every(isTextPart);
return message?.content == null || typeof message.content === "string";
});
}
function encodeHistoryMessage(message) {
const content = textFromContent(message?.content);
if (!content) return null;
const extras = [];
if (message?.role === ROLE.ASSISTANT && message.tool_calls?.length) {
for (const tc of message.tool_calls) {
extras.push(`[tool_call id=${tc.id || ""} name=${tc.function?.name || "tool"} args=${tc.function?.arguments || "{}"}]`);
}
}
if (message?.role === ROLE.TOOL) {
extras.push(`[tool_result id=${message.tool_call_id || ""}]`);
}
const textBody = [content, ...extras].filter(Boolean).join("\n");
if (!textBody) return null;
// ConversationHistoryMessage.user / .assistant -> repeated content -> text.
const text = agentString(1, content);
if (message.role === "assistant") {
const text = agentString(1, textBody);
if (message.role === ROLE.ASSISTANT) {
return agentMessage(2, agentMessage(1, agentMessage(1, text)));
}
return agentMessage(1, agentMessage(1, agentMessage(1, text)));
}
function buildAgentRunFrame(messages, model) {
export function buildAgentRunFrame(messages, model, tools = []) {
// custom_system_prompt (RunRequest field 8) makes AgentService return an
// empty turn. Fold system text into the current user message instead.
const system = messages
.filter((message) => message?.role === "system")
.filter((message) => message?.role === ROLE.SYSTEM)
.map((message) => textFromContent(message.content))
.filter(Boolean)
.join("\n\n");
const chatMessages = messages.filter((message) => message?.role !== "system");
const currentIndex = [...chatMessages].map((message) => message?.role).lastIndexOf("user");
const chatMessages = messages.filter((message) => message?.role !== ROLE.SYSTEM);
const currentIndex = [...chatMessages].map((message) => message?.role).lastIndexOf(ROLE.USER);
const current = currentIndex >= 0 ? chatMessages[currentIndex] : chatMessages.at(-1);
const history = chatMessages
.slice(0, currentIndex >= 0 ? currentIndex : -1)
.map(encodeHistoryMessage)
.filter(Boolean);
const userText = textFromContent(current?.content) || "Continue.";
const rawUser = textFromContent(current?.content) || "Continue.";
const userText = system ? `${system}\n\n${rawUser}` : rawUser;
// agent.v1.UserMessageAction.user_message and its optional history.
// selected_context (3) + mode=1 (4) match cursor-agent's wire format; without
// them the server may accept the RPC and stream an empty turn.
const userMessage = concatBuffers(
agentString(1, userText),
agentString(2, crypto.randomUUID()),
agentMessage(3, new Uint8Array()),
encodeField(4, PROTOBUF_VARINT, 1),
);
const conversationHistory = history.length
? concatBuffers(...history.map((entry) => agentMessage(1, entry)))
@@ -124,11 +146,20 @@ function buildAgentRunFrame(messages, model) {
);
const conversationAction = agentMessage(1, userAction);
const requestedModel = concatBuffers(agentString(1, model), agentBool(7, true));
// ModelDetails (field 3): thinking variants (Composer, Grok, *-thinking)
// return an empty turn when only RequestedModel (field 9) is set.
const modelDetails = concatBuffers(
agentString(1, model),
agentString(3, model),
agentString(4, model),
);
const mcpTools = encodeMcpTools(tools);
const runRequest = concatBuffers(
// An empty ConversationStateStructure starts a fresh local agent session.
agentMessage(1, new Uint8Array()),
agentMessage(2, conversationAction),
...(system ? [agentString(8, system)] : []),
agentMessage(3, modelDetails),
...(mcpTools.length ? [agentMessage(4, mcpTools)] : []),
agentMessage(9, requestedModel),
);
@@ -157,13 +188,51 @@ function decodeAgentFrames(buffer, onFrame) {
return pending;
}
function createRequestContextResponse() {
// AgentService asks every run for client context. 9router has no IDE file
// context, so acknowledge with an empty RequestContext.
function execIds(execRequest) {
const id = Number(execRequest?.get(1)?.[0]?.value || 0);
const execId = extractAgentString(execRequest, 15);
return { id, execId };
}
function wrapExecClientMessage(execMsgId, execId, resultField, resultPayload) {
const parts = [];
if (execMsgId) parts.push(encodeField(1, PROTOBUF_VARINT, execMsgId));
parts.push(agentString(15, execId || ""));
parts.push(encodeField(resultField, PROTOBUF_LEN, resultPayload || new Uint8Array()));
return wrapConnectRPCFrame(agentMessage(2, concatBuffers(...parts)));
}
function createRequestContextResponse(execRequest) {
// Tools already go out on AgentRunRequest.mcp_tools. Echoing them again on
// this ack makes AgentService stall silently (0 SSE bytes until abort).
const { id, execId } = execIds(execRequest);
const requestContextSuccess = agentMessage(1, new Uint8Array());
const requestContextResult = agentMessage(1, requestContextSuccess);
const execClientMessage = agentMessage(10, requestContextResult);
return wrapConnectRPCFrame(agentMessage(2, execClientMessage));
return wrapExecClientMessage(id, execId, 10, requestContextResult);
}
// ExecServerMessage variant → ExecClientMessage result field (same numbers).
const EXEC_RESULT_FIELD = {
2: 2, 3: 3, 4: 4, 5: 5, 7: 7, 8: 8, 9: 9, 16: 16, 20: 20, 23: 23,
};
function rejectExecRequest(execRequest) {
const { id, execId } = execIds(execRequest);
const variant = [...(execRequest?.keys?.() || [])].find((field) => field !== 1 && field !== 15);
const resultField = EXEC_RESULT_FIELD[variant];
if (!resultField) return null;
// Diagnostics has no rejected variant — empty success unblocks the stream.
if (variant === 9) return wrapExecClientMessage(id, execId, 9, new Uint8Array());
const rejected = agentMessage(2, agentString(2, "Tool not available in this environment. Use the MCP tools provided instead."));
return wrapExecClientMessage(id, execId, resultField, rejected);
}
function encodeKvClientMessage(kvId, resultField, resultPayload, metadata) {
const parts = [];
if (kvId) parts.push(encodeField(1, PROTOBUF_VARINT, kvId));
parts.push(encodeField(resultField, PROTOBUF_LEN, resultPayload || new Uint8Array()));
if (metadata && metadata.length) parts.push(encodeField(4, PROTOBUF_LEN, metadata));
return wrapConnectRPCFrame(agentMessage(3, concatBuffers(...parts)));
}
const CURSOR_STREAM_DEBUG = process.env.CURSOR_STREAM_DEBUG === "1";
@@ -479,7 +548,7 @@ export class CursorExecutor extends BaseExecutor {
};
}
async executeAgent({ model, body, stream, credentials, signal }) {
async executeAgent({ model, body, stream, credentials, signal, log }) {
const agentEndpoint = PROVIDER_OAUTH.cursor?.agentEndpoint;
if (!agentEndpoint) throw new Error("Cursor AgentService endpoint is not configured");
@@ -491,9 +560,10 @@ export class CursorExecutor extends BaseExecutor {
}
let session;
const tools = body.tools || [];
try {
session = this.openAgentHttp2Stream(url, headers, requestController.signal);
session.write(buildAgentRunFrame(body.messages || [], model));
session.write(buildAgentRunFrame(body.messages || [], model, tools));
} catch (error) {
throw new Error(`Cursor AgentService request failed: ${error.message}`);
}
@@ -533,8 +603,23 @@ export class CursorExecutor extends BaseExecutor {
// so strict clients such as Claude Code accept the completed stream.
const responseId = `chatcmpl-msg_${Date.now()}`;
const created = Math.floor(Date.now() / 1000);
const composerModel = isComposerModel(model);
let pending = Buffer.alloc(0);
let finished = false;
let thinkingAcc = "";
let emittedVisible = 0;
let emittedText = false;
const flushThinkingFallback = (onEvent) => {
if (emittedText || !thinkingAcc) return;
const fallback = composerModel
? visibleComposerContentFromThinking(thinkingAcc)
: thinkingAcc.trim();
if (fallback) {
emittedText = true;
onEvent({ type: "text", value: fallback });
}
};
const consume = async (onEvent) => {
try {
@@ -553,32 +638,87 @@ export class CursorExecutor extends BaseExecutor {
const update = decodeMessage(serverMessage.get(1)[0].value);
if (update.has(1)) {
const textDelta = extractAgentString(decodeMessage(update.get(1)[0].value), 1);
if (textDelta) onEvent({ type: "text", value: textDelta });
if (textDelta) {
emittedText = true;
onEvent({ type: "text", value: textDelta });
}
}
// Cursor's AgentService emits internal reasoning without the
// cryptographic signature required by Anthropic thinking blocks.
// Forwarding it makes strict Anthropic clients (Claude Code)
// discard or wait on an otherwise complete response. Keep the
// reasoning upstream-only and emit the normal answer text.
// thinking_delta (field 4). Composer (and some Grok variants) put
// the visible answer after </think> here and never send text_delta.
if (update.has(4)) {
const thinkingDelta = extractAgentString(decodeMessage(update.get(4)[0].value), 1);
if (thinkingDelta) {
thinkingAcc += thinkingDelta;
if (composerModel) {
const visible = visibleComposerContentFromThinking(thinkingAcc);
if (visible.length > emittedVisible) {
const deltaContent = visible.slice(emittedVisible);
emittedVisible = visible.length;
emittedText = true;
onEvent({ type: "text", value: deltaContent });
}
}
}
}
// Keep unsigned reasoning upstream-only for Anthropic clients.
if (update.has(14)) {
flushThinkingFallback(onEvent);
finished = true;
onEvent({ type: "done" });
}
}
// KvServerMessage (field 4): get/set blob. Ack so the stream proceeds.
if (serverMessage.has(4)) {
const kv = decodeMessage(serverMessage.get(4)[0].value);
const kvId = kv.get(1)?.[0]?.value || 0;
const metadata = kv.get(4)?.[0]?.value || null;
if (kv.has(2)) {
session.write(encodeKvClientMessage(kvId, 2, agentMessage(1, new Uint8Array()), metadata));
} else if (kv.has(3)) {
session.write(encodeKvClientMessage(kvId, 3, new Uint8Array(), metadata));
}
}
// AgentService requests IDE context before producing a response.
// Return an empty context; 9router is not coupled to an editor.
if (serverMessage.has(2)) {
const execRequest = decodeMessage(serverMessage.get(2)[0].value);
if (execRequest.has(10)) {
session.write(createRequestContextResponse());
log?.info?.("CURSOR", "AgentService request_context ack");
session.write(createRequestContextResponse(execRequest));
} else if (execRequest.has(11)) {
const mcp = decodeMcpArgs(execRequest.get(11)[0].value);
const name = mcp.toolName || mcp.name;
if (name) {
log?.info?.("CURSOR", `AgentService MCP tool_call ${name}`);
finished = true;
onEvent({
type: "tool_call",
value: {
id: mcp.toolCallId || `call_${crypto.randomUUID()}`,
name,
arguments: JSON.stringify(mcp.args || {}),
},
});
onEvent({ type: "done", finishReason: "tool_calls" });
} else {
debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`);
finished = true;
onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" });
}
} else {
// Every other ExecServerMessage variant is an editor-backed tool
// (shell, read, write, …) that 9router cannot service. Fail the
// turn rather than narrating protocol state as assistant text.
debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`);
finished = true;
onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" });
// Auto/Composer often probe IDE builtins (shell/read/…). Reject
// them so the model can continue with MCP tools or a text answer
// instead of stalling the h2 stream.
const rejection = rejectExecRequest(execRequest);
if (rejection) {
log?.info?.("CURSOR", `AgentService rejected IDE exec fields=${[...execRequest.keys()].join(",")}`);
session.write(rejection);
} else {
debugLog(`[CURSOR AGENT] Unsupported exec request fields: ${[...execRequest.keys()].join(",")}`);
finished = true;
onEvent({ type: "error", value: "Cursor AgentService requested an unsupported IDE tool" });
}
}
}
});
@@ -586,7 +726,10 @@ export class CursorExecutor extends BaseExecutor {
} finally {
try { session.end(); } catch {}
try { session.close(); } catch {}
if (!finished) onEvent({ type: "done" });
if (!finished) {
flushThinkingFallback(onEvent);
onEvent({ type: "done" });
}
}
};
@@ -594,10 +737,21 @@ export class CursorExecutor extends BaseExecutor {
let content = "";
let reasoning = "";
let agentError = null;
const toolCalls = [];
let finishReason = "stop";
await consume((event) => {
if (event.type === "text") content += event.value;
else if (event.type === "thinking") reasoning += event.value;
else if (event.type === "tool_call") {
toolCalls.push({
id: event.value.id,
type: "function",
function: { name: event.value.name, arguments: event.value.arguments },
});
finishReason = "tool_calls";
}
else if (event.type === "error") agentError = event.value;
else if (event.type === "done" && event.finishReason) finishReason = event.finishReason;
});
if (agentError) {
return {
@@ -611,13 +765,19 @@ export class CursorExecutor extends BaseExecutor {
responseFormat: FORMATS.OPENAI,
};
}
const message = {
role: "assistant",
content: content || null,
...(reasoning ? { reasoning_content: reasoning } : {}),
...(toolCalls.length ? { tool_calls: toolCalls } : {}),
};
return {
response: new Response(JSON.stringify({
id: responseId,
object: "chat.completion",
created,
model,
choices: [{ index: 0, message: { role: "assistant", content: content || null, ...(reasoning ? { reasoning_content: reasoning } : {}) }, finish_reason: "stop" }],
choices: [{ index: 0, message, finish_reason: finishReason }],
usage: estimateUsage(body, content.length, FORMATS.OPENAI),
}), { headers: { "Content-Type": "application/json" } }),
url,
@@ -635,6 +795,18 @@ export class CursorExecutor extends BaseExecutor {
controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: { content: event.value } })));
} else if (event.type === "thinking") {
controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: { reasoning_content: event.value } })));
} else if (event.type === "tool_call") {
controller.enqueue(encoder.encode(chatChunkSse({
id: responseId, created, model,
delta: {
tool_calls: [{
index: 0,
id: event.value.id,
type: "function",
function: { name: event.value.name, arguments: event.value.arguments },
}],
},
})));
} else if (event.type === "error") {
// An SSE error frame, not a content delta: a protocol failure must not
// be rendered to the user as the assistant's reply, and downstream
@@ -643,7 +815,10 @@ export class CursorExecutor extends BaseExecutor {
controller.enqueue(encoder.encode(SSE_DONE));
controller.close();
} else if (event.type === "done") {
controller.enqueue(encoder.encode(chatChunkSse({ id: responseId, created, model, delta: {}, finishReason: "stop" })));
controller.enqueue(encoder.encode(chatChunkSse({
id: responseId, created, model, delta: {},
finishReason: event.finishReason || "stop",
})));
controller.enqueue(encoder.encode(SSE_DONE));
controller.close();
}
@@ -664,9 +839,9 @@ export class CursorExecutor extends BaseExecutor {
}
async execute({ model, body, stream, credentials, signal, log, proxyOptions = null }) {
if (isAgentTextRequest(body)) {
if (isAgentCapableRequest(body)) {
try {
return await this.executeAgent({ model, body, stream, credentials, signal });
return await this.executeAgent({ model, body, stream, credentials, signal, log });
} catch (error) {
return {
response: new Response(JSON.stringify({