From daf2680298853c3dacbe3b2d5563bba3bf6ddb76 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sun, 26 Jul 2026 15:16:10 +0200 Subject: [PATCH] fix(cursor): bound the dispatch drain by the abort signal and settle todo calls when the host sync callback throws Two defects flagged by review on the final head, both in the new drain/settle machinery: - drainInFlightDispatches looped on Promise.all unconditionally. Exec handlers have no cancellation contract (the bridge calls tool.execute with no signal), so a hung or long-running tool held the terminal error hostage after Ctrl+C. The drain now returns once the signal aborts; late results were discarded after agent_end regardless. - A host todoSync callback that throws (session persistence on disk failure) skipped both the paired result and toolcall_end, stranding the live card and leaving the resolved block to be stripped from rebuilt transcripts. The callback is now caught and the call settles as a failure carrying the thrown message. Both regression tests fail without their fix (the first by hanging the stream). --- packages/ai/src/providers/cursor.ts | 44 +++++++++++-- .../ai/test/cursor-terminal-error.test.ts | 62 ++++++++++++++++++- packages/ai/test/cursor-todo-bridge.test.ts | 24 +++++++ 3 files changed, 124 insertions(+), 6 deletions(-) diff --git a/packages/ai/src/providers/cursor.ts b/packages/ai/src/providers/cursor.ts index 6d8df0211..0e604a5c2 100644 --- a/packages/ai/src/providers/cursor.ts +++ b/packages/ai/src/providers/cursor.ts @@ -387,9 +387,28 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( // A dispatch can spawn another (a handler that decodes a nested frame), so // re-check rather than awaiting one snapshot. Each dispatch already // swallows its own rejection, so this only waits. + // + // The wait is bounded by the abort signal: exec handlers have no + // cancellation contract (the coding-agent bridge invokes `tool.execute` + // with no signal), so a hung or long-running tool would otherwise hold + // the terminal event hostage after the user already gave up on the turn. + // Once aborted, the Agent finalizes from the abort error and discards + // late results regardless, so skipping the rest of the drain loses + // nothing that could still be delivered. + let abortSettled: Promise | undefined; const drainInFlightDispatches = async (): Promise => { + const signal = options?.signal; while (inFlightDispatches.size > 0) { - await Promise.all([...inFlightDispatches]); + if (signal?.aborted) return; + const settled = Promise.all([...inFlightDispatches]); + if (!signal) { + await settled; + continue; + } + abortSettled ??= new Promise(resolve => + signal.addEventListener("abort", () => resolve(), { once: true }), + ); + await Promise.race([settled, abortSettled]); } }; @@ -682,8 +701,9 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( // call from this terminal error and clears its Cursor result buffer, so // a handler still running would land its real result after `agent_end` // and be discarded — even though the tool may already have run side - // effects. Wait for it first. An abort has already closed the transport, - // and every dispatch settles rather than hanging on it. + // effects. Wait for it first; on abort the drain returns immediately + // (handlers have no cancellation contract and must not delay the + // terminal error the user asked for). await drainInFlightDispatches(); const result = await AIError.finalize(error, { api: model.api, signal: options?.signal }); output.stopReason = result.stopReason; @@ -2765,8 +2785,22 @@ export function processInteractionUpdate( // carries the `details.phases` the todo renderer replays the list // from — with the provider's summary standing in when the host has // nothing to add. - const persisted = state.onTodoSnapshot?.(snapshot, state.currentToolCall.id, error) ?? undefined; - state.onToolResult?.(persisted ?? buildTodoToolResult(state.currentToolCall.id, snapshot, error)); + let persisted: ToolResultMessage | undefined; + let hostError: string | null = null; + try { + persisted = state.onTodoSnapshot?.(snapshot, state.currentToolCall.id, error) ?? undefined; + } catch (callbackError) { + // A throwing host callback (e.g. session persistence failing on + // disk error) must not leave the resolved block unpaired: the + // exception would skip both the paired result and `toolcall_end`, + // stranding the live card and stripping the call from every + // rebuilt transcript. Settle it as a failure instead. + hostError = callbackError instanceof Error ? callbackError.message : String(callbackError); + log("error", "onTodoSnapshot", { error: hostError }); + } + state.onToolResult?.( + persisted ?? buildTodoToolResult(state.currentToolCall.id, snapshot, hostError ?? error), + ); } const idx = output.content.indexOf(state.currentToolCall); clearStreamingPartialJson(state.currentToolCall); diff --git a/packages/ai/test/cursor-terminal-error.test.ts b/packages/ai/test/cursor-terminal-error.test.ts index 4d4fc64f7..fc986e809 100644 --- a/packages/ai/test/cursor-terminal-error.test.ts +++ b/packages/ai/test/cursor-terminal-error.test.ts @@ -22,7 +22,8 @@ type Scenario = | { kind: "end-before-turn" } | { kind: "hang-after-turn" } | { kind: "exec-in-final-chunk"; responseFinished: PromiseWithResolvers } - | { kind: "exec-then-transport-error"; responseFinished: PromiseWithResolvers }; + | { kind: "exec-then-transport-error"; responseFinished: PromiseWithResolvers } + | { kind: "exec-then-hang" }; let server: http2.Http2Server | undefined; const sessions = new Set(); @@ -173,6 +174,13 @@ async function startServer(): Promise { return; } + if (scenario.kind === "exec-then-hang") { + // Exec request, then the stream stays open: the only way this turn + // ends is the client aborting. + stream.write(execRequestFrame()); + return; + } + stream.write(Buffer.concat([textDeltaFrame("hello"), turnEndedFrame()])); if (scenario.kind === "connect-error-after-turn") { @@ -446,4 +454,56 @@ describe("Cursor terminal lifecycle after turnEnded", () => { expect(result.errorMessage).toContain("mid-exec transport failure"); expect(paired).toEqual(["call-final"]); }); + + it("does not hold the abort hostage to a hung exec handler", async () => { + // Exec handlers have no cancellation contract — the coding-agent bridge + // invokes `tool.execute` with no signal — so a hung or long-running tool + // cannot be interrupted. Once the user aborts, the drain must not wait + // for it: the Agent finalizes from the abort error and discards late + // results regardless, so waiting only delays the terminal event the + // user asked for. Without the abort-bounded drain this test times out + // with the stream never settling. + scenario = { kind: "exec-then-hang" }; + const baseUrl = await startServer(); + const controller = new AbortController(); + const handlerStarted = Promise.withResolvers(); + const handlerDone = Promise.withResolvers(); + const stream = streamCursor(makeModel(baseUrl), context, { + apiKey: "test-token", + signal: controller.signal, + execHandlers: { + async read() { + handlerStarted.resolve(); + await handlerDone.promise; + return { + role: "toolResult", + toolCallId: "call-final", + toolName: "read", + content: [{ type: "text", text: "late result" }], + isError: false, + timestamp: 1, + }; + }, + }, + }); + + const gate = (async () => { + await handlerStarted.promise; + controller.abort(); + })(); + + const eventTypes: string[] = []; + for await (const event of stream) { + eventTypes.push(event.type); + } + await gate; + const result = await stream.result(); + // Released only AFTER the stream settled: reaching this line at all + // proves the terminal error did not wait for the handler. + handlerDone.resolve(); + + expect(eventTypes.at(-1)).toBe("error"); + expect(eventTypes).not.toContain("done"); + expect(result.stopReason).toBe("aborted"); + }); }); diff --git a/packages/ai/test/cursor-todo-bridge.test.ts b/packages/ai/test/cursor-todo-bridge.test.ts index b9bee53e6..1f184e9d0 100644 --- a/packages/ai/test/cursor-todo-bridge.test.ts +++ b/packages/ai/test/cursor-todo-bridge.test.ts @@ -876,6 +876,30 @@ describe("cursor native todo bridge (wire-encoded protobuf)", () => { }); }); + it("settles the call as a failure when the host sync callback throws", () => { + // The host callback persists to the session branch and can throw + // synchronously (e.g. a disk failure). The exception must not skip the + // paired result and `toolcall_end`: the block is already marked + // resolved, so left unpaired it is stripped from every rebuilt + // transcript and the live card never resolves. + const h = newHarness(); + h.state.onTodoSnapshot = () => { + throw new Error("session persistence failed"); + }; + const toolCall = updateCall(items([["1", "task", 3]]), 1); + for (const kind of ["toolCallStarted", "toolCallCompleted"] as const) { + processInteractionUpdate(wireUpdate(kind, toolCall) as never, h.output, h.stream, h.state, h.usageState); + } + + expect(h.state.currentToolCall).toBeNull(); + expect(h.toolResults).toHaveLength(1); + expect(h.toolResults[0]).toMatchObject({ + toolCallId: todoBlocks(h)[0].id, + isError: true, + content: [{ type: "text", text: "session persistence failed" }], + }); + }); + it("recognizes an MCP call through the wire-encoded oneof, start and completion", () => { // Same wire-shape trap the native todo calls fell into: `ToolCall.tool` is // a protobuf oneof, so a decoded message exposes the variant as