From 271e7ba892bfd61a6be787596abfffcb6b93df98 Mon Sep 17 00:00:00 2001 From: bnivanov <251194612+bnivanov@users.noreply.github.com> Date: Tue, 18 Aug 2026 14:48:45 +0200 Subject: [PATCH] fix(cursor): answer interactionQuery so hosted fetch can continue Cursor hosted web search / Exa / unnamed field-9 WebFetch send interaction_query and block the Run RPC until the client writes interaction_response. Dropping the frame leaves the HTTP/2 stream alive on heartbeats that are not semantic progress, so the 300s idle watchdog aborts with "Provider stream stalled while waiting for the next event". Approve network permission gates and reject interactive ask / switch-mode / create-plan. Leave VM setup unanswered rather than inventing a success result. --- docs/provider-quirks.md | 8 +- packages/ai/CHANGELOG.md | 4 + packages/ai/src/providers/cursor.ts | 3 + .../src/providers/cursor/interaction-query.ts | 200 ++++++++++++++++ .../ai/test/cursor-interaction-query.test.ts | 219 ++++++++++++++++++ 5 files changed, 432 insertions(+), 2 deletions(-) create mode 100644 packages/ai/src/providers/cursor/interaction-query.ts create mode 100644 packages/ai/test/cursor-interaction-query.test.ts diff --git a/docs/provider-quirks.md b/docs/provider-quirks.md index 08a1a4c60..1dfba11f5 100644 --- a/docs/provider-quirks.md +++ b/docs/provider-quirks.md @@ -450,8 +450,12 @@ Cursor's integration in `packages/ai` operates over an HTTP/2 Connect RPC transp - **Trailer & Transport Error Handling**: - Monitors HTTP/2 trailers (`grpc-status`, `grpc-message`) and maps socket or TLS disconnects using `mapH2TransportError`. - **Bi-Directional RPC Dispatch**: - - Server streams `AgentServerMessage` (`interactionUpdate`, `execServerMessage`, `kvServerMessage`). - - Client writes `AgentClientMessage` (`runRequest`, periodic `clientHeartbeat` every 5 seconds) and `ExecClientMessage` tool responses (`readResult`, `writeResult`, `execClientThrow`, `requestContextResult`). + - Server streams `AgentServerMessage` (`interactionUpdate`, `execServerMessage`, `kvServerMessage`, `interactionQuery`). + - Client writes `AgentClientMessage` (`runRequest`, periodic `clientHeartbeat` every 5 seconds, `interactionResponse`) and `ExecClientMessage` tool responses (`readResult`, `writeResult`, `execClientThrow`, `requestContextResult`). +- **Interaction Query Handshake**: + - Hosted web search / Exa / unnamed field-9 WebFetch send `interactionQuery` and block the turn until the client writes `interactionResponse`. + - Heartbeats keep HTTP/2 alive but are not semantic progress; an unanswered query sits silent until the 300s idle watchdog (`Provider stream stalled while waiting for the next event`). + - `handleInteractionQuery` approves network permission gates and rejects interactive ask / switch-mode / create-plan. VM setup is left unanswered because its result oneof is success-only. - **Async Execution Drain & Turn Completion**: - `handleServerMessage` processes frames asynchronously so the socket continues draining. Dispatches are tracked in `inFlightDispatches` and bounded by `options.signal` abort handling before finalizing stream completion. - Stream completion verifies `turnEnded` (`sawTurnEnded`) or throws `incomplete-stream`. diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index e469ca4f6..20dff55ac 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Answer Cursor `interaction_query` permission gates (hosted web search, Exa, unnamed field-9 WebFetch) so the Run RPC continues instead of sitting silent until the 300s idle watchdog. + ## [17.3.7] - 2026-08-17 ### Changed diff --git a/packages/ai/src/providers/cursor.ts b/packages/ai/src/providers/cursor.ts index de00062c8..829733834 100644 --- a/packages/ai/src/providers/cursor.ts +++ b/packages/ai/src/providers/cursor.ts @@ -220,6 +220,7 @@ import { piReadPathHasRange, piTimeout, } from "./cursor/exec-modern"; +import { handleInteractionQuery } from "./cursor/interaction-query"; export const CURSOR_API_URL = "https://api2.cursor.sh"; export const CURSOR_CLIENT_VERSION = "cli-2026.07.23-e383d2b"; @@ -1024,6 +1025,8 @@ export async function handleServerMessage( state, ), ); + } else if (msgCase === "interactionQuery") { + handleInteractionQuery(msg.message.value, h2Request); } else if (msgCase === "conversationCheckpointUpdate") { handleConversationCheckpointUpdate(msg.message.value, output, usageState, onConversationCheckpoint); } diff --git a/packages/ai/src/providers/cursor/interaction-query.ts b/packages/ai/src/providers/cursor/interaction-query.ts new file mode 100644 index 000000000..5b5f454fd --- /dev/null +++ b/packages/ai/src/providers/cursor/interaction-query.ts @@ -0,0 +1,200 @@ +import type http2 from "node:http2"; +import { create, toBinary } from "@bufbuild/protobuf"; +import { + AgentClientMessageSchema, + AskQuestionInteractionResponseSchema, + AskQuestionRejectedSchema, + AskQuestionResultSchema, + CreatePlanErrorSchema, + CreatePlanRequestResponseSchema, + CreatePlanResultSchema, + ExaFetchRequestResponse_ApprovedSchema, + ExaFetchRequestResponseSchema, + ExaSearchRequestResponse_ApprovedSchema, + ExaSearchRequestResponseSchema, + type InteractionQuery, + type InteractionResponse, + InteractionResponseSchema, + SwitchModeRequestResponse_RejectedSchema, + SwitchModeRequestResponseSchema, + WebSearchRequestResponse_ApprovedSchema, + WebSearchRequestResponseSchema, +} from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb"; +import { $env } from "@oh-my-pi/pi-utils"; + +const NOT_IMPLEMENTED_SUFFIX = "not implemented by this client"; + +type ProtoUnknownField = { no: number; wireType: number; data: Uint8Array }; +type ProtoUnknownBag = { $unknown?: ProtoUnknownField[] }; + +type InteractionQueryCase = NonNullable; +type InteractionResult = Exclude; + +function frameConnectMessage(data: Uint8Array, flags = 0): Buffer { + const frame = Buffer.alloc(5 + data.length); + frame[0] = flags; + frame.writeUInt32BE(data.length, 1); + frame.set(data, 5); + return frame; +} + +function isProtoUnknownField(value: unknown): value is ProtoUnknownField { + if (!value || typeof value !== "object") return false; + if (!("no" in value) || !("wireType" in value) || !("data" in value)) return false; + return typeof value.no === "number" && typeof value.wireType === "number" && value.data instanceof Uint8Array; +} + +function protoUnknownFields(message: object): ProtoUnknownField[] { + if (!("$unknown" in message) || !Array.isArray(message.$unknown)) return []; + return message.$unknown.filter(isProtoUnknownField); +} + +function attachUnknownApprovedField(response: InteractionResponse, fieldNo: number): void { + const bag: ProtoUnknownBag = response; // protobuf-es unnamed oneof members live on $unknown + // protobuf-es writes tag + raw(data); LEN fields need the length prefix in data. + const field: ProtoUnknownField = { no: fieldNo, wireType: 2, data: new Uint8Array([0x02, 0x0a, 0x00]) }; + const existing = bag.$unknown; + if (Array.isArray(existing)) { + existing.push(field); + return; + } + bag.$unknown = [field]; +} + +function log(type: string, subtype?: string, data?: unknown): void { + if (!$env.DEBUG_CURSOR) return; + const verbose = $env.DEBUG_CURSOR === "2" || $env.DEBUG_CURSOR === "verbose"; + const dataStr = verbose && data ? ` ${JSON.stringify(data)?.slice(0, 500)}` : ""; + console.error(`[CURSOR] ${type}${subtype ? `: ${subtype}` : ""}${dataStr}`); +} + +function sendInteractionResponse(h2Request: http2.ClientHttp2Stream, queryId: number, result: InteractionResult): void { + const response = create(InteractionResponseSchema, { id: queryId, result }); + const clientMessage = create(AgentClientMessageSchema, { + message: { case: "interactionResponse", value: response }, + }); + h2Request.write(frameConnectMessage(toBinary(AgentClientMessageSchema, clientMessage))); + log("interactionResponse", result.case, { id: queryId }); +} + +function sendUnknownApprovedInteractionResponse( + h2Request: http2.ClientHttp2Stream, + queryId: number, + fieldNo: number, +): void { + // `approved {}` on the matching response oneof: field N, empty message whose + // first member is field 1 (`approved`) with an empty length-delimited payload. + const response = create(InteractionResponseSchema, { id: queryId }); + attachUnknownApprovedField(response, fieldNo); + const clientMessage = create(AgentClientMessageSchema, { + message: { case: "interactionResponse", value: response }, + }); + h2Request.write(frameConnectMessage(toBinary(AgentClientMessageSchema, clientMessage))); + log("interactionResponse", "unknownApproved", { id: queryId, field: fieldNo }); +} + +/** + * Answer a Cursor `interaction_query` so the Run RPC can continue. + * + * Hosted web search / Exa / unnamed permission gates (field 9 = WebFetch on + * current Cursor builds) block the turn until the client writes an + * `interaction_response`. Dropping the frame leaves the HTTP/2 stream alive + * on heartbeats that are not semantic progress; the lazy idle watchdog then + * aborts with "Provider stream stalled while waiting for the next event". + * + * Unsupported interactive queries are rejected so the server is not stranded. + * VM setup is left unanswered rather than reporting a fake success. + */ +export function handleInteractionQuery(query: InteractionQuery, h2Request: http2.ClientHttp2Stream): void { + const queryCase = query.query.case; + log("interactionQuery", queryCase, query.query.value); + if (!queryCase) { + const unknown = protoUnknownFields(query).find(field => field.wireType === 2 && field.no >= 2); + if (unknown) { + log("warn", "unknownInteractionQueryApproved", { id: query.id, field: unknown.no }); + sendUnknownApprovedInteractionResponse(h2Request, query.id, unknown.no); + return; + } + log("warn", "unknownInteractionQuery", { id: query.id }); + return; + } + + switch (queryCase) { + case "webSearchRequestQuery": + sendInteractionResponse(h2Request, query.id, { + case: "webSearchRequestResponse", + value: create(WebSearchRequestResponseSchema, { + result: { case: "approved", value: create(WebSearchRequestResponse_ApprovedSchema, {}) }, + }), + }); + return; + case "exaSearchRequestQuery": + sendInteractionResponse(h2Request, query.id, { + case: "exaSearchRequestResponse", + value: create(ExaSearchRequestResponseSchema, { + result: { case: "approved", value: create(ExaSearchRequestResponse_ApprovedSchema, {}) }, + }), + }); + return; + case "exaFetchRequestQuery": + sendInteractionResponse(h2Request, query.id, { + case: "exaFetchRequestResponse", + value: create(ExaFetchRequestResponseSchema, { + result: { case: "approved", value: create(ExaFetchRequestResponse_ApprovedSchema, {}) }, + }), + }); + return; + case "askQuestionInteractionQuery": + sendInteractionResponse(h2Request, query.id, { + case: "askQuestionInteractionResponse", + value: create(AskQuestionInteractionResponseSchema, { + result: create(AskQuestionResultSchema, { + result: { + case: "rejected", + value: create(AskQuestionRejectedSchema, { + reason: `Interactive questions are ${NOT_IMPLEMENTED_SUFFIX}`, + }), + }, + }), + }), + }); + return; + case "switchModeRequestQuery": + sendInteractionResponse(h2Request, query.id, { + case: "switchModeRequestResponse", + value: create(SwitchModeRequestResponseSchema, { + result: { + case: "rejected", + value: create(SwitchModeRequestResponse_RejectedSchema, { + reason: `Mode switches are ${NOT_IMPLEMENTED_SUFFIX}`, + }), + }, + }), + }); + return; + case "createPlanRequestQuery": + sendInteractionResponse(h2Request, query.id, { + case: "createPlanRequestResponse", + value: create(CreatePlanRequestResponseSchema, { + result: create(CreatePlanResultSchema, { + result: { + case: "error", + value: create(CreatePlanErrorSchema, { + error: `Plan files are ${NOT_IMPLEMENTED_SUFFIX}`, + }), + }, + }), + }), + }); + return; + case "setupVmEnvironmentArgs": + // Result oneof is success-only. Do not invent a VM; silence is better + // than a false SetupVmEnvironmentSuccess (review of #8047). + log("warn", "unhandledInteractionQuery", { queryCase, id: query.id }); + return; + default: { + const _exhaustive: InteractionQueryCase = queryCase; + log("warn", "unhandledInteractionQuery", { queryCase: _exhaustive, id: query.id }); + } + } +} diff --git a/packages/ai/test/cursor-interaction-query.test.ts b/packages/ai/test/cursor-interaction-query.test.ts new file mode 100644 index 000000000..f8b8a89e4 --- /dev/null +++ b/packages/ai/test/cursor-interaction-query.test.ts @@ -0,0 +1,219 @@ +import { describe, expect, it } from "bun:test"; +import { create, fromBinary } from "@bufbuild/protobuf"; +import { type BlockState, handleServerMessage, type ToolCallState } from "@oh-my-pi/pi-ai/providers/cursor"; +import type { AssistantMessage } from "@oh-my-pi/pi-ai/types"; +import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream"; +import { + type AgentClientMessage, + AgentClientMessageSchema, + AgentServerMessageSchema, + AskQuestionInteractionQuerySchema, + CreatePlanRequestQuerySchema, + ExaFetchRequestQuerySchema, + ExaSearchRequestQuerySchema, + type InteractionQuery, + InteractionQuerySchema, + SetupVmEnvironmentArgsSchema, + SwitchModeRequestQuerySchema, + WebSearchRequestQuerySchema, +} from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb"; + +type ProtoUnknownField = { no: number; wireType: number; data: Uint8Array }; +type ProtoUnknownBag = { $unknown?: ProtoUnknownField[] }; + +function cursorAssistantMessage(): AssistantMessage { + return { + role: "assistant", + api: "cursor-agent", + provider: "cursor", + model: "cursor-composer-2.5", + content: [], + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop", + timestamp: 1, + }; +} + +function emptyBlockState(): BlockState { + let textBlock: BlockState["currentTextBlock"] = null; + let thinkingBlock: BlockState["currentThinkingBlock"] = null; + let toolCall: ToolCallState | null = null; + return { + get currentTextBlock() { + return textBlock; + }, + get currentThinkingBlock() { + return thinkingBlock; + }, + get currentToolCall() { + return toolCall; + }, + openToolCalls: new Map(), + resolvedMcpToolCallIds: new Set(), + firstTokenTime: undefined, + setTextBlock: b => { + textBlock = b; + }, + setThinkingBlock: b => { + thinkingBlock = b; + }, + setToolCall: t => { + toolCall = t; + }, + setFirstTokenTime: () => {}, + }; +} + +function decodeConnectFrame(frame: Buffer): AgentClientMessage { + return fromBinary(AgentClientMessageSchema, frame.subarray(5)); +} + +async function dispatchQuery(query: InteractionQuery): Promise { + const frames: Buffer[] = []; + const h2Request = { + write(chunk: Buffer) { + frames.push(chunk); + return true; + }, + } as unknown as Parameters[5]; + const serverMsg = create(AgentServerMessageSchema, { + message: { case: "interactionQuery", value: query }, + }); + await handleServerMessage( + serverMsg, + cursorAssistantMessage(), + new AssistantMessageEventStream(), + emptyBlockState(), + new Map(), + h2Request, + undefined, + undefined, + { sawTokenDelta: false }, + [], + ); + return frames; +} + +describe("cursor interaction query handshake", () => { + it("approves hosted web search so the Run RPC can continue", async () => { + const frames = await dispatchQuery( + create(InteractionQuerySchema, { + id: 11, + query: { case: "webSearchRequestQuery", value: create(WebSearchRequestQuerySchema, {}) }, + }), + ); + expect(frames).toHaveLength(1); + const client = decodeConnectFrame(frames[0]!); + expect(client.message.case).toBe("interactionResponse"); + expect(client.message.value).toMatchObject({ + id: 11, + result: { case: "webSearchRequestResponse", value: { result: { case: "approved" } } }, + }); + }); + + it("approves hosted Exa search and fetch permission queries", async () => { + const search = await dispatchQuery( + create(InteractionQuerySchema, { + id: 12, + query: { case: "exaSearchRequestQuery", value: create(ExaSearchRequestQuerySchema, {}) }, + }), + ); + const fetch = await dispatchQuery( + create(InteractionQuerySchema, { + id: 13, + query: { case: "exaFetchRequestQuery", value: create(ExaFetchRequestQuerySchema, {}) }, + }), + ); + expect(decodeConnectFrame(search[0]!).message.value).toMatchObject({ + id: 12, + result: { case: "exaSearchRequestResponse", value: { result: { case: "approved" } } }, + }); + expect(decodeConnectFrame(fetch[0]!).message.value).toMatchObject({ + id: 13, + result: { case: "exaFetchRequestResponse", value: { result: { case: "approved" } } }, + }); + }); + + it("approves unnamed field-9 permission queries used by hosted WebFetch", async () => { + const query = create(InteractionQuerySchema, { id: 18 }); + const bag: ProtoUnknownBag = query; // protobuf-es unnamed query oneof (field 9) + bag.$unknown = [{ no: 9, wireType: 2, data: new Uint8Array([0x02, 0x0a, 0x00]) }]; + const frames = await dispatchQuery(query); + expect(frames).toHaveLength(1); + const client = decodeConnectFrame(frames[0]!); + expect(client.message.case).toBe("interactionResponse"); + if (client.message.case !== "interactionResponse") { + throw new Error("expected interactionResponse"); + } + expect(client.message.value.id).toBe(18); + const responseBag: ProtoUnknownBag = client.message.value; + expect(responseBag.$unknown?.some(field => field.no === 9 && field.wireType === 2)).toBe(true); + }); + + it("rejects interactive ask / switch-mode / create-plan queries", async () => { + const ask = decodeConnectFrame( + ( + await dispatchQuery( + create(InteractionQuerySchema, { + id: 14, + query: { case: "askQuestionInteractionQuery", value: create(AskQuestionInteractionQuerySchema, {}) }, + }), + ) + )[0]!, + ); + const mode = decodeConnectFrame( + ( + await dispatchQuery( + create(InteractionQuerySchema, { + id: 15, + query: { case: "switchModeRequestQuery", value: create(SwitchModeRequestQuerySchema, {}) }, + }), + ) + )[0]!, + ); + const plan = decodeConnectFrame( + ( + await dispatchQuery( + create(InteractionQuerySchema, { + id: 16, + query: { case: "createPlanRequestQuery", value: create(CreatePlanRequestQuerySchema, {}) }, + }), + ) + )[0]!, + ); + expect(ask.message.value).toMatchObject({ + id: 14, + result: { + case: "askQuestionInteractionResponse", + value: { result: { result: { case: "rejected" } } }, + }, + }); + expect(mode.message.value).toMatchObject({ + id: 15, + result: { case: "switchModeRequestResponse", value: { result: { case: "rejected" } } }, + }); + expect(plan.message.value).toMatchObject({ + id: 16, + result: { case: "createPlanRequestResponse", value: { result: { result: { case: "error" } } } }, + }); + }); + + it("does not invent a VM success or a reply for an empty query", async () => { + const vm = await dispatchQuery( + create(InteractionQuerySchema, { + id: 19, + query: { case: "setupVmEnvironmentArgs", value: create(SetupVmEnvironmentArgsSchema, {}) }, + }), + ); + const empty = await dispatchQuery(create(InteractionQuerySchema, { id: 17 })); + expect(vm).toEqual([]); + expect(empty).toEqual([]); + }); +});