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