Merge PR #8830: fix(ai): answer Cursor hosted WebFetch permission queries (@Unravl)

# Conflicts:
#	packages/ai/src/providers/cursor.ts
#	packages/ai/test/cursor-interaction-query.test.ts
This commit is contained in:
can1357
2026-08-19 01:38:07 +02:00
6 changed files with 1162 additions and 523 deletions
+1
View File
@@ -24,6 +24,7 @@
- Fixed the OpenCode Go login prompting for an "OpenCode Zen API key": the shared login flow now names the provider you selected, so connecting OpenCode Go asks for an OpenCode Go key (the `opencode.ai/auth` console is still shared, as documented upstream) ([#8738](https://github.com/can1357/oh-my-pi/issues/8738)).
- Fixed Anthropic-compatible endpoints with strict prompt validation (e.g. Z.AI GLM `api.z.ai/api/anthropic`, which rejects the whole request with `400 code 1213 "The prompt parameter was not received normally"`) failing sessions once a tool returned empty output on a vision-capable model: empty successful `tool_result` blocks now encode as `content: ""` instead of `content: []`, which both the official API and strict compatible endpoints accept.
- Fixed `retry.usageReservePct` (Reserve Margin) ignoring Claude Fable/Mythos weekly tier usage until it hit 100%, so a Fable model kept serving turns past the configured reserve; reserve health now honors the mapped tier row while credential-wide hard blocks still require confirmed exhaustion ([#8773](https://github.com/can1357/oh-my-pi/issues/8773)).
- Fixed `cursor-agent` streams stalling with "Provider stream stalled while waiting for the next event" when Cursor asked the client to approve a hosted WebFetch / web search (reproduced on `cursor-grok-4.6-xhigh` after "I'll fetch the page…"). Those `interaction_query` frames — including the newer WebFetch field 9 this proto did not name — were dropped, so the server waited forever and the idle watchdog aborted a live connection. Permission queries are now answered; hosted search/fetch is approved, unnamed permission fields get an `approved` reply on the same field number, and prompts this client cannot serve are rejected so the turn can continue.
## [17.3.5] - 2026-08-16
+257 -1
View File
@@ -12,6 +12,9 @@ import {
AgentServerMessageSchema,
AgentStoreConflictErrorSchema,
AgentStoreConflictResultSchema,
AskQuestionInteractionResponseSchema,
AskQuestionRejectedSchema,
AskQuestionResultSchema,
AssistantMessageSchema,
BackgroundShellSpawnResultSchema,
CanvasDiagnosticsErrorSchema,
@@ -26,6 +29,9 @@ import {
ConversationStateStructureSchema,
ConversationStepSchema,
ConversationTurnStructureSchema,
CreatePlanErrorSchema,
CreatePlanRequestResponseSchema,
CreatePlanResultSchema,
DeleteErrorSchema,
DeleteRejectedSchema,
DeleteResultSchema,
@@ -34,6 +40,10 @@ import {
DiagnosticsRejectedSchema,
DiagnosticsResultSchema,
DiagnosticsSuccessSchema,
ExaFetchRequestResponse_ApprovedSchema,
ExaFetchRequestResponseSchema,
ExaSearchRequestResponse_ApprovedSchema,
ExaSearchRequestResponseSchema,
ExecClientControlMessageSchema,
type ExecClientMessage,
ExecClientMessageSchema,
@@ -59,6 +69,9 @@ import {
GrepSuccessSchema,
type GrepUnionResult,
GrepUnionResultSchema,
type InteractionQuery,
type InteractionResponse,
InteractionResponseSchema,
KvClientMessageSchema,
type KvServerMessage,
ListMcpResourcesErrorSchema,
@@ -109,6 +122,8 @@ import {
SelectedContextSchema,
SelectedImageSchema,
SetBlobResultSchema,
SetupVmEnvironmentResultSchema,
SetupVmEnvironmentSuccessSchema,
ShellAllowlistPrecheckResultSchema,
type ShellArgs,
ShellFailureSchema,
@@ -128,11 +143,17 @@ import {
SubagentAwaitResultSchema,
SubagentErrorSchema,
SubagentResultSchema,
SwitchModeRequestResponse_RejectedSchema,
SwitchModeRequestResponseSchema,
ThinkingMessageSchema,
ToolCallSchema,
UserMessageActionSchema,
UserMessageSchema,
WebFetchAllowlistPrecheckResultSchema,
WebFetchRequestResponse_ApprovedSchema,
WebFetchRequestResponseSchema,
WebSearchRequestResponse_ApprovedSchema,
WebSearchRequestResponseSchema,
WriteErrorSchema,
WriteRejectedSchema,
WriteResultSchema,
@@ -941,7 +962,7 @@ export type ToolCallState = ToolCall & {
[kStreamingBlockIndex]: number;
[kStreamingPartialJson]?: string;
[kStreamingLastParseLen]?: number;
[kStreamingBlockKind]: "mcp" | "todo" | "cursor-exec" | "connect-scm";
[kStreamingBlockKind]: "mcp" | "todo" | "cursor-exec" | "connect-scm" | "web-fetch";
[kStreamingEnvelopeId]?: string;
[kCursorExecResolved]?: true;
};
@@ -1026,12 +1047,212 @@ export async function handleServerMessage(
),
);
} else if (msgCase === "interactionQuery") {
// Cursor asks the client to approve native web search / Exa fetch / etc.
// before it will continue the turn. Dropping the frame leaves the server
// waiting on a reply that never comes; the lazy idle watchdog then
// aborts a live stream with "Provider stream stalled while waiting for
// the next event" (cursor-grok-4.6-xhigh after a WebFetch/WebSearch
// permission prompt).
handleInteractionQuery(msg.message.value, h2Request);
} else if (msgCase === "conversationCheckpointUpdate") {
handleConversationCheckpointUpdate(msg.message.value, output, usageState, onConversationCheckpoint);
}
}
function handleInteractionQuery(query: InteractionQuery, h2Request: http2.ClientHttp2Stream): void {
const queryCase = query.query.case;
log("interactionQuery", queryCase, query.query.value);
if (!queryCase) {
// Newer Cursor builds add query variants this proto has not named yet
// (WebFetch was field 9). The server still blocks on a same-number
// InteractionResponse. Permission-shaped queries use Approved/Rejected
// like Exa fetch; answering `approved` unblocks the turn.
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":
// Permission gate, not "please run the search". Approve so Cursor
// performs the hosted search and the turn continues.
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 "webFetchRequestQuery":
// Hosted WebFetch permission prompt. Field 9 is what cursor-grok-4.6-xhigh
// sends after "I'll fetch the page…"; answering lets the server continue.
sendInteractionResponse(h2Request, query.id, {
case: "webFetchRequestResponse",
value: create(WebFetchRequestResponseSchema, {
result: { case: "approved", value: create(WebFetchRequestResponse_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 has only `success`. Answering is still required: silence
// strands the query id the same way an unanswered search approval does.
log("warn", "setupVmEnvironmentApprovedEmpty", { id: query.id });
sendInteractionResponse(h2Request, query.id, {
case: "setupVmEnvironmentResult",
value: create(SetupVmEnvironmentResultSchema, {
result: { case: "success", value: create(SetupVmEnvironmentSuccessSchema, {}) },
}),
});
return;
default: {
const _exhaustive: never = queryCase;
log("warn", "unhandledInteractionQuery", { queryCase: _exhaustive, id: query.id });
}
}
}
type ProtoUnknownField = { no: number; wireType: number; data: Uint8Array };
type HostedFetchCall = {
args?: { url?: string; toolCallId?: string };
result?: { result?: { case?: string; value?: { content?: string; error?: string; url?: string } } };
};
function selectHostedFetchCall(
toolCall: { tool?: { case?: string; value?: HostedFetchCall } } | undefined,
): HostedFetchCall | undefined {
const oneof = toolCall?.tool;
if (oneof?.case === "fetchToolCall" || oneof?.case === "webFetchToolCall") return oneof.value;
return undefined;
}
function hostedFetchUnknown(toolCall: object | undefined): boolean {
if (!toolCall) return false;
return (
protoUnknownFields(toolCall).some(field => field.no === 37) ||
protoUnknownFields((toolCall as { tool?: object }).tool ?? {}).some(field => field.no === 37)
);
}
function extractHttpUrlFromUnknown(message: object): string | undefined {
for (const field of protoUnknownFields(message)) {
const match = new TextDecoder().decode(field.data).match(/https?:\/\/[^\x00-\x1f]+/);
if (match) return match[0];
}
const nested = (message as { tool?: object }).tool;
return nested ? extractHttpUrlFromUnknown(nested) : undefined;
}
function describeHostedFetchResult(call: HostedFetchCall | undefined): { text: string; isError: boolean } {
const result = call?.result?.result;
if (result?.case === "success") {
return { text: result.value?.content || result.value?.url || "Fetched", isError: false };
}
if (result?.case === "error") {
return { text: result.value?.error || "Fetch failed", isError: true };
}
return { text: "Fetch completed", isError: false };
}
function protoUnknownFields(message: object): ProtoUnknownField[] {
const raw = (message as { $unknown?: ProtoUnknownField[] }).$unknown;
return Array.isArray(raw) ? raw : [];
}
function sendUnknownApprovedInteractionResponse(
h2Request: http2.ClientHttp2Stream,
queryId: number,
fieldNo: number,
): void {
// `approved {}` on the matching response oneof: field 1, empty message.
const response = create(InteractionResponseSchema, { id: queryId });
(response as { $unknown?: ProtoUnknownField[] }).$unknown = [
{ no: fieldNo, wireType: 2, data: new Uint8Array([0x0a, 0x00]) },
];
const clientMessage = create(AgentClientMessageSchema, {
message: { case: "interactionResponse", value: response },
});
h2Request.write(frameConnectMessage(toBinary(AgentClientMessageSchema, clientMessage)));
log("interactionResponse", "unknownApproved", { id: queryId, field: fieldNo });
}
function sendInteractionResponse(
h2Request: http2.ClientHttp2Stream,
queryId: number,
result: InteractionResponse["result"],
): 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 handleKvServerMessage(
kvMsg: KvServerMessage,
blobStore: Map<string, Uint8Array>,
@@ -3908,6 +4129,28 @@ export function processInteractionUpdate(
output.content.push(block);
retainStreamedCall(state, block, update.message.value.callId);
stream.push({ type: "toolcall_start", contentIndex: output.content.length - 1, partial: output });
return;
}
const fetchCall = selectHostedFetchCall(toolCall);
if (fetchCall || hostedFetchUnknown(toolCall)) {
// Hosted WebFetch / Fetch is permission-gated via InteractionQuery, then
// run server-side. Stamp resolved so agent-loop does not try a local tool.
const url = fetchCall?.args?.url || extractHttpUrlFromUnknown(toolCall);
const callId = fetchCall?.args?.toolCallId || update.message.value.callId || crypto.randomUUID();
const block: ToolCallState = {
type: "toolCall",
id: callId,
name: "web_fetch",
arguments: url ? { url } : {},
[kStreamingBlockIndex]: output.content.length,
[kStreamingBlockKind]: "web-fetch",
[kStreamingEnvelopeId]: update.message.value.callId || undefined,
[kCursorExecResolved]: true,
};
output.content.push(block);
retainStreamedCall(state, block, update.message.value.callId);
stream.push({ type: "toolcall_start", contentIndex: output.content.length - 1, partial: output });
}
}
} else if (updateCase === "toolCallDelta" || updateCase === "partialToolCall") {
@@ -3985,6 +4228,19 @@ export function processInteractionUpdate(
isError,
timestamp: Date.now(),
});
} else if (settled[kStreamingBlockKind] === "web-fetch") {
const fetchCall = selectHostedFetchCall(toolCall);
const url = fetchCall?.args?.url || extractHttpUrlFromUnknown(toolCall ?? {});
if (url) settled.arguments = { url };
const { text, isError } = describeHostedFetchResult(fetchCall);
state.onToolResult?.({
role: "toolResult",
toolCallId: settled.id,
toolName: "web_fetch",
content: [{ type: "text", text }],
isError,
timestamp: Date.now(),
});
} else if (settled[kStreamingBlockKind] === "todo") {
// Only the server's success snapshot is authoritative: the request args
// may differ from what was actually stored after a merge, and on
@@ -426,11 +426,13 @@ message ToolCall {
PiLsToolCall pi_ls_tool_call = 67;
ConnectScmToolCall connect_scm_tool_call = 68;
SearchConversationsToolCall search_conversations_tool_call = 69;
FetchToolCall web_fetch_tool_call = 37;
}
// Modern builds carry the call id on the envelope instead of inside each
// variant's args. Fields 37-53, 55, 56, 58 (further tool variants) and 54
// (hook_additional_contexts), 59/60 (started_at_ms/completed_at_ms) are not
// modelled and decode into unknown fields; do not reuse those numbers.
// variant's args. Field 37 is the hosted WebFetch tool. Fields 38-53, 55,
// 56, 58 (further tool variants) and 54 (hook_additional_contexts),
// 59/60 (started_at_ms/completed_at_ms) are not modelled and decode into
// unknown fields; do not reuse those numbers.
optional string tool_call_id = 57;
}
@@ -883,6 +885,7 @@ message InteractionQuery {
ExaFetchRequestQuery exa_fetch_request_query = 6;
CreatePlanRequestQuery create_plan_request_query = 7;
SetupVmEnvironmentArgs setup_vm_environment_args = 8;
WebFetchRequestQuery web_fetch_request_query = 9;
}
}
@@ -896,9 +899,28 @@ message InteractionResponse {
ExaFetchRequestResponse exa_fetch_request_response = 6;
CreatePlanRequestResponse create_plan_request_response = 7;
SetupVmEnvironmentResult setup_vm_environment_result = 8;
WebFetchRequestResponse web_fetch_request_response = 9;
}
}
message WebFetchRequestQuery {
FetchArgs args = 1;
}
message WebFetchRequestResponse {
oneof result {
WebFetchRequestResponse_Approved approved = 1;
WebFetchRequestResponse_Rejected rejected = 2;
}
}
message WebFetchRequestResponse_Approved {
}
message WebFetchRequestResponse_Rejected {
string reason = 1;
}
message AskQuestionInteractionQuery {
AskQuestionArgs args = 1;
string tool_call_id = 2;
@@ -0,0 +1,262 @@
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 { InteractionQuery, InteractionResponse } from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb";
import {
type AgentClientMessage,
AgentClientMessageSchema,
AgentServerMessageSchema,
AskQuestionInteractionQuerySchema,
CreatePlanRequestQuerySchema,
ExaFetchArgsSchema,
ExaFetchRequestQuerySchema,
ExaSearchArgsSchema,
ExaSearchRequestQuerySchema,
FetchArgsSchema,
InteractionQuerySchema,
SwitchModeRequestQuerySchema,
WebFetchRequestQuerySchema,
WebSearchArgsSchema,
WebSearchRequestQuerySchema,
} from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb";
function cursorAssistantMessage(): AssistantMessage {
return {
role: "assistant",
content: [],
api: "cursor-agent",
provider: "cursor",
model: "cursor-grok-4.6-xhigh-fast",
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: 0,
};
}
function newBlockState(): 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 decodeClientFrame(frame: Buffer): AgentClientMessage {
const length = frame.readUInt32BE(1);
return fromBinary(AgentClientMessageSchema, frame.subarray(5, 5 + length));
}
function expectInteractionResponse(frames: AgentClientMessage[]): InteractionResponse {
expect(frames).toHaveLength(1);
const frame = frames[0];
if (frame?.message.case !== "interactionResponse") {
throw new Error("expected an interactionResponse frame");
}
return frame.message.value;
}
async function dispatchQuery(query: InteractionQuery): Promise<AgentClientMessage[]> {
const written: Buffer[] = [];
const h2Request = {
write: (chunk: Buffer) => {
written.push(chunk);
return true;
},
} as unknown as Parameters<typeof handleServerMessage>[5];
await handleServerMessage(
create(AgentServerMessageSchema, {
message: {
case: "interactionQuery",
value: query,
},
}),
cursorAssistantMessage(),
new AssistantMessageEventStream(),
newBlockState(),
new Map(),
h2Request,
undefined,
undefined,
{ sawTokenDelta: false },
[],
);
return written.map(decodeClientFrame);
}
describe("Cursor interaction queries", () => {
it("approves hosted web search so the turn is not left waiting on permission", async () => {
const response = expectInteractionResponse(
await dispatchQuery(
create(InteractionQuerySchema, {
id: 11,
query: {
case: "webSearchRequestQuery",
value: create(WebSearchRequestQuerySchema, {
args: create(WebSearchArgsSchema, { searchTerm: "Grok Bot use cases", toolCallId: "ws-1" }),
}),
},
}),
),
);
expect(response.id).toBe(11);
expect(response.result.case).toBe("webSearchRequestResponse");
if (response.result.case !== "webSearchRequestResponse") return;
expect(response.result.value.result.case).toBe("approved");
});
it("approves Exa fetch, the permission prompt that stalled cursor-grok-4.6-xhigh", async () => {
const response = expectInteractionResponse(
await dispatchQuery(
create(InteractionQuerySchema, {
id: 12,
query: {
case: "exaFetchRequestQuery",
value: create(ExaFetchRequestQuerySchema, {
args: create(ExaFetchArgsSchema, {
ids: ["https://docs.x.ai/grok-bot/use-cases"],
toolCallId: "fetch-1",
}),
}),
},
}),
),
);
expect(response.id).toBe(12);
expect(response.result.case).toBe("exaFetchRequestResponse");
if (response.result.case !== "exaFetchRequestResponse") return;
expect(response.result.value.result.case).toBe("approved");
});
it("approves Exa search", async () => {
const response = expectInteractionResponse(
await dispatchQuery(
create(InteractionQuerySchema, {
id: 13,
query: {
case: "exaSearchRequestQuery",
value: create(ExaSearchRequestQuerySchema, {
args: create(ExaSearchArgsSchema, {
query: "Grok Bot",
type: "auto",
numResults: 5,
toolCallId: "es-1",
}),
}),
},
}),
),
);
expect(response.result.case).toBe("exaSearchRequestResponse");
if (response.result.case !== "exaSearchRequestResponse") return;
expect(response.result.value.result.case).toBe("approved");
});
it("rejects ask-question instead of leaving the query unanswered", async () => {
const response = expectInteractionResponse(
await dispatchQuery(
create(InteractionQuerySchema, {
id: 14,
query: {
case: "askQuestionInteractionQuery",
value: create(AskQuestionInteractionQuerySchema, { toolCallId: "ask-1" }),
},
}),
),
);
expect(response.result.case).toBe("askQuestionInteractionResponse");
if (response.result.case !== "askQuestionInteractionResponse") return;
expect(response.result.value.result?.result.case).toBe("rejected");
});
it("rejects mode switches", async () => {
const response = expectInteractionResponse(
await dispatchQuery(
create(InteractionQuerySchema, {
id: 15,
query: {
case: "switchModeRequestQuery",
value: create(SwitchModeRequestQuerySchema, {}),
},
}),
),
);
expect(response.result.case).toBe("switchModeRequestResponse");
if (response.result.case !== "switchModeRequestResponse") return;
expect(response.result.value.result.case).toBe("rejected");
});
it("errors create-plan so the server is not blocked on a plan URI", async () => {
const response = expectInteractionResponse(
await dispatchQuery(
create(InteractionQuerySchema, {
id: 16,
query: {
case: "createPlanRequestQuery",
value: create(CreatePlanRequestQuerySchema, { toolCallId: "plan-1" }),
},
}),
),
);
expect(response.result.case).toBe("createPlanRequestResponse");
if (response.result.case !== "createPlanRequestResponse") return;
expect(response.result.value.result?.result.case).toBe("error");
});
it("approves hosted WebFetch permission queries", async () => {
const response = expectInteractionResponse(
await dispatchQuery(
create(InteractionQuerySchema, {
id: 18,
query: {
case: "webFetchRequestQuery",
value: create(WebFetchRequestQuerySchema, {
args: create(FetchArgsSchema, { url: "https://example.com", toolCallId: "fetch-2" }),
}),
},
}),
),
);
expect(response.id).toBe(18);
expect(response.result.case).toBe("webFetchRequestResponse");
if (response.result.case !== "webFetchRequestResponse") return;
expect(response.result.value.result.case).toBe("approved");
});
it("does not invent a reply for an unknown query variant", async () => {
const frames = await dispatchQuery(create(InteractionQuerySchema, { id: 17 }));
expect(frames).toHaveLength(0);
});
});