Merge pull request #525 from kaitranntt/kai/feat/518-cursor-protobuf-executor

feat(cursor): add ConnectRPC protobuf encoder/decoder and HTTP/2 executor
This commit is contained in:
Kai (Tam Nhu) Tran authored and GitHub committed 2026-02-12 03:48:23 +07:00
commit 5df2965642
7 files changed
+2559

No files matched your search

+788
View File
@@ -0,0 +1,788 @@
/**
* Cursor Executor
* Handles HTTP/2 requests to Cursor API with protobuf encoding/decoding
*/
import * as crypto from 'crypto';
import * as zlib from 'zlib';
import type { IncomingHttpHeaders } from 'http';
import { generateCursorBody, extractTextFromResponse } from './cursor-protobuf.js';
import { buildCursorRequest } from './cursor-translator.js';
import type { CursorTool, CursorCredentials } from './cursor-protobuf-schema.js';
import { COMPRESS_FLAG } from './cursor-protobuf-schema.js';
/** Executor parameters */
interface ExecutorParams {
model: string;
body: {
messages: Array<{
role: string;
content: string | Array<{ type: string; text?: string }>;
name?: string;
tool_call_id?: string;
tool_calls?: Array<{
id: string;
type: string;
function: { name: string; arguments: string };
}>;
}>;
tools?: CursorTool[];
reasoning_effort?: string;
};
stream: boolean;
credentials: CursorCredentials;
signal?: AbortSignal;
}
/** HTTP/2 response structure */
interface Http2Response {
status: number;
headers: IncomingHttpHeaders;
body: Buffer;
}
/** Lazy import http2 */
let http2Module: typeof import('http2') | null = null;
async function getHttp2() {
if (http2Module) return http2Module;
try {
http2Module = await import('http2');
return http2Module;
} catch (err) {
if (process.env.CCS_DEBUG) {
console.error('[cursor] http2 module not available, falling back to fetch:', err);
}
return null;
}
}
/**
* Decompress payload if needed
* NOTE: Uses synchronous gzip for single-request CLI tool. Async not warranted for small payloads.
*/
function decompressPayload(payload: Buffer, flags: number): Buffer {
// Check if payload is JSON error
if (payload.length > 10 && payload[0] === 0x7b && payload[1] === 0x22) {
try {
const text = payload.toString('utf-8');
if (text.startsWith('{"error"')) {
return payload;
}
} catch (err) {
if (process.env.CCS_DEBUG) {
console.error('[cursor] JSON error detection failed:', err);
}
}
}
if (
flags === COMPRESS_FLAG.GZIP ||
flags === COMPRESS_FLAG.GZIP_ALT ||
flags === COMPRESS_FLAG.GZIP_BOTH
) {
try {
return zlib.gunzipSync(payload);
} catch (err) {
if (process.env.CCS_DEBUG) {
console.error('[cursor] gzip decompression failed:', err);
}
return Buffer.alloc(0);
}
}
return payload;
}
/**
* Create error response from JSON error
*/
function createErrorResponse(jsonError: {
error?: {
code?: string;
message?: string;
details?: Array<{ debug?: { details?: { title?: string; detail?: string }; error?: string } }>;
};
}): Response {
const errorMsg =
jsonError?.error?.details?.[0]?.debug?.details?.title ||
jsonError?.error?.details?.[0]?.debug?.details?.detail ||
jsonError?.error?.message ||
'API Error';
const isRateLimit = jsonError?.error?.code === 'resource_exhausted';
return new Response(
JSON.stringify({
error: {
message: errorMsg,
type: isRateLimit ? 'rate_limit_error' : 'api_error',
code: jsonError?.error?.details?.[0]?.debug?.error || 'unknown',
},
}),
{
status: isRateLimit ? 429 : 400,
headers: { 'Content-Type': 'application/json' },
}
);
}
export class CursorExecutor {
private readonly baseUrl = 'https://api2.cursor.sh';
private readonly chatPath = '/aiserver.v1.AiService/StreamChat';
private readonly CURSOR_CLIENT_VERSION = '2.3.41';
private readonly CURSOR_USER_AGENT = 'connect-es/1.6.1';
buildUrl(): string {
return `${this.baseUrl}${this.chatPath}`;
}
/**
* Generate checksum using Jyh cipher (time-based XOR with rolling key seed=165)
*/
generateChecksum(machineId: string): string {
const timestamp = Math.floor(Date.now() / 1000000);
// JS bitwise shifts wrap modulo 32, so >>40 and >>32 give wrong results.
// Use Math.trunc division for upper bytes that exceed 32-bit range.
const byteArray = new Uint8Array([
Math.trunc(timestamp / 2 ** 40) & 0xff,
Math.trunc(timestamp / 2 ** 32) & 0xff,
(timestamp >>> 24) & 0xff,
(timestamp >>> 16) & 0xff,
(timestamp >>> 8) & 0xff,
timestamp & 0xff,
]);
let t = 165;
for (let i = 0; i < byteArray.length; i++) {
byteArray[i] = ((byteArray[i] ^ t) + (i % 256)) & 0xff;
t = byteArray[i];
}
const alphabet = 'ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_';
let encoded = '';
for (let i = 0; i < byteArray.length; i += 3) {
const a = byteArray[i];
const b = i + 1 < byteArray.length ? byteArray[i + 1] : 0;
const c = i + 2 < byteArray.length ? byteArray[i + 2] : 0;
encoded += alphabet[a >> 2];
encoded += alphabet[((a & 3) << 4) | (b >> 4)];
if (i + 1 < byteArray.length) {
encoded += alphabet[((b & 15) << 2) | (c >> 6)];
}
if (i + 2 < byteArray.length) {
encoded += alphabet[c & 63];
}
}
return `${encoded}${machineId}`;
}
buildHeaders(credentials: CursorCredentials): Record<string, string> {
const accessToken = credentials.accessToken;
const machineId = credentials.machineId;
const ghostMode = credentials.ghostMode !== false;
if (!machineId) {
throw new Error('Machine ID is required for Cursor API');
}
const delimIdx = accessToken.indexOf('::');
const cleanToken = delimIdx !== -1 ? accessToken.slice(delimIdx + 2) : accessToken;
if (!cleanToken) {
throw new Error('Access token is empty after parsing');
}
return {
authorization: `Bearer ${cleanToken}`,
'connect-accept-encoding': 'gzip',
'connect-protocol-version': '1',
'content-type': 'application/connect+proto',
'user-agent': this.CURSOR_USER_AGENT,
'x-amzn-trace-id': `Root=${crypto.randomUUID()}`,
'x-client-key': crypto.createHash('sha256').update(cleanToken).digest('hex'),
'x-cursor-checksum': this.generateChecksum(machineId),
'x-cursor-client-version': this.CURSOR_CLIENT_VERSION,
'x-cursor-client-type': 'ide',
'x-cursor-client-os':
process.platform === 'win32'
? 'windows'
: process.platform === 'darwin'
? 'macos'
: 'linux',
'x-cursor-client-arch': process.arch === 'arm64' ? 'aarch64' : 'x64',
'x-cursor-client-device-type': 'desktop',
'x-cursor-config-version': crypto.randomUUID(),
'x-cursor-timezone': Intl.DateTimeFormat().resolvedOptions().timeZone || 'UTC',
'x-ghost-mode': ghostMode ? 'true' : 'false',
'x-request-id': crypto.randomUUID(),
'x-session-id': crypto.createHash('sha256').update(cleanToken).digest('hex').substring(0, 36),
};
}
transformRequest(
model: string,
body: ExecutorParams['body'],
stream: boolean,
credentials: CursorCredentials
): Uint8Array {
const translatedBody = buildCursorRequest(model, body, stream, credentials);
const messages = translatedBody.messages || [];
const tools = (translatedBody.tools || body.tools || []) as CursorTool[];
const reasoningEffort = body.reasoning_effort || null;
return generateCursorBody(messages, model, tools, reasoningEffort);
}
async makeFetchRequest(
url: string,
headers: Record<string, string>,
body: Uint8Array,
signal?: AbortSignal
): Promise<Http2Response> {
const response = await fetch(url, {
method: 'POST',
headers,
body,
signal,
});
const responseHeaders: Record<string, string> = {};
response.headers.forEach((value, key) => {
responseHeaders[key] = value;
});
return {
status: response.status,
headers: responseHeaders,
body: Buffer.from(await response.arrayBuffer()),
};
}
async makeHttp2Request(
url: string,
headers: Record<string, string>,
body: Uint8Array,
signal?: AbortSignal
): Promise<Http2Response> {
const http2 = await getHttp2();
if (!http2) {
throw new Error('http2 module not available');
}
return new Promise((resolve, reject) => {
const urlObj = new URL(url);
const client = http2.connect(`https://${urlObj.host}`);
const chunks: Buffer[] = [];
let responseHeaders: IncomingHttpHeaders = {};
client.on('error', (err) => {
client.close();
reject(err);
});
const req = client.request({
':method': 'POST',
':path': urlObj.pathname,
':authority': urlObj.host,
':scheme': 'https',
...headers,
});
req.on('response', (hdrs) => {
responseHeaders = hdrs;
});
req.on('data', (chunk: Buffer) => {
chunks.push(chunk);
});
req.on('end', () => {
client.close();
resolve({
status: Number(responseHeaders[':status']) || 500,
headers: responseHeaders,
body: Buffer.concat(chunks),
});
});
req.on('error', (err) => {
client.close();
reject(err);
});
if (signal) {
const onAbort = () => {
req.close();
client.close();
reject(new Error('Request aborted'));
};
signal.addEventListener('abort', onAbort, { once: true });
const cleanup = () => {
signal.removeEventListener('abort', onAbort);
};
req.on('end', cleanup);
req.on('error', cleanup);
}
req.write(body);
req.end();
});
}
async execute(params: ExecutorParams): Promise<{
response: Response;
url: string;
headers: Record<string, string>;
transformedBody: ExecutorParams['body'];
}> {
const { model, body, stream, credentials, signal } = params;
const url = this.buildUrl();
const headers = this.buildHeaders(credentials);
const transformedBody = this.transformRequest(model, body, stream, credentials);
try {
const http2 = await getHttp2();
const response = http2
? await this.makeHttp2Request(url, headers, transformedBody, signal)
: await this.makeFetchRequest(url, headers, transformedBody, signal);
if (response.status !== 200) {
const errorText = response.body?.toString() || 'Unknown error';
const errorResponse = new Response(
JSON.stringify({
error: {
message: `[${response.status}]: ${errorText}`,
type: 'invalid_request_error',
code: '',
},
}),
{
status: response.status,
headers: { 'Content-Type': 'application/json' },
}
);
return { response: errorResponse, url, headers, transformedBody: body };
}
const transformedResponse =
stream === true
? this.transformProtobufToSSE(response.body, model, body)
: this.transformProtobufToJSON(response.body, model, body);
return { response: transformedResponse, url, headers, transformedBody: body };
} catch (error) {
const errorResponse = new Response(
JSON.stringify({
error: {
message: (error as Error).message,
type: 'connection_error',
code: '',
},
}),
{
status: 500,
headers: { 'Content-Type': 'application/json' },
}
);
return { response: errorResponse, url, headers, transformedBody: body };
}
}
/**
* Parse protobuf buffer into frames and extract text/toolcalls.
* Shared logic between JSON and SSE transformers.
*/
private *parseProtobufFrames(buffer: Buffer): Generator<
| { type: 'error'; response: Response }
| { type: 'text'; text: string }
| {
type: 'toolCall';
toolCall: {
id: string;
type: string;
function: { name: string; arguments: string };
isLast: boolean;
};
}
> {
let offset = 0;
while (offset < buffer.length) {
if (offset + 5 > buffer.length) break;
const flags = buffer[offset];
const length = buffer.readUInt32BE(offset + 1);
if (offset + 5 + length > buffer.length) break;
let payload = buffer.slice(offset + 5, offset + 5 + length);
offset += 5 + length;
payload = decompressPayload(payload, flags);
// Check for JSON error format
try {
const text = payload.toString('utf-8');
if (text.startsWith('{') && text.includes('"error"')) {
yield { type: 'error', response: createErrorResponse(JSON.parse(text)) };
return;
}
} catch (err) {
if (process.env.CCS_DEBUG) {
console.error('[cursor] parseProtobufFrames error parsing failed:', err);
}
}
const result = extractTextFromResponse(new Uint8Array(payload));
// Check for protobuf-decoded error
if (result.error) {
const errorLower = result.error.toLowerCase();
const isRateLimit =
errorLower.includes('rate limit') ||
errorLower.includes('resource_exhausted') ||
errorLower.includes('too many requests');
yield {
type: 'error',
response: new Response(
JSON.stringify({
error: {
message: result.error,
type: isRateLimit ? 'rate_limit_error' : 'server_error',
code: isRateLimit ? 'rate_limited' : 'cursor_error',
},
}),
{
status: isRateLimit ? 429 : 400,
headers: { 'Content-Type': 'application/json' },
}
),
};
return;
}
if (result.toolCall) {
yield { type: 'toolCall', toolCall: result.toolCall };
}
if (result.text) {
yield { type: 'text', text: result.text };
}
}
}
transformProtobufToJSON(buffer: Buffer, model: string, _body: ExecutorParams['body']): Response {
const responseId = `chatcmpl-cursor-${Date.now()}`;
const created = Math.floor(Date.now() / 1000);
let totalContent = '';
const toolCalls: Array<{
id: string;
type: string;
function: { name: string; arguments: string };
}> = [];
const toolCallsMap = new Map<
string,
{
id: string;
type: string;
function: { name: string; arguments: string };
isLast: boolean;
index: number;
}
>();
for (const frame of this.parseProtobufFrames(buffer)) {
if (frame.type === 'error') {
return frame.response;
}
if (frame.type === 'toolCall') {
const tc = frame.toolCall;
if (toolCallsMap.has(tc.id)) {
const existing = toolCallsMap.get(tc.id);
if (!existing) continue;
existing.function.arguments += tc.function.arguments;
existing.isLast = tc.isLast;
} else {
toolCallsMap.set(tc.id, {
...tc,
index: toolCallsMap.size,
});
}
if (tc.isLast) {
const finalToolCall = toolCallsMap.get(tc.id);
if (!finalToolCall) continue;
toolCalls.push({
id: finalToolCall.id,
type: finalToolCall.type,
function: {
name: finalToolCall.function.name,
arguments: finalToolCall.function.arguments,
},
});
}
}
if (frame.type === 'text') {
totalContent += frame.text;
}
}
// Finalize remaining tool calls
for (const id of Array.from(toolCallsMap.keys())) {
const tc = toolCallsMap.get(id);
if (!tc) continue;
if (!toolCalls.find((t) => t.id === id)) {
toolCalls.push({
id: tc.id,
type: tc.type,
function: {
name: tc.function.name,
arguments: tc.function.arguments,
},
});
}
}
const message: {
role: string;
content: string | null;
tool_calls?: Array<{
id: string;
type: string;
function: { name: string; arguments: string };
}>;
} = {
role: 'assistant',
content: totalContent || null,
};
if (toolCalls.length > 0) {
message.tool_calls = toolCalls;
}
const completion = {
id: responseId,
object: 'chat.completion',
created,
model,
choices: [
{
index: 0,
message,
finish_reason: toolCalls.length > 0 ? 'tool_calls' : 'stop',
},
],
usage: {
prompt_tokens: 0,
completion_tokens: 0,
total_tokens: 0,
},
};
return new Response(JSON.stringify(completion), {
status: 200,
headers: { 'Content-Type': 'application/json' },
});
}
transformProtobufToSSE(buffer: Buffer, model: string, _body: ExecutorParams['body']): Response {
// TODO(#531): Implement true streaming — currently buffers entire response before transforming.
// This should pipe HTTP/2 data events through a TransformStream for incremental SSE output.
// See: https://github.com/kaitranntt/ccs/issues/531
// NOTE: Chunk boundary splits may emit duplicate SSE messages if a frame spans multiple chunks.
const responseId = `chatcmpl-cursor-${Date.now()}`;
const created = Math.floor(Date.now() / 1000);
const chunks: string[] = [];
const toolCalls: Array<{
id: string;
type: string;
function: { name: string; arguments: string };
index: number;
}> = [];
const toolCallsMap = new Map<
string,
{
id: string;
type: string;
function: { name: string; arguments: string };
isLast: boolean;
index: number;
}
>();
for (const frame of this.parseProtobufFrames(buffer)) {
if (frame.type === 'error') {
return frame.response;
}
if (frame.type === 'toolCall') {
const tc = frame.toolCall;
if (chunks.length === 0) {
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: 'chat.completion.chunk',
created,
model,
choices: [
{
index: 0,
delta: { role: 'assistant', content: '' },
finish_reason: null,
},
],
})}\n\n`
);
}
if (toolCallsMap.has(tc.id)) {
const existing = toolCallsMap.get(tc.id);
if (!existing) continue;
existing.function.arguments += tc.function.arguments;
existing.isLast = tc.isLast;
if (tc.function.arguments) {
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: 'chat.completion.chunk',
created,
model,
choices: [
{
index: 0,
delta: {
tool_calls: [
{
index: existing.index,
id: tc.id,
type: 'function',
function: {
name: tc.function.name,
arguments: tc.function.arguments,
},
},
],
},
finish_reason: null,
},
],
})}\n\n`
);
}
} else {
const toolCallIndex = toolCalls.length;
toolCalls.push({ ...tc, index: toolCallIndex });
toolCallsMap.set(tc.id, { ...tc, index: toolCallIndex });
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: 'chat.completion.chunk',
created,
model,
choices: [
{
index: 0,
delta: {
tool_calls: [
{
index: toolCallIndex,
id: tc.id,
type: 'function',
function: {
name: tc.function.name,
arguments: tc.function.arguments,
},
},
],
},
finish_reason: null,
},
],
})}\n\n`
);
}
}
if (frame.type === 'text') {
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: 'chat.completion.chunk',
created,
model,
choices: [
{
index: 0,
delta:
chunks.length === 0 && toolCalls.length === 0
? { role: 'assistant', content: frame.text }
: { content: frame.text },
finish_reason: null,
},
],
})}\n\n`
);
}
}
if (chunks.length === 0 && toolCalls.length === 0) {
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: 'chat.completion.chunk',
created,
model,
choices: [
{
index: 0,
delta: { role: 'assistant', content: '' },
finish_reason: null,
},
],
})}\n\n`
);
}
chunks.push(
`data: ${JSON.stringify({
id: responseId,
object: 'chat.completion.chunk',
created,
model,
choices: [
{
index: 0,
delta: {},
finish_reason: toolCalls.length > 0 ? 'tool_calls' : 'stop',
},
],
usage: {
prompt_tokens: 0,
completion_tokens: 0,
total_tokens: 0,
},
})}\n\n`
);
chunks.push('data: [DONE]\n\n');
return new Response(chunks.join(''), {
status: 200,
headers: {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
Connection: 'keep-alive',
},
});
}
}
export default CursorExecutor;
+336
View File
@@ -0,0 +1,336 @@
/**
* Cursor Protobuf Decoder
* Implements ConnectRPC protobuf wire format decoding
*/
import * as zlib from 'zlib';
import { WIRE_TYPE, FIELD, COMPRESS_FLAG, type WireType } from './cursor-protobuf-schema.js';
/**
* Decode a varint from buffer
* Returns [value, newOffset]
*/
export function decodeVarint(buffer: Uint8Array, offset: number): [number, number] {
let result = 0;
let shift = 0;
let pos = offset;
const maxBytes = 5;
while (pos < buffer.length && pos - offset < maxBytes) {
const b = buffer[pos];
result |= (b & 0x7f) << shift;
pos++;
if (!(b & 0x80)) break;
shift += 7;
}
return [result >>> 0, pos]; // Ensure unsigned
}
/**
* Decode a single protobuf field
* Returns [fieldNum, wireType, value, newOffset]
*/
export function decodeField(
buffer: Uint8Array,
offset: number
): [number | null, WireType | null, Uint8Array | number | null, number] {
if (offset >= buffer.length) {
return [null, null, null, offset];
}
const [tag, pos1] = decodeVarint(buffer, offset);
const fieldNum = tag >> 3;
const wireType = (tag & 0x07) as WireType;
let value: Uint8Array | number | null;
let pos = pos1;
if (wireType === WIRE_TYPE.VARINT) {
[value, pos] = decodeVarint(buffer, pos);
} else if (wireType === WIRE_TYPE.LEN) {
const [length, pos2] = decodeVarint(buffer, pos);
if (pos2 + length > buffer.length) {
return [null, null, null, buffer.length];
}
value = buffer.slice(pos2, pos2 + length);
pos = pos2 + length;
} else if (wireType === WIRE_TYPE.FIXED64) {
if (pos + 8 > buffer.length) {
return [null, null, null, buffer.length];
}
value = buffer.slice(pos, pos + 8);
pos += 8;
} else if (wireType === WIRE_TYPE.FIXED32) {
if (pos + 4 > buffer.length) {
return [null, null, null, buffer.length];
}
value = buffer.slice(pos, pos + 4);
pos += 4;
} else {
value = null;
}
return [fieldNum, wireType, value, pos];
}
/**
* Decode a protobuf message into a map of fields
*/
export function decodeMessage(
data: Uint8Array
): Map<number, Array<{ wireType: WireType; value: Uint8Array | number }>> {
const fields = new Map<number, Array<{ wireType: WireType; value: Uint8Array | number }>>();
let pos = 0;
// NOTE: If two fields share the same field number but different wire types, later values overwrite earlier ones.
while (pos < data.length) {
const [fieldNum, wireType, value, newPos] = decodeField(data, pos);
if (fieldNum === null || wireType === null || value === null) break;
if (!fields.has(fieldNum)) {
fields.set(fieldNum, []);
}
const fieldArray = fields.get(fieldNum);
if (fieldArray) {
fieldArray.push({ wireType, value: value as Uint8Array | number });
}
pos = newPos;
}
return fields;
}
/**
* Parse ConnectRPC frame from buffer
* Returns frame data or null if incomplete
*/
export function parseConnectRPCFrame(buffer: Buffer): {
flags: number;
length: number;
payload: Uint8Array;
consumed: number;
} | null {
if (buffer.length < 5) return null;
const flags = buffer[0];
const length = (buffer[1] << 24) | (buffer[2] << 16) | (buffer[3] << 8) | buffer[4];
if (buffer.length < 5 + length) return null;
let payload = buffer.slice(5, 5 + length);
// Decompress if gzip
if (
flags === COMPRESS_FLAG.GZIP ||
flags === COMPRESS_FLAG.GZIP_ALT ||
flags === COMPRESS_FLAG.GZIP_BOTH
) {
try {
payload = Buffer.from(zlib.gunzipSync(payload));
} catch (err) {
if (process.env.CCS_DEBUG) {
console.error('[cursor] parseConnectRPCFrame decompression failed:', err);
}
// Decompression failed, use raw payload
}
}
return {
flags,
length,
payload: new Uint8Array(payload),
consumed: 5 + length,
};
}
/**
* Extract tool call from protobuf data
*/
function extractToolCall(toolCallData: Uint8Array): {
id: string;
type: string;
function: { name: string; arguments: string };
isLast: boolean;
} | null {
const toolCall = decodeMessage(toolCallData);
let toolCallId = '';
let toolName = '';
let rawArgs = '';
let isLast = false;
// Extract tool call ID
if (toolCall.has(FIELD.TOOL_ID)) {
const idField = toolCall.get(FIELD.TOOL_ID);
if (idField && idField[0]) {
const fullId = new TextDecoder().decode(idField[0].value as Uint8Array);
toolCallId = fullId.split('\n')[0]; // Take first line
}
}
// Extract tool name
if (toolCall.has(FIELD.TOOL_NAME)) {
const nameField = toolCall.get(FIELD.TOOL_NAME);
if (nameField && nameField[0]) {
toolName = new TextDecoder().decode(nameField[0].value as Uint8Array);
}
}
// Extract is_last flag
if (toolCall.has(FIELD.TOOL_IS_LAST)) {
const lastField = toolCall.get(FIELD.TOOL_IS_LAST);
if (lastField && lastField[0]) {
isLast = (lastField[0].value as number) !== 0;
}
}
// Extract MCP params - nested real tool info
if (toolCall.has(FIELD.TOOL_MCP_PARAMS)) {
try {
const mcpField = toolCall.get(FIELD.TOOL_MCP_PARAMS);
if (!mcpField || !mcpField[0]) return null;
const mcpParams = decodeMessage(mcpField[0].value as Uint8Array);
if (mcpParams.has(FIELD.MCP_TOOLS_LIST)) {
const toolsList = mcpParams.get(FIELD.MCP_TOOLS_LIST);
if (!toolsList || !toolsList[0]) return null;
const tool = decodeMessage(toolsList[0].value as Uint8Array);
if (tool.has(FIELD.MCP_NESTED_NAME)) {
const nestedName = tool.get(FIELD.MCP_NESTED_NAME);
if (nestedName && nestedName[0]) {
toolName = new TextDecoder().decode(nestedName[0].value as Uint8Array);
}
}
if (tool.has(FIELD.MCP_NESTED_PARAMS)) {
const nestedParams = tool.get(FIELD.MCP_NESTED_PARAMS);
if (nestedParams && nestedParams[0]) {
rawArgs = new TextDecoder().decode(nestedParams[0].value as Uint8Array);
}
}
}
} catch (err) {
if (process.env.CCS_DEBUG) {
console.error('[cursor] extractToolCall MCP parsing failed:', err);
}
// MCP parse error, continue
}
}
// Fallback to raw_args
if (!rawArgs && toolCall.has(FIELD.TOOL_RAW_ARGS)) {
const rawArgsField = toolCall.get(FIELD.TOOL_RAW_ARGS);
if (rawArgsField && rawArgsField[0]) {
rawArgs = new TextDecoder().decode(rawArgsField[0].value as Uint8Array);
}
}
if (toolCallId && toolName) {
return {
id: toolCallId,
type: 'function',
function: {
name: toolName,
arguments: rawArgs || '{}',
},
isLast,
};
}
return null;
}
/**
* Extract text and thinking from response data
*/
function extractTextAndThinking(responseData: Uint8Array): {
text: string | null;
thinking: string | null;
} {
const nested = decodeMessage(responseData);
let text: string | null = null;
let thinking: string | null = null;
// Extract text
if (nested.has(FIELD.RESPONSE_TEXT)) {
const textField = nested.get(FIELD.RESPONSE_TEXT);
if (textField && textField[0]) {
text = new TextDecoder().decode(textField[0].value as Uint8Array);
}
}
// Extract thinking
if (nested.has(FIELD.THINKING)) {
try {
const thinkingField = nested.get(FIELD.THINKING);
if (thinkingField && thinkingField[0]) {
const thinkingMsg = decodeMessage(thinkingField[0].value as Uint8Array);
if (thinkingMsg.has(FIELD.THINKING_TEXT)) {
const thinkingTextField = thinkingMsg.get(FIELD.THINKING_TEXT);
if (thinkingTextField && thinkingTextField[0]) {
thinking = new TextDecoder().decode(thinkingTextField[0].value as Uint8Array);
}
}
}
} catch (err) {
if (process.env.CCS_DEBUG) {
console.error('[cursor] extractTextAndThinking parsing failed:', err);
}
// Thinking parse error, continue
}
}
return { text, thinking };
}
/**
* Extract text and tool calls from response payload
*/
export function extractTextFromResponse(payload: Uint8Array): {
text: string | null;
error: string | null;
toolCall: {
id: string;
type: string;
function: { name: string; arguments: string };
isLast: boolean;
} | null;
thinking: string | null;
} {
try {
const fields = decodeMessage(payload);
// Field 1: ClientSideToolV2Call
if (fields.has(FIELD.TOOL_CALL)) {
const toolCallField = fields.get(FIELD.TOOL_CALL);
if (toolCallField && toolCallField[0]) {
const toolCall = extractToolCall(toolCallField[0].value as Uint8Array);
if (toolCall) {
return { text: null, error: null, toolCall, thinking: null };
}
}
}
// Field 2: StreamUnifiedChatResponse
if (fields.has(FIELD.RESPONSE)) {
const responseField = fields.get(FIELD.RESPONSE);
if (responseField && responseField[0]) {
const { text, thinking } = extractTextAndThinking(responseField[0].value as Uint8Array);
if (text || thinking) {
return { text, error: null, toolCall: null, thinking };
}
}
}
return { text: null, error: null, toolCall: null, thinking: null };
} catch (err) {
if (process.env.CCS_DEBUG) {
console.error('[cursor] extractTextFromResponse parsing failed:', err);
}
return { text: null, error: null, toolCall: null, thinking: null };
}
}
+224
View File
@@ -0,0 +1,224 @@
/**
* Cursor Protobuf Encoder
* Implements ConnectRPC protobuf wire format encoding
*/
import * as zlib from 'zlib';
import {
WIRE_TYPE,
FIELD,
COMPRESS_FLAG,
UNIFIED_MODE,
type WireType,
type RoleType,
type CursorTool,
type CursorToolResult,
} from './cursor-protobuf-schema.js';
/**
* Encode a varint (variable-length integer)
*/
export function encodeVarint(value: number): Uint8Array {
const bytes: number[] = [];
let val = value >>> 0; // Ensure unsigned
while (val >= 0x80) {
bytes.push((val & 0x7f) | 0x80);
val >>>= 7;
}
bytes.push(val & 0x7f);
return new Uint8Array(bytes);
}
/**
* Encode a protobuf field (tag + value)
*/
export function encodeField(
fieldNum: number,
wireType: WireType,
value: number | string | Uint8Array
): Uint8Array {
const tag = (fieldNum << 3) | wireType;
const tagBytes = encodeVarint(tag);
if (wireType === WIRE_TYPE.VARINT) {
const valueBytes = encodeVarint(value as number);
return concatArrays(tagBytes, valueBytes);
}
if (wireType === WIRE_TYPE.LEN) {
const dataBytes =
typeof value === 'string'
? new TextEncoder().encode(value)
: value instanceof Uint8Array
? value
: new Uint8Array(0);
const lengthBytes = encodeVarint(dataBytes.length);
return concatArrays(tagBytes, lengthBytes, dataBytes);
}
return new Uint8Array(0);
}
/**
* Concatenate multiple Uint8Arrays
*/
export function concatArrays(...arrays: Uint8Array[]): Uint8Array {
const totalLength = arrays.reduce((sum, arr) => sum + arr.length, 0);
const result = new Uint8Array(totalLength);
let offset = 0;
for (const arr of arrays) {
result.set(arr, offset);
offset += arr.length;
}
return result;
}
/**
* Encode a tool result
*/
export function encodeToolResult(toolResult: CursorToolResult): Uint8Array {
const toolCallId = toolResult.tool_call_id || '';
const toolName = toolResult.name || '';
const toolIndex = toolResult.index || 0;
const rawArgs = toolResult.raw_args || '{}';
return concatArrays(
encodeField(FIELD.TOOL_RESULT_CALL_ID, WIRE_TYPE.LEN, toolCallId),
encodeField(FIELD.TOOL_RESULT_NAME, WIRE_TYPE.LEN, toolName),
encodeField(FIELD.TOOL_RESULT_INDEX, WIRE_TYPE.VARINT, toolIndex),
encodeField(FIELD.TOOL_RESULT_RAW_ARGS, WIRE_TYPE.LEN, rawArgs)
);
}
/**
* Encode a conversation message
*/
export function encodeMessage(
content: string,
role: RoleType,
messageId: string,
isLast: boolean,
hasTools: boolean,
toolResults: CursorToolResult[]
): Uint8Array {
return concatArrays(
encodeField(FIELD.MSG_CONTENT, WIRE_TYPE.LEN, content),
encodeField(FIELD.MSG_ROLE, WIRE_TYPE.VARINT, role),
encodeField(FIELD.MSG_ID, WIRE_TYPE.LEN, messageId),
...(toolResults.length > 0
? toolResults.map((tr) =>
encodeField(FIELD.MSG_TOOL_RESULTS, WIRE_TYPE.LEN, encodeToolResult(tr))
)
: []),
encodeField(FIELD.MSG_IS_AGENTIC, WIRE_TYPE.VARINT, hasTools ? 1 : 0),
encodeField(
FIELD.MSG_UNIFIED_MODE,
WIRE_TYPE.VARINT,
hasTools ? UNIFIED_MODE.AGENT : UNIFIED_MODE.CHAT
),
...(isLast && hasTools
? [encodeField(FIELD.MSG_SUPPORTED_TOOLS, WIRE_TYPE.LEN, encodeVarint(1))]
: [])
);
}
/**
* Encode instruction text
*/
export function encodeInstruction(text: string): Uint8Array {
return text ? encodeField(FIELD.INSTRUCTION_TEXT, WIRE_TYPE.LEN, text) : new Uint8Array(0);
}
/**
* Encode model information
*/
export function encodeModel(modelName: string): Uint8Array {
return concatArrays(
encodeField(FIELD.MODEL_NAME, WIRE_TYPE.LEN, modelName),
encodeField(FIELD.MODEL_EMPTY, WIRE_TYPE.LEN, new Uint8Array(0))
);
}
/**
* Encode cursor settings
*/
export function encodeCursorSetting(): Uint8Array {
const unknown6 = concatArrays(
encodeField(FIELD.SETTING6_FIELD_1, WIRE_TYPE.LEN, new Uint8Array(0)),
encodeField(FIELD.SETTING6_FIELD_2, WIRE_TYPE.LEN, new Uint8Array(0))
);
return concatArrays(
encodeField(FIELD.SETTING_PATH, WIRE_TYPE.LEN, 'cursor\\aisettings'),
encodeField(FIELD.SETTING_UNKNOWN_3, WIRE_TYPE.LEN, new Uint8Array(0)),
encodeField(FIELD.SETTING_UNKNOWN_6, WIRE_TYPE.LEN, unknown6),
encodeField(FIELD.SETTING_UNKNOWN_8, WIRE_TYPE.VARINT, 1),
encodeField(FIELD.SETTING_UNKNOWN_9, WIRE_TYPE.VARINT, 1)
);
}
/**
* Encode metadata
*/
export function encodeMetadata(): Uint8Array {
return concatArrays(
encodeField(FIELD.META_PLATFORM, WIRE_TYPE.LEN, process.platform || 'linux'),
encodeField(FIELD.META_ARCH, WIRE_TYPE.LEN, process.arch || 'x64'),
encodeField(FIELD.META_VERSION, WIRE_TYPE.LEN, process.version || 'v20.0.0'),
encodeField(FIELD.META_CWD, WIRE_TYPE.LEN, process.cwd() || '/'),
encodeField(FIELD.META_TIMESTAMP, WIRE_TYPE.LEN, new Date().toISOString())
);
}
/**
* Encode message ID
*/
export function encodeMessageId(messageId: string, role: RoleType, summaryId?: string): Uint8Array {
return concatArrays(
encodeField(FIELD.MSGID_ID, WIRE_TYPE.LEN, messageId),
...(summaryId ? [encodeField(FIELD.MSGID_SUMMARY, WIRE_TYPE.LEN, summaryId)] : []),
encodeField(FIELD.MSGID_ROLE, WIRE_TYPE.VARINT, role)
);
}
/**
* Encode MCP tool
*/
export function encodeMcpTool(tool: CursorTool): Uint8Array {
const toolName = tool.function?.name || tool.name || '';
const toolDesc = tool.function?.description || tool.description || '';
const inputSchema = tool.function?.parameters || tool.input_schema || {};
return concatArrays(
...(toolName ? [encodeField(FIELD.MCP_TOOL_NAME, WIRE_TYPE.LEN, toolName)] : []),
...(toolDesc ? [encodeField(FIELD.MCP_TOOL_DESC, WIRE_TYPE.LEN, toolDesc)] : []),
...(Object.keys(inputSchema).length > 0
? [encodeField(FIELD.MCP_TOOL_PARAMS, WIRE_TYPE.LEN, JSON.stringify(inputSchema))]
: []),
encodeField(FIELD.MCP_TOOL_SERVER, WIRE_TYPE.LEN, 'custom')
);
}
/**
* Wrap payload in ConnectRPC frame (5-byte header + payload)
*/
export function wrapConnectRPCFrame(payload: Uint8Array, compress = false): Uint8Array {
let finalPayload = payload;
let flags: number = COMPRESS_FLAG.NONE;
if (compress) {
finalPayload = new Uint8Array(zlib.gzipSync(Buffer.from(payload)));
flags = COMPRESS_FLAG.GZIP;
}
const frame = new Uint8Array(5 + finalPayload.length);
frame[0] = flags;
frame[1] = (finalPayload.length >> 24) & 0xff;
frame[2] = (finalPayload.length >> 16) & 0xff;
frame[3] = (finalPayload.length >> 8) & 0xff;
frame[4] = finalPayload.length & 0xff;
frame.set(finalPayload, 5);
return frame;
}
+213
View File
@@ -0,0 +1,213 @@
/**
* Cursor Protobuf Schema Constants
* Field definitions and wire types for ConnectRPC protocol
*/
/** Wire types for protobuf encoding */
export const WIRE_TYPE = {
VARINT: 0,
FIXED64: 1,
LEN: 2,
FIXED32: 5,
} as const;
/** Message role constants */
export const ROLE = {
USER: 1,
ASSISTANT: 2,
} as const;
/** Unified mode constants */
export const UNIFIED_MODE = {
CHAT: 1,
AGENT: 2,
} as const;
/** Thinking level constants */
export const THINKING_LEVEL = {
UNSPECIFIED: 0,
MEDIUM: 1,
HIGH: 2,
} as const;
/** Field numbers for all protobuf messages */
export const FIELD = {
// ===== StreamUnifiedChatRequestWithTools (top level) =====
REQUEST: 1,
// ===== StreamUnifiedChatRequest =====
MESSAGES: 1,
UNKNOWN_2: 2,
INSTRUCTION: 3,
UNKNOWN_4: 4,
MODEL: 5,
WEB_TOOL: 8,
UNKNOWN_13: 13,
CURSOR_SETTING: 15,
UNKNOWN_19: 19,
CONVERSATION_ID: 23,
METADATA: 26,
IS_AGENTIC: 27,
SUPPORTED_TOOLS: 29,
MESSAGE_IDS: 30,
MCP_TOOLS: 34,
LARGE_CONTEXT: 35,
UNKNOWN_38: 38,
UNIFIED_MODE: 46,
UNKNOWN_47: 47,
SHOULD_DISABLE_TOOLS: 48,
THINKING_LEVEL: 49,
UNKNOWN_51: 51,
UNKNOWN_53: 53,
UNIFIED_MODE_NAME: 54,
// ===== ConversationMessage =====
MSG_CONTENT: 1,
MSG_ROLE: 2,
MSG_ID: 13,
MSG_TOOL_RESULTS: 18,
MSG_IS_AGENTIC: 29,
MSG_UNIFIED_MODE: 47,
MSG_SUPPORTED_TOOLS: 51,
// ===== ConversationMessage.ToolResult =====
TOOL_RESULT_CALL_ID: 1,
TOOL_RESULT_NAME: 2,
TOOL_RESULT_INDEX: 3,
TOOL_RESULT_RAW_ARGS: 5,
TOOL_RESULT_RESULT: 8, // Reserved for future tool result parsing
// ===== Model =====
MODEL_NAME: 1,
MODEL_EMPTY: 4,
// ===== Instruction =====
INSTRUCTION_TEXT: 1,
// ===== CursorSetting =====
SETTING_PATH: 1,
SETTING_UNKNOWN_3: 3,
SETTING_UNKNOWN_6: 6,
SETTING_UNKNOWN_8: 8,
SETTING_UNKNOWN_9: 9,
// ===== CursorSetting.Unknown6 =====
SETTING6_FIELD_1: 1,
SETTING6_FIELD_2: 2,
// ===== Metadata =====
META_PLATFORM: 1,
META_ARCH: 2,
META_VERSION: 3,
META_CWD: 4,
META_TIMESTAMP: 5,
// ===== MessageId =====
MSGID_ID: 1,
MSGID_SUMMARY: 2,
MSGID_ROLE: 3,
// ===== MCPTool =====
MCP_TOOL_NAME: 1,
MCP_TOOL_DESC: 2,
MCP_TOOL_PARAMS: 3,
MCP_TOOL_SERVER: 4,
// ===== StreamUnifiedChatResponseWithTools (response) =====
TOOL_CALL: 1,
RESPONSE: 2,
// ===== ClientSideToolV2Call =====
TOOL_ID: 3,
TOOL_NAME: 9,
TOOL_RAW_ARGS: 10,
TOOL_IS_LAST: 11,
TOOL_MCP_PARAMS: 27,
// ===== MCPParams =====
MCP_TOOLS_LIST: 1,
// ===== MCPParams.Tool (nested) =====
MCP_NESTED_NAME: 1,
MCP_NESTED_PARAMS: 3,
// ===== StreamUnifiedChatResponse =====
RESPONSE_TEXT: 1,
THINKING: 25,
// ===== Thinking =====
THINKING_TEXT: 1,
} as const;
/** Type definitions */
export type WireType = (typeof WIRE_TYPE)[keyof typeof WIRE_TYPE];
export type RoleType = (typeof ROLE)[keyof typeof ROLE];
export type UnifiedModeType = (typeof UNIFIED_MODE)[keyof typeof UNIFIED_MODE];
export type ThinkingLevelType = (typeof THINKING_LEVEL)[keyof typeof THINKING_LEVEL];
export type FieldNumber = (typeof FIELD)[keyof typeof FIELD];
/** Cursor credentials structure */
export interface CursorCredentials {
accessToken: string;
machineId: string;
ghostMode?: boolean;
}
/** Cursor tool definition */
export interface CursorTool {
function?: {
name?: string;
description?: string;
parameters?: Record<string, unknown>;
};
name?: string;
description?: string;
input_schema?: Record<string, unknown>;
}
/** Cursor tool result */
export interface CursorToolResult {
tool_call_id?: string;
name?: string;
index?: number;
raw_args?: string;
}
/** Cursor message format */
export interface CursorMessage {
role: string;
content: string;
tool_results?: CursorToolResult[];
tool_calls?: Array<{
id: string;
type: string;
function: {
name: string;
arguments: string;
};
}>;
}
/** Formatted message for encoding */
export interface FormattedMessage {
content: string;
role: RoleType;
messageId: string;
isLast: boolean;
hasTools: boolean;
toolResults: CursorToolResult[];
}
/** Message ID structure */
export interface MessageId {
messageId: string;
role: RoleType;
}
/** Compression flags for ConnectRPC frames */
export const COMPRESS_FLAG = {
NONE: 0x00,
GZIP: 0x01,
GZIP_ALT: 0x02,
GZIP_BOTH: 0x03,
} as const;
+186
View File
@@ -0,0 +1,186 @@
/**
* Cursor Protobuf Main Module
* Exports encoder/decoder functions and builds complete requests
*/
import { randomUUID } from 'crypto';
import {
ROLE,
UNIFIED_MODE,
THINKING_LEVEL,
FIELD,
type CursorMessage,
type CursorTool,
type FormattedMessage,
type MessageId,
type ThinkingLevelType,
} from './cursor-protobuf-schema.js';
import {
encodeField,
encodeVarint,
encodeMessage,
encodeInstruction,
encodeModel,
encodeCursorSetting,
encodeMetadata,
encodeMessageId,
encodeMcpTool,
wrapConnectRPCFrame,
concatArrays,
} from './cursor-protobuf-encoder.js';
import {
decodeVarint,
decodeField,
decodeMessage,
parseConnectRPCFrame,
extractTextFromResponse,
} from './cursor-protobuf-decoder.js';
import { WIRE_TYPE } from './cursor-protobuf-schema.js';
/**
* Build complete chat request protobuf
*/
export function encodeRequest(
messages: CursorMessage[],
modelName: string,
tools: CursorTool[] = [],
reasoningEffort: string | null = null
): Uint8Array {
if (messages.length === 0) {
throw new Error('Messages array must not be empty');
}
const hasTools = tools?.length > 0;
const isAgentic = hasTools;
const formattedMessages: FormattedMessage[] = [];
const messageIds: MessageId[] = [];
// Prepare messages
for (let i = 0; i < messages.length; i++) {
const msg = messages[i];
const role = msg.role === 'user' ? ROLE.USER : ROLE.ASSISTANT;
const msgId = randomUUID();
const isLast = i === messages.length - 1;
formattedMessages.push({
content: msg.content,
role,
messageId: msgId,
isLast,
hasTools,
toolResults: msg.tool_results || [],
});
messageIds.push({ messageId: msgId, role });
}
// Map reasoning effort to thinking level
let thinkingLevel: ThinkingLevelType = THINKING_LEVEL.UNSPECIFIED;
if (reasoningEffort === 'medium') thinkingLevel = THINKING_LEVEL.MEDIUM;
else if (reasoningEffort === 'high') thinkingLevel = THINKING_LEVEL.HIGH;
// Build arrays for messages and tools
const messageFields = formattedMessages.map((fm) =>
encodeField(
FIELD.MESSAGES,
WIRE_TYPE.LEN,
encodeMessage(fm.content, fm.role, fm.messageId, fm.isLast, fm.hasTools, fm.toolResults)
)
);
const messageIdFields = messageIds.map((mid) =>
encodeField(FIELD.MESSAGE_IDS, WIRE_TYPE.LEN, encodeMessageId(mid.messageId, mid.role))
);
const toolFields =
tools?.length > 0
? tools.map((tool) => encodeField(FIELD.MCP_TOOLS, WIRE_TYPE.LEN, encodeMcpTool(tool)))
: [];
const supportedToolsField = isAgentic
? [encodeField(FIELD.SUPPORTED_TOOLS, WIRE_TYPE.LEN, encodeVarint(1))]
: [];
// Concatenate all parts
const parts: Uint8Array[] = [
...messageFields,
encodeField(FIELD.UNKNOWN_2, WIRE_TYPE.VARINT, 1),
encodeField(FIELD.INSTRUCTION, WIRE_TYPE.LEN, encodeInstruction('')),
encodeField(FIELD.UNKNOWN_4, WIRE_TYPE.VARINT, 1),
encodeField(FIELD.MODEL, WIRE_TYPE.LEN, encodeModel(modelName)),
encodeField(FIELD.WEB_TOOL, WIRE_TYPE.LEN, ''),
encodeField(FIELD.UNKNOWN_13, WIRE_TYPE.VARINT, 1),
encodeField(FIELD.CURSOR_SETTING, WIRE_TYPE.LEN, encodeCursorSetting()),
encodeField(FIELD.UNKNOWN_19, WIRE_TYPE.VARINT, 1),
encodeField(FIELD.CONVERSATION_ID, WIRE_TYPE.LEN, randomUUID()),
encodeField(FIELD.METADATA, WIRE_TYPE.LEN, encodeMetadata()),
encodeField(FIELD.IS_AGENTIC, WIRE_TYPE.VARINT, isAgentic ? 1 : 0),
...supportedToolsField,
...messageIdFields,
...toolFields,
encodeField(FIELD.LARGE_CONTEXT, WIRE_TYPE.VARINT, 0),
encodeField(FIELD.UNKNOWN_38, WIRE_TYPE.VARINT, 0),
encodeField(
FIELD.UNIFIED_MODE,
WIRE_TYPE.VARINT,
isAgentic ? UNIFIED_MODE.AGENT : UNIFIED_MODE.CHAT
),
encodeField(FIELD.UNKNOWN_47, WIRE_TYPE.LEN, ''),
encodeField(FIELD.SHOULD_DISABLE_TOOLS, WIRE_TYPE.VARINT, isAgentic ? 0 : 1),
encodeField(FIELD.THINKING_LEVEL, WIRE_TYPE.VARINT, thinkingLevel),
encodeField(FIELD.UNKNOWN_51, WIRE_TYPE.VARINT, 0),
encodeField(FIELD.UNKNOWN_53, WIRE_TYPE.VARINT, 1),
encodeField(FIELD.UNIFIED_MODE_NAME, WIRE_TYPE.LEN, isAgentic ? 'Agent' : 'Ask'),
];
return concatArrays(...parts);
}
/**
* Build chat request wrapped in top-level message
*/
export function buildChatRequest(
messages: CursorMessage[],
modelName: string,
tools: CursorTool[] = [],
reasoningEffort: string | null = null
): Uint8Array {
return encodeField(
FIELD.REQUEST,
WIRE_TYPE.LEN,
encodeRequest(messages, modelName, tools, reasoningEffort)
);
}
/**
* Generate complete Cursor request body with ConnectRPC framing
*/
export function generateCursorBody(
messages: CursorMessage[],
modelName: string,
tools: CursorTool[] = [],
reasoningEffort: string | null = null
): Uint8Array {
const protobuf = buildChatRequest(messages, modelName, tools, reasoningEffort);
const framed = wrapConnectRPCFrame(protobuf, false); // Cursor doesn't support compressed requests
return framed;
}
// Re-export all functions
export {
encodeVarint,
encodeField,
encodeMessage,
encodeInstruction,
encodeModel,
encodeCursorSetting,
encodeMetadata,
encodeMessageId,
encodeMcpTool,
wrapConnectRPCFrame,
decodeVarint,
decodeField,
decodeMessage,
parseConnectRPCFrame,
extractTextFromResponse,
};
+155
View File
@@ -0,0 +1,155 @@
/**
* OpenAI to Cursor Request Translator
* Converts OpenAI messages to Cursor format
*/
import type { CursorMessage, CursorToolResult, CursorTool } from './cursor-protobuf-schema.js';
/** OpenAI message format */
interface OpenAIMessage {
role: string;
content: string | Array<{ type: string; text?: string }>;
name?: string;
tool_call_id?: string;
tool_calls?: Array<{
id: string;
type: string;
function: { name: string; arguments: string };
}>;
}
/** OpenAI request body */
interface OpenAIRequestBody {
messages: OpenAIMessage[];
tools?: CursorTool[];
reasoning_effort?: string;
}
/**
* Convert OpenAI messages to Cursor format with native tool_results support
* - system → user with [System Instructions] prefix
* - tool → accumulate into tool_results array for next user/assistant message
* - assistant with tool_calls → keep tool_calls structure (Cursor supports it natively)
*/
function convertMessages(messages: OpenAIMessage[]): CursorMessage[] {
const result: CursorMessage[] = [];
let pendingToolResults: CursorToolResult[] = [];
for (let i = 0; i < messages.length; i++) {
const msg = messages[i];
if (msg.role === 'system') {
let content = '';
if (typeof msg.content === 'string') {
content = msg.content;
} else if (Array.isArray(msg.content)) {
for (const part of msg.content) {
if (part.type === 'text' && part.text) content += part.text;
}
}
result.push({
role: 'user',
content: `[System Instructions]\n${content}`,
});
continue;
}
if (msg.role === 'tool') {
let toolContent = '';
if (typeof msg.content === 'string') {
toolContent = msg.content;
} else if (Array.isArray(msg.content)) {
for (const part of msg.content) {
if (part.type === 'text' && part.text) {
toolContent += part.text;
}
}
}
const toolName = msg.name || 'tool';
const toolCallId = msg.tool_call_id || '';
// Accumulate tool result
pendingToolResults.push({
tool_call_id: toolCallId,
name: toolName,
index: pendingToolResults.length,
raw_args: toolContent,
});
continue;
}
if (msg.role === 'user' || msg.role === 'assistant') {
let content = '';
if (typeof msg.content === 'string') {
content = msg.content;
} else if (Array.isArray(msg.content)) {
for (const part of msg.content) {
if (part.type === 'text' && part.text) {
content += part.text;
}
}
}
// Keep tool_calls structure for assistant messages
if (msg.role === 'assistant' && msg.tool_calls && msg.tool_calls.length > 0) {
const assistantMsg: CursorMessage = { role: 'assistant', content: '' };
if (content) {
assistantMsg.content = content;
}
assistantMsg.tool_calls = msg.tool_calls;
// Attach pending tool results to assistant message with tool_calls
if (pendingToolResults.length > 0) {
assistantMsg.tool_results = pendingToolResults;
pendingToolResults = [];
}
result.push(assistantMsg);
} else if (content || pendingToolResults.length > 0) {
const msgObj: CursorMessage = {
role: msg.role,
content: content || '',
};
// Attach pending tool results to this message
if (pendingToolResults.length > 0) {
msgObj.tool_results = pendingToolResults;
pendingToolResults = [];
}
result.push(msgObj);
}
continue;
}
// Unknown role - skip with debug warning
if (process.env.CCS_DEBUG) {
console.error(`[cursor] Unknown message role: ${msg.role}, skipping`);
}
}
return result;
}
/**
* Transform OpenAI request to Cursor format
* Returns modified body with converted messages
*/
export function buildCursorRequest(
_model: string,
body: OpenAIRequestBody,
_stream: boolean,
_credentials: unknown
): {
messages: CursorMessage[];
tools?: CursorTool[];
} {
const messages = convertMessages(body.messages || []);
return {
...body,
messages,
};
}
+657
View File
@@ -0,0 +1,657 @@
/**
* Cursor Protobuf Module Unit Tests
* Tests encoder, decoder, translator, and executor components
*/
import { describe, it, expect } from 'bun:test';
import {
encodeVarint,
encodeField,
wrapConnectRPCFrame,
concatArrays,
} from '../../../src/cursor/cursor-protobuf-encoder';
import {
decodeVarint,
decodeField,
parseConnectRPCFrame,
} from '../../../src/cursor/cursor-protobuf-decoder';
import { buildCursorRequest } from '../../../src/cursor/cursor-translator';
import { generateCursorBody } from '../../../src/cursor/cursor-protobuf';
import { CursorExecutor } from '../../../src/cursor/cursor-executor';
import { WIRE_TYPE, FIELD } from '../../../src/cursor/cursor-protobuf-schema';
describe('Protobuf Encoding/Decoding', () => {
describe('encodeVarint / decodeVarint round-trip', () => {
it('should encode and decode 0', () => {
const encoded = encodeVarint(0);
const [decoded, offset] = decodeVarint(encoded, 0);
expect(decoded).toBe(0);
expect(offset).toBe(1);
});
it('should encode and decode 1', () => {
const encoded = encodeVarint(1);
const [decoded, offset] = decodeVarint(encoded, 0);
expect(decoded).toBe(1);
expect(offset).toBe(1);
});
it('should encode and decode 127', () => {
const encoded = encodeVarint(127);
const [decoded, offset] = decodeVarint(encoded, 0);
expect(decoded).toBe(127);
expect(offset).toBe(1);
});
it('should encode and decode 128', () => {
const encoded = encodeVarint(128);
const [decoded, offset] = decodeVarint(encoded, 0);
expect(decoded).toBe(128);
expect(offset).toBe(2);
});
it('should encode and decode 16383', () => {
const encoded = encodeVarint(16383);
const [decoded, offset] = decodeVarint(encoded, 0);
expect(decoded).toBe(16383);
expect(offset).toBe(2);
});
it('should encode and decode 0xFFFFFFFF', () => {
const encoded = encodeVarint(0xffffffff);
const [decoded, offset] = decodeVarint(encoded, 0);
expect(decoded).toBe(0xffffffff);
expect(offset).toBe(5);
});
});
describe('encodeField / decodeField round-trip', () => {
it('should encode and decode VARINT field', () => {
const fieldNum = 5;
const value = 42;
const encoded = encodeField(fieldNum, WIRE_TYPE.VARINT, value);
const [decodedFieldNum, wireType, decodedValue, offset] = decodeField(encoded, 0);
expect(decodedFieldNum).toBe(fieldNum);
expect(wireType).toBe(WIRE_TYPE.VARINT);
expect(decodedValue).toBe(value);
expect(offset).toBe(encoded.length);
});
it('should encode and decode LEN field with string', () => {
const fieldNum = 10;
const value = 'Hello, World!';
const encoded = encodeField(fieldNum, WIRE_TYPE.LEN, value);
const [decodedFieldNum, wireType, decodedValue, offset] = decodeField(encoded, 0);
expect(decodedFieldNum).toBe(fieldNum);
expect(wireType).toBe(WIRE_TYPE.LEN);
expect(new TextDecoder().decode(decodedValue as Uint8Array)).toBe(value);
expect(offset).toBe(encoded.length);
});
it('should encode and decode LEN field with binary data', () => {
const fieldNum = 15;
const value = new Uint8Array([1, 2, 3, 4, 5]);
const encoded = encodeField(fieldNum, WIRE_TYPE.LEN, value);
const [decodedFieldNum, wireType, decodedValue, offset] = decodeField(encoded, 0);
expect(decodedFieldNum).toBe(fieldNum);
expect(wireType).toBe(WIRE_TYPE.LEN);
expect(decodedValue).toEqual(value);
expect(offset).toBe(encoded.length);
});
});
describe('wrapConnectRPCFrame / parseConnectRPCFrame round-trip', () => {
it('should wrap and parse uncompressed frame', () => {
const payload = new Uint8Array([1, 2, 3, 4, 5, 6, 7, 8, 9, 10]);
const frame = wrapConnectRPCFrame(payload, false);
const parsed = parseConnectRPCFrame(Buffer.from(frame));
expect(parsed).not.toBeNull();
expect(parsed!.flags).toBe(0x00);
expect(parsed!.length).toBe(payload.length);
expect(parsed!.payload).toEqual(payload);
expect(parsed!.consumed).toBe(5 + payload.length);
});
it('should wrap and parse compressed frame', () => {
const payload = new Uint8Array([1, 2, 3, 4, 5, 6, 7, 8, 9, 10]);
const frame = wrapConnectRPCFrame(payload, true);
const parsed = parseConnectRPCFrame(Buffer.from(frame));
expect(parsed).not.toBeNull();
expect(parsed!.flags).toBe(0x01); // GZIP flag
expect(parsed!.payload).toEqual(payload); // Should be decompressed
});
it('should handle incomplete frame', () => {
const partial = new Uint8Array([0x00, 0x00, 0x00]); // Only 3 bytes
const parsed = parseConnectRPCFrame(Buffer.from(partial));
expect(parsed).toBeNull();
});
});
describe('concatArrays', () => {
it('should concatenate multiple arrays', () => {
const arr1 = new Uint8Array([1, 2, 3]);
const arr2 = new Uint8Array([4, 5]);
const arr3 = new Uint8Array([6, 7, 8, 9]);
const result = concatArrays(arr1, arr2, arr3);
expect(result).toEqual(new Uint8Array([1, 2, 3, 4, 5, 6, 7, 8, 9]));
});
it('should handle empty arrays', () => {
const arr1 = new Uint8Array([1, 2]);
const arr2 = new Uint8Array([]);
const arr3 = new Uint8Array([3, 4]);
const result = concatArrays(arr1, arr2, arr3);
expect(result).toEqual(new Uint8Array([1, 2, 3, 4]));
});
});
});
describe('Message Translation', () => {
describe('buildCursorRequest', () => {
it('should convert system message to user with prefix', () => {
const result = buildCursorRequest(
'gpt-4',
{
messages: [{ role: 'system', content: 'You are a helpful assistant.' }],
},
false,
{}
);
expect(result.messages).toHaveLength(1);
expect(result.messages[0].role).toBe('user');
expect(result.messages[0].content).toContain('[System Instructions]');
expect(result.messages[0].content).toContain('You are a helpful assistant.');
});
it('should keep user and assistant messages', () => {
const result = buildCursorRequest(
'gpt-4',
{
messages: [
{ role: 'user', content: 'Hello' },
{ role: 'assistant', content: 'Hi there!' },
],
},
false,
{}
);
expect(result.messages).toHaveLength(2);
expect(result.messages[0].role).toBe('user');
expect(result.messages[0].content).toBe('Hello');
expect(result.messages[1].role).toBe('assistant');
expect(result.messages[1].content).toBe('Hi there!');
});
it('should handle assistant messages with tool_calls', () => {
const result = buildCursorRequest(
'gpt-4',
{
messages: [
{
role: 'assistant',
content: '',
tool_calls: [
{
id: 'call_123',
type: 'function',
function: { name: 'get_weather', arguments: '{"city":"NYC"}' },
},
],
},
],
},
false,
{}
);
expect(result.messages).toHaveLength(1);
expect(result.messages[0].role).toBe('assistant');
expect(result.messages[0].tool_calls).toHaveLength(1);
expect(result.messages[0].tool_calls![0].id).toBe('call_123');
expect(result.messages[0].tool_calls![0].function.name).toBe('get_weather');
});
it('should accumulate tool results', () => {
const result = buildCursorRequest(
'gpt-4',
{
messages: [
{
role: 'assistant',
content: '',
tool_calls: [
{
id: 'call_123',
type: 'function',
function: { name: 'get_weather', arguments: '{"city":"NYC"}' },
},
],
},
{
role: 'tool',
content: '{"temperature": 72}',
name: 'get_weather',
tool_call_id: 'call_123',
},
{ role: 'user', content: 'What is the weather?' },
],
},
false,
{}
);
expect(result.messages).toHaveLength(2);
// Tool result should be attached to next message
expect(result.messages[1].tool_results).toBeDefined();
expect(result.messages[1].tool_results).toHaveLength(1);
expect(result.messages[1].tool_results![0].tool_call_id).toBe('call_123');
});
it('should handle array content format', () => {
const result = buildCursorRequest(
'gpt-4',
{
messages: [
{
role: 'user',
content: [
{ type: 'text', text: 'Hello' },
{ type: 'text', text: ' World' },
],
},
],
},
false,
{}
);
expect(result.messages).toHaveLength(1);
expect(result.messages[0].content).toBe('Hello World');
});
it('should handle system message with array content format', () => {
const result = buildCursorRequest(
'gpt-4',
{
messages: [
{
role: 'system',
content: [
{ type: 'text', text: 'System instruction part 1' },
{ type: 'text', text: ' part 2' },
],
},
],
},
false,
{}
);
expect(result.messages).toHaveLength(1);
expect(result.messages[0].role).toBe('user');
expect(result.messages[0].content).toBe('[System Instructions]\nSystem instruction part 1 part 2');
});
});
});
describe('Request Encoding', () => {
describe('generateCursorBody', () => {
it('should encode basic text message', () => {
const result = generateCursorBody([{ role: 'user', content: 'Hello' }], 'gpt-4', [], null);
expect(result).toBeInstanceOf(Uint8Array);
expect(result.length).toBeGreaterThan(0);
});
it('should encode message with tools', () => {
const tools = [
{
type: 'function' as const,
function: {
name: 'get_weather',
description: 'Get weather data',
parameters: {
type: 'object',
properties: {
city: { type: 'string' },
},
required: ['city'],
},
},
},
];
const result = generateCursorBody([{ role: 'user', content: 'What is the weather?' }], 'gpt-4', tools, null);
expect(result).toBeInstanceOf(Uint8Array);
expect(result.length).toBeGreaterThan(0);
});
});
describe('Edge cases', () => {
it('should handle malformed frame gracefully', () => {
const executor = new CursorExecutor();
// Incomplete frame header (only 3 bytes instead of 5)
const incompleteFrame = Buffer.from([0x00, 0x00, 0x00]);
const result = executor.transformProtobufToJSON(incompleteFrame, 'gpt-4', {
messages: [],
});
// Should return valid response even with malformed input
expect(result.status).toBe(200);
});
it('should handle truncated payload', () => {
const executor = new CursorExecutor();
// Frame header says payload is 100 bytes but only 5 bytes follow
const truncatedFrame = Buffer.from([0x00, 0x00, 0x00, 0x00, 0x64, 0x01, 0x02, 0x03, 0x04, 0x05]);
const result = executor.transformProtobufToJSON(truncatedFrame, 'gpt-4', {
messages: [],
});
// Should handle gracefully
expect(result.status).toBe(200);
});
it('should handle multi-frame buffer', () => {
const executor = new CursorExecutor();
// Create two simple frames
const frame1 = wrapConnectRPCFrame(
encodeField(FIELD.RESPONSE_TEXT, WIRE_TYPE.LEN, 'Frame 1'),
false
);
const frame2 = wrapConnectRPCFrame(
encodeField(FIELD.RESPONSE_TEXT, WIRE_TYPE.LEN, ' Frame 2'),
false
);
// Concatenate them
const multiFrame = Buffer.concat([Buffer.from(frame1), Buffer.from(frame2)]);
const result = executor.transformProtobufToJSON(multiFrame, 'gpt-4', {
messages: [],
});
expect(result.status).toBe(200);
});
});
});
describe('CursorExecutor', () => {
const executor = new CursorExecutor();
describe('generateChecksum', () => {
it('should generate valid checksum format', () => {
const machineId = 'test-machine-id';
const checksum = executor.generateChecksum(machineId);
// Should end with machine ID
expect(checksum.endsWith(machineId)).toBe(true);
// Should have base64url-like prefix (8 chars from 6 bytes)
const prefix = checksum.slice(0, -machineId.length);
expect(prefix.length).toBe(8);
expect(/^[A-Za-z0-9_-]+$/.test(prefix)).toBe(true);
});
it('should generate valid checksums at different call times', async () => {
const machineId = 'test-machine-id';
const checksum1 = executor.generateChecksum(machineId);
// Wait to ensure timestamp may change (though timestamp granularity is ~16 min)
await new Promise((resolve) => setTimeout(resolve, 10));
const checksum2 = executor.generateChecksum(machineId);
// Verify both checksums are valid (may be same due to timestamp granularity)
expect(checksum1.endsWith(machineId)).toBe(true);
expect(checksum2.endsWith(machineId)).toBe(true);
});
});
describe('buildHeaders', () => {
it('should generate all required headers', () => {
const credentials = {
accessToken: 'test-token',
machineId: 'test-machine-id',
};
const headers = executor.buildHeaders(credentials);
expect(headers).toHaveProperty('authorization');
expect(headers.authorization).toContain('Bearer');
expect(headers).toHaveProperty('connect-accept-encoding', 'gzip');
expect(headers).toHaveProperty('connect-protocol-version', '1');
expect(headers).toHaveProperty('content-type', 'application/connect+proto');
expect(headers).toHaveProperty('user-agent', 'connect-es/1.6.1');
expect(headers).toHaveProperty('x-cursor-checksum');
expect(headers).toHaveProperty('x-cursor-client-version', '2.3.41');
expect(headers).toHaveProperty('x-cursor-client-type', 'ide');
expect(headers).toHaveProperty('x-ghost-mode', 'true');
});
it('should handle token with :: delimiter', () => {
const credentials = {
accessToken: 'prefix::actual-token',
machineId: 'test-machine-id',
};
const headers = executor.buildHeaders(credentials);
expect(headers.authorization).toBe('Bearer actual-token');
});
it('should respect ghostMode flag', () => {
const credentialsGhost = {
accessToken: 'test-token',
machineId: 'test-machine-id',
ghostMode: true,
};
const credentialsNoGhost = {
accessToken: 'test-token',
machineId: 'test-machine-id',
ghostMode: false,
};
const headersGhost = executor.buildHeaders(credentialsGhost);
const headersNoGhost = executor.buildHeaders(credentialsNoGhost);
expect(headersGhost['x-ghost-mode']).toBe('true');
expect(headersNoGhost['x-ghost-mode']).toBe('false');
});
it('should throw error if machineId missing', () => {
const credentials = {
accessToken: 'test-token',
machineId: '',
};
expect(() => executor.buildHeaders(credentials)).toThrow('Machine ID is required');
});
});
describe('buildUrl', () => {
it('should return correct API endpoint', () => {
const url = executor.buildUrl();
expect(url).toBe('https://api2.cursor.sh/aiserver.v1.AiService/StreamChat');
});
});
describe('transformProtobufToJSON', () => {
it('should handle basic text response', async () => {
// Create minimal protobuf response with text
const textContent = 'Hello, world!';
const responseField = encodeField(FIELD.RESPONSE_TEXT, WIRE_TYPE.LEN, textContent);
const responseMsg = encodeField(FIELD.RESPONSE, WIRE_TYPE.LEN, responseField);
const frame = wrapConnectRPCFrame(responseMsg, false);
const result = executor.transformProtobufToJSON(Buffer.from(frame), 'gpt-4', {
messages: [],
});
expect(result.status).toBe(200);
const bodyText = await result.text();
const body = JSON.parse(bodyText);
expect(body.choices[0].message.content).toBe(textContent);
expect(body.choices[0].finish_reason).toBe('stop');
});
it('should handle JSON error response', async () => {
const errorJson = JSON.stringify({
error: {
code: 'resource_exhausted',
message: 'Rate limit exceeded',
},
});
const frame = wrapConnectRPCFrame(new TextEncoder().encode(errorJson), false);
const result = executor.transformProtobufToJSON(Buffer.from(frame), 'gpt-4', {
messages: [],
});
expect(result.status).toBe(429);
const bodyText = await result.text();
const body = JSON.parse(bodyText);
expect(body.error.type).toBe('rate_limit_error');
});
});
describe('transformProtobufToSSE', () => {
it('should output SSE format', async () => {
// Create minimal protobuf response with text
const textContent = 'Hello';
const responseField = encodeField(FIELD.RESPONSE_TEXT, WIRE_TYPE.LEN, textContent);
const responseMsg = encodeField(FIELD.RESPONSE, WIRE_TYPE.LEN, responseField);
const frame = wrapConnectRPCFrame(responseMsg, false);
const result = executor.transformProtobufToSSE(Buffer.from(frame), 'gpt-4', {
messages: [],
});
expect(result.status).toBe(200);
expect(result.headers.get('content-type')).toBe('text/event-stream');
const bodyText = await result.text();
expect(bodyText).toContain('data: ');
expect(bodyText).toContain('data: [DONE]');
expect(bodyText).toContain(textContent);
});
it('should handle JSON error response', async () => {
const errorJson = JSON.stringify({
error: {
code: 'resource_exhausted',
message: 'Rate limit exceeded',
},
});
const frame = wrapConnectRPCFrame(new TextEncoder().encode(errorJson), false);
const result = executor.transformProtobufToSSE(Buffer.from(frame), 'gpt-4', {
messages: [],
});
expect(result.status).toBe(429);
const bodyText = await result.text();
const body = JSON.parse(bodyText);
expect(body.error.type).toBe('rate_limit_error');
});
});
describe('decompressPayload error handling', () => {
it('should return empty buffer on decompression failure', () => {
// Create invalid gzip data
const invalidGzip = Buffer.from([0x1f, 0x8b, 0x08, 0x00, 0xff, 0xff]);
const frame = new Uint8Array(5 + invalidGzip.length);
frame[0] = 0x01; // GZIP flag
frame[1] = 0;
frame[2] = 0;
frame[3] = 0;
frame[4] = invalidGzip.length;
frame.set(invalidGzip, 5);
const result = executor.transformProtobufToJSON(Buffer.from(frame), 'gpt-4', {
messages: [],
});
// Should handle gracefully and return valid response
expect(result.status).toBe(200);
});
});
describe('error handling', () => {
it('should return empty buffer on decompression failure', () => {
const executor = new CursorExecutor();
// Invalid compressed payload (not actually gzipped)
const invalidGzipPayload = new Uint8Array([1, 2, 3, 4, 5]);
const flags = 0x01; // GZIP flag
// Wrap with ConnectRPC frame header (flags + length)
const length = invalidGzipPayload.length;
const frame = new Uint8Array(5 + length);
frame[0] = flags;
frame[1] = (length >> 24) & 0xff;
frame[2] = (length >> 16) & 0xff;
frame[3] = (length >> 8) & 0xff;
frame[4] = length & 0xff;
frame.set(invalidGzipPayload, 5);
const buffer = Buffer.from(frame);
// Should not crash - decompression failure returns empty buffer
const result = executor.transformProtobufToJSON(buffer, 'test-model', {
messages: [],
stream: false,
});
expect(result.status).toBe(200);
});
it('should log unknown message roles in debug mode', () => {
const originalDebug = process.env.CCS_DEBUG;
process.env.CCS_DEBUG = '1';
const consoleSpy: string[] = [];
const originalError = console.error;
console.error = (...args: unknown[]) => {
const msg = args.map((a) => String(a)).join(' ');
consoleSpy.push(msg);
};
try {
const messages = [
{
role: 'unknown_role' as 'user', // Type assertion to bypass TS
content: 'test',
},
];
// buildCursorRequest expects (model, body, stream, credentials)
buildCursorRequest('test-model', { messages }, false, { machineId: '12345', accessToken: 'test' });
// Should have logged warning
const hasWarning = consoleSpy.some((log) => log.includes('Unknown message role'));
expect(hasWarning).toBe(true);
} finally {
console.error = originalError;
process.env.CCS_DEBUG = originalDebug;
}
});
});
});