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([]); + }); +});