From c0f1f69e16d0ff0bf83739eae5afca05c46fa931 Mon Sep 17 00:00:00 2001 From: Diogo Soares Rodrigues Date: Sat, 25 Jul 2026 20:20:10 -0300 Subject: [PATCH] fix(cursor): await in-flight dispatches before emitting done The h2 data loop parses every frame in a chunk synchronously and starts each handleServerMessage with `void`, so the socket keeps draining while a handler runs. Nothing tracked those promises: after h2Completion the provider went straight to `done`. When an exec request, turnEnded and the stream close arrive in ONE chunk - routine, since the server has no reason to split them - the transport completes while the exec handler (and any onToolResult transformer) is still resolving. The Agent drains its Cursor result buffer on the terminal event, so the result is reserved after the snapshot and never persisted, leaving the synthesized (already resolved) toolCall block unpaired and stripped on replay. Dispatches are tracked in a Set and awaited after h2Completion. Each already swallows its own rejection, so this only waits. The test drives a real h2 server whose final chunk carries readArgs + turnEnded, and releases the handler only once the server flushed its response and the handler is known to be running - no wall-clock delay. Verified it fails 5/5 without the barrier and passes 5/5 with it. A microtask drain in place of the yield does NOT discriminate: the client end handler is IO, not a microtask. --- packages/ai/CHANGELOG.md | 1 + packages/ai/src/providers/cursor.ts | 20 +++- .../ai/test/cursor-terminal-error.test.ts | 104 +++++++++++++++++- 3 files changed, 123 insertions(+), 2 deletions(-) diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 4b23c1afe..3b8b44159 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -16,6 +16,7 @@ - Fixed local Cursor exec calls (`read`/`write`/`grep`/`delete`/`bash`/`lsp`/MCP) vanishing from rebuilt transcripts when the tool produced no result. The assistant block is synthesized and marked server-resolved before the handler runs, so the three result-less paths — no handler installed, a handler returning nothing, and a thrown handler — left the call unpaired. Each now pairs a result carrying the same text the server receives. - Fixed Cursor MCP tool calls being unrecognized on the wire. `ToolCall.tool` is a protobuf oneof, so a decoded message exposes the variant as `{ case, value }` and never as a flat `mcpToolCall` property — the same trap that made native todo calls invisible while hand-shaped fixtures kept passing. Both the streamed start and the completion arg merge now go through a shared selector. - Fixed a streamed Cursor MCP block being named from `name` while its paired result used `toolName`, so the two disagreed whenever the server sent different values. Both now prefer `toolName`. +- Fixed the Cursor stream emitting `done` while a tool handler decoded from the final chunk was still running. Server messages are dispatched fire-and-forget so the socket keeps draining, but nothing waited for them: when an exec request, `turnEnded` and the stream close arrived in one chunk, the turn finished before the handler produced its result, and the result missed the buffer drain that pairs it with its call. In-flight dispatches are now awaited after the transport completes. - Fixed a server-resolved Cursor todo call leaving its transcript block stuck pending: the synthetic completion was emitted under a freshly generated id instead of the streamed call id the interactive transcript filed the block under, so the card animated indefinitely. The settled call id is now handed to the sync handler. - Fixed server-resolved Cursor todo blocks disappearing from rebuilt transcripts: nothing produced a `toolResult` for them, and `buildSessionContext` strips any `toolCall` left unpaired, so the interaction vanished on reload, branch switch, or transcript rebuild. The result the host builds is now persisted verbatim — it carries the `details.phases` the todo renderer rebuilds the list from, which a summary-only result would have replayed as `0 tasks`. - Fixed a refused or failed Cursor todo call leaving its card animating forever. Only a successful snapshot settled the block, so a `read_todos` narrowed by a filter and a server `UpdateTodosError` both went unanswered — no `tool_execution_end`, and no `toolResult` to keep the block from being stripped on rebuild. Every completed native todo call now settles. A server error is carried through as a failed result rather than collapsed into the benign "nothing to mirror" case, which would have replayed the failure as a success. diff --git a/packages/ai/src/providers/cursor.ts b/packages/ai/src/providers/cursor.ts index 0abb83d59..8d4685e1a 100644 --- a/packages/ai/src/providers/cursor.ts +++ b/packages/ai/src/providers/cursor.ts @@ -475,6 +475,7 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( let currentToolCall: ToolCallState | null = null; const resolvedMcpToolCallIds = new Set(); const usageState: UsageState = { sawTokenDelta: false }; + const inFlightDispatches = new Set>(); const state: BlockState = { get currentTextBlock() { @@ -547,7 +548,13 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( const isTurnEnded = serverMessage.message.case === "interactionUpdate" && serverMessage.message.value.message?.case === "turnEnded"; - void handleServerMessage( + // Dispatch is fire-and-forget so the socket keeps draining while a + // handler runs, but the promise is tracked: `done` must not be + // pushed while an exec handler is still resolving, or the Agent + // drains its Cursor result buffer before the handler reserved its + // entry and the call is left unpaired. Awaited after + // `h2Completion` below. + const dispatch = handleServerMessage( serverMessage, output, stream, @@ -562,6 +569,8 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( ).catch(error => { log("error", "handleServerMessage", { error: String(error) }); }); + inFlightDispatches.add(dispatch); + void dispatch.finally(() => inFlightDispatches.delete(dispatch)); // Application completion is not protocol success; wait for a clean HTTP/2 end. if (isTurnEnded) { @@ -623,6 +632,15 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( h2Request.write(frameConnectMessage(requestBytes)); heartbeatTimer = setInterval(sendHeartbeat, 5000); await h2Completion.promise; + // The transport is done, but a handler decoded from the last chunk may + // still be running: exec handlers and `onToolResult` transformers are + // async. Pushing `done` now would let the Agent drain its Cursor result + // buffer before such a handler reserves its entry, leaving the call + // unpaired and stripped from every rebuilt transcript. Each dispatch + // already swallows its own rejection, so this only waits. + while (inFlightDispatches.size > 0) { + await Promise.all([...inFlightDispatches]); + } endCurrentTextBlock(output, stream, state); endCurrentThinkingBlock(output, stream, state); diff --git a/packages/ai/test/cursor-terminal-error.test.ts b/packages/ai/test/cursor-terminal-error.test.ts index b25abdec6..7f15b8902 100644 --- a/packages/ai/test/cursor-terminal-error.test.ts +++ b/packages/ai/test/cursor-terminal-error.test.ts @@ -6,7 +6,9 @@ import type { Context, Model } from "@oh-my-pi/pi-ai/types"; import { buildModel } from "@oh-my-pi/pi-catalog/build"; import { AgentServerMessageSchema, + ExecServerMessageSchema, InteractionUpdateSchema, + ReadArgsSchema, TextDeltaUpdateSchema, TurnEndedUpdateSchema, } from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb"; @@ -18,7 +20,8 @@ type Scenario = | { kind: "connect-error-after-turn" } | { kind: "grpc-trailer-after-turn" } | { kind: "end-before-turn" } - | { kind: "hang-after-turn" }; + | { kind: "hang-after-turn" } + | { kind: "exec-in-final-chunk"; responseFinished: PromiseWithResolvers }; let server: http2.Http2Server | undefined; const sessions = new Set(); @@ -67,6 +70,31 @@ function connectEndErrorFrame(code: string, message: string): Buffer { return frameConnectMessage(payload, CONNECT_END_STREAM_FLAG); } +/** + * A `read` exec request, `turnEnded` and the stream close, all in ONE chunk. + * + * The provider parses every frame in a chunk synchronously and dispatches each + * `handleServerMessage` fire-and-forget, so the exec handler is still running + * when the transport completes. Without a barrier before `done`, the Agent + * drains its Cursor result buffer first and the call is never paired. + */ +function execAndTurnEndedFrame(): Buffer { + const message = create(AgentServerMessageSchema, { + message: { + case: "execServerMessage", + value: create(ExecServerMessageSchema, { + id: 1, + execId: "exec-final", + message: { + case: "readArgs", + value: create(ReadArgsSchema, { path: "/tmp/final", toolCallId: "call-final" }), + }, + }), + }, + }); + return Buffer.concat([frameConnectMessage(toBinary(AgentServerMessageSchema, message)), turnEndedFrame()]); +} + async function startServer(): Promise { server = http2.createServer(); server.on("session", session => { @@ -113,6 +141,16 @@ async function startServer(): Promise { return; } + if (scenario.kind === "exec-in-final-chunk") { + const { responseFinished } = scenario; + // Resolves once the server has flushed the whole response, so the test + // never guesses at timing. + stream.on("finish", () => responseFinished.resolve()); + stream.write(execAndTurnEndedFrame()); + stream.end(); + return; + } + stream.write(Buffer.concat([textDeltaFrame("hello"), turnEndedFrame()])); if (scenario.kind === "connect-error-after-turn") { @@ -254,4 +292,68 @@ describe("Cursor terminal lifecycle after turnEnded", () => { expect(eventTypes).not.toContain("done"); expect(result.stopReason).toBe("aborted"); }); + + it("waits for an exec handler decoded from the final chunk before done", async () => { + // The provider dispatches every decoded message fire-and-forget so the + // socket keeps draining. When the exec request, `turnEnded` and the close + // arrive in ONE chunk, the transport completes while the handler is still + // running. `done` must not be pushed first: the Agent drains its Cursor + // result buffer on the terminal event, so a result reserved afterwards + // misses the drain and the synthesized (already resolved) toolCall block + // is stripped from every rebuilt transcript as dangling. + // + // No wall-clock delay. The handler is released only after the server has + // flushed its whole response AND the handler is known to be running, so + // the transport has genuinely completed while the handler is in flight. + const responseFinished = Promise.withResolvers(); + scenario = { kind: "exec-in-final-chunk", responseFinished }; + const baseUrl = await startServer(); + const paired: string[] = []; + const handlerStarted = Promise.withResolvers(); + const handlerDone = Promise.withResolvers(); + const stream = streamCursor(makeModel(baseUrl), context, { + apiKey: "test-token", + execHandlers: { + async read() { + handlerStarted.resolve(); + await handlerDone.promise; + return { + role: "toolResult", + toolCallId: "call-final", + toolName: "read", + content: [{ type: "text", text: "file body" }], + isError: false, + timestamp: 1, + }; + }, + }, + onToolResult: result => { + paired.push(result.toolCallId); + return result; + }, + }); + + const gate = (async () => { + await Promise.all([handlerStarted.promise, responseFinished.promise]); + // `finish` means the server flushed its bytes, not that the client has + // processed the end. Yield so the client's `end` handler and every + // queued continuation run first: a provider that does not await the + // handler settles the stream in exactly that window. + await Bun.sleep(0); + expect(stream.resultSettled).toBe(false); + expect(paired).toEqual([]); + handlerDone.resolve(); + })(); + + const eventTypes: string[] = []; + for await (const event of stream) { + // The result must already be paired by the time `done` is observed. + if (event.type === "done") expect(paired).toEqual(["call-final"]); + eventTypes.push(event.type); + } + await gate; + + expect(eventTypes).toContain("done"); + expect(paired).toEqual(["call-final"]); + }); });