diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 3b8b44159..4f97fec14 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -9,6 +9,8 @@ - Added `AuthStorage.revalidateCredentials()` and the optional `AuthCredentialStore.refreshSnapshot` hook: remote broker stores re-fetch `GET /v1/snapshot` on demand so callers pairing live per-credential data with stored identities (`omp usage`) never render against the up-to-an-hour-stale disk-cached snapshot; local SQLite stores are always current and only reload. ### Fixed +- Cursor no longer discards a local tool result when the transport fails mid-execution. The provider waits for in-flight exec dispatches before pushing `done`, but the error path skipped that wait, so a handler decoded from the last chunk landed its result after the Agent had already finalized the call from the terminal error and cleared its buffer — losing the real outcome of a tool that may already have run side effects. Both exits now drain the same barrier. +- Cursor exec handlers returning the bare-result form no longer record a failed call as successful. When an SDK handler returns only a protocol result (no paired `toolResult`), the synthesized transcript entry was always `"Tool produced no transcript result"` with `isError: false`, even for a `rejected` or `error` result — so Cursor saw a failure while the rebuilt transcript showed success. The synthesized entry now derives its state and message from the result's own oneof variant. - Fixed Cursor models silently failing to maintain the todo list. Cursor resolves its native `update_todos`/`read_todos` tools server-side, but the bridge looked for them under flattened `updateTodosToolCall`/`readTodosToolCall` properties, which a decoded `agent.v1.ToolCall` never has — the variant only arrives through the `tool` oneof — so no native todo call was ever recognized. The synthesized `todo` tool call was also emitted as locally runnable with a `{todos}` payload the local tool's schema rejects, so any update that did surface ended as a validation error and local todo state never followed Cursor's. Todo calls are now read from the oneof, both native todo blocks are marked as already-resolved, and local state is mirrored from the server's confirmed success snapshot (leaving state untouched on `UpdateTodosError`). `TODO_STATUS_CANCELLED` now maps to `abandoned` instead of reverting the task to `pending`. - Hardened Cursor todo mirroring against partial `read_todos` responses: a read narrowed by `status_filter`/`id_filter`, or one returning fewer rows than the server's own `total_count`, is a subset rather than the list, and is no longer treated as authoritative. Previously such a response would have deleted every task it omitted. - Fixed an empty `update_todos` response whose `total_count` is nonzero being mirrored as an authoritative clear, deleting every local task at once. The count-mismatch guard skipped empty responses entirely; only a matching zero count is a genuine clear now. An empty `read_todos` stays refused outright, since proto3 decodes an unset `total_count` as `0` and it cannot be told apart from a filtered read that matched nothing. diff --git a/packages/ai/src/providers/cursor.ts b/packages/ai/src/providers/cursor.ts index 8d4685e1a..bc3b1e53b 100644 --- a/packages/ai/src/providers/cursor.ts +++ b/packages/ai/src/providers/cursor.ts @@ -379,6 +379,20 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( timestamp: Date.now(), }; + // Declared outside the `try` because BOTH exits must drain it: an exec + // handler decoded from the last chunk can still be running when the + // transport fails, and the error path finalizes the synthesized call just + // like the success path does. + const inFlightDispatches = new Set>(); + // 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. + const drainInFlightDispatches = async (): Promise => { + while (inFlightDispatches.size > 0) { + await Promise.all([...inFlightDispatches]); + } + }; + let h2Client: http2.ClientHttp2Session | null = null; let h2Request: http2.ClientHttp2Stream | null = null; let heartbeatTimer: NodeJS.Timeout | null = null; @@ -475,7 +489,6 @@ 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() { @@ -638,9 +651,7 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( // 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]); - } + await drainInFlightDispatches(); endCurrentTextBlock(output, stream, state); endCurrentThinkingBlock(output, stream, state); @@ -667,6 +678,13 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( }); stream.end(); } catch (error) { + // Same reason as the success path: the Agent finalizes the synthesized + // 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. + await drainInFlightDispatches(); const result = await AIError.finalize(error, { api: model.api, signal: options?.signal }); output.stopReason = result.stopReason; output.errorStatus = result.status; @@ -1521,9 +1539,14 @@ export async function resolveExecHandler( const finalToolResult = await applyToolResultHandler(toolResult, onToolResult); if (execResult) { + // TResult-only is a supported return form, so the transcript entry has to + // be synthesized here. Deriving its state from the raw result keeps the + // two views consistent: every exec result is a proto oneof whose only + // non-failure variant is `success`, so a `rejected`/`error`/ + // `file_not_found`/... result must not be recorded as a successful call. return { execResult, - toolResult: finalToolResult ?? (await pair("Tool produced no transcript result", false)), + toolResult: finalToolResult ?? (await pair(...describeExecResult(execResult))), }; } if (finalToolResult) { @@ -1537,6 +1560,25 @@ export async function resolveExecHandler( } } +/** + * Derive the transcript state of an exec result the SDK handler returned in the + * TResult-only form, which carries no `toolResult` to copy it from. + * + * Every exec result in `agent.proto` is a `oneof result` whose success variant + * is named `success` — the rest (`error`, `rejected`, `file_not_found`, + * `permission_denied`, `invalid_file`, ...) are failures. Recording those as a + * successful call would show the user a green entry for a call Cursor was told + * failed. The variant's own `error`/`reason` text is the same string the server + * receives, so it is reused verbatim as the transcript body. + */ +function describeExecResult(execResult: unknown): [text: string, isError: boolean] { + const result = (execResult as { result?: { case?: string; value?: unknown } } | null)?.result; + const variant = result?.case; + if (!variant || variant === "success") return ["Tool produced no transcript result", false]; + const value = result?.value as { error?: string; reason?: string } | undefined; + return [value?.error || value?.reason || `Tool call ${variant}`, true]; +} + function splitExecHandlerResult(result: CursorExecHandlerResult): { execResult?: TResult; toolResult?: ToolResultMessage; diff --git a/packages/ai/test/cursor-exec-handlers.test.ts b/packages/ai/test/cursor-exec-handlers.test.ts index da4b8a286..cc515f866 100644 --- a/packages/ai/test/cursor-exec-handlers.test.ts +++ b/packages/ai/test/cursor-exec-handlers.test.ts @@ -16,12 +16,17 @@ import type { AssistantMessage, Context, CursorExecHandlers, Model, ToolResultMe import { kCursorExecResolved } from "@oh-my-pi/pi-ai/utils/block-symbols"; import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream"; import { buildModel } from "@oh-my-pi/pi-catalog/build"; +import type { ReadResult } from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb"; import { type AgentRunRequest, AgentServerMessageSchema, ExecServerMessageSchema, McpArgsSchema, ReadArgsSchema, + ReadErrorSchema, + ReadRejectedSchema, + ReadResultSchema, + ReadSuccessSchema, } from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb"; const cursorModel: Model<"cursor-agent"> = buildModel({ @@ -250,6 +255,78 @@ describe("Cursor resolveExecHandler execHandlers binding", () => { expect(toolResult).toMatchObject({ toolCallId: "exec-1", isError: false }); }); + it("records a rejected TResult-only return as a failed call", async () => { + // TResult-only is a supported handler form, so the transcript entry has + // to be synthesized. A `rejected` result means Cursor was told the call + // failed - recording it as successful hides that from the user and from + // downstream lifecycle logic. + const rejected = create(ReadResultSchema, { + result: { case: "rejected", value: create(ReadRejectedSchema, { path: "/tmp/foo", reason: "denied" }) }, + }); + // Explicit TResult: `ReadResult` has its own `result` field, so inference + // would otherwise match the `{ result?: TResult }` handler-return variant + // and unwrap the oneof as the exec result. + const { execResult, toolResult } = await resolveExecHandler<{ path: string }, ReadResult>( + { path: "/tmp/foo" }, + async () => rejected, + undefined, + () => rejected, + () => rejected, + () => rejected, + pairing, + ); + + expect(execResult).toBe(rejected); + // The variant's own text is what the server received, so reuse it. + expect(toolResult).toMatchObject({ + toolCallId: "exec-1", + content: [{ type: "text", text: "denied" }], + isError: true, + }); + }); + + it("records an errored TResult-only return as a failed call", async () => { + const errored = create(ReadResultSchema, { + result: { case: "error", value: create(ReadErrorSchema, { path: "/tmp/foo", error: "EIO" }) }, + }); + const { toolResult } = await resolveExecHandler<{ path: string }, ReadResult>( + { path: "/tmp/foo" }, + async () => errored, + undefined, + () => errored, + () => errored, + () => errored, + pairing, + ); + + expect(toolResult).toMatchObject({ content: [{ type: "text", text: "EIO" }], isError: true }); + }); + + it("keeps a successful TResult-only return successful", async () => { + // `success` is the only non-failure variant; the placeholder text still + // applies because the handler gave the transcript nothing to show. + const ok = create(ReadResultSchema, { + result: { + case: "success", + value: create(ReadSuccessSchema, { path: "/tmp/foo", output: { case: "content", value: "hi" } }), + }, + }); + const { toolResult } = await resolveExecHandler<{ path: string }, ReadResult>( + { path: "/tmp/foo" }, + async () => ok, + undefined, + () => ok, + () => ok, + () => ok, + pairing, + ); + + expect(toolResult).toMatchObject({ + content: [{ type: "text", text: "Tool produced no transcript result" }], + isError: false, + }); + }); + it("routes a synthesized result through onToolResult, like a real one", async () => { const seen: string[] = []; const { toolResult } = await resolveExecHandler( diff --git a/packages/ai/test/cursor-terminal-error.test.ts b/packages/ai/test/cursor-terminal-error.test.ts index b1fcc93ee..4d4fc64f7 100644 --- a/packages/ai/test/cursor-terminal-error.test.ts +++ b/packages/ai/test/cursor-terminal-error.test.ts @@ -21,7 +21,8 @@ type Scenario = | { kind: "grpc-trailer-after-turn" } | { kind: "end-before-turn" } | { kind: "hang-after-turn" } - | { kind: "exec-in-final-chunk"; responseFinished: PromiseWithResolvers }; + | { kind: "exec-in-final-chunk"; responseFinished: PromiseWithResolvers } + | { kind: "exec-then-transport-error"; responseFinished: PromiseWithResolvers }; let server: http2.Http2Server | undefined; const sessions = new Set(); @@ -71,14 +72,12 @@ function connectEndErrorFrame(code: string, message: string): Buffer { } /** - * 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. + * A `read` exec request. The provider parses every frame in a chunk + * synchronously and dispatches each `handleServerMessage` fire-and-forget, so + * pairing this with a terminal frame in ONE chunk leaves the exec handler + * running while the transport settles. */ -function execAndTurnEndedFrame(): Buffer { +function execRequestFrame(): Buffer { const message = create(AgentServerMessageSchema, { message: { case: "execServerMessage", @@ -92,7 +91,16 @@ function execAndTurnEndedFrame(): Buffer { }), }, }); - return Buffer.concat([frameConnectMessage(toBinary(AgentServerMessageSchema, message)), turnEndedFrame()]); + return frameConnectMessage(toBinary(AgentServerMessageSchema, message)); +} + +/** + * Exec request + `turnEnded` in one chunk: the clean-completion race. Without a + * barrier before `done`, the Agent drains its Cursor result buffer first and + * the call is never paired. + */ +function execAndTurnEndedFrame(): Buffer { + return Buffer.concat([execRequestFrame(), turnEndedFrame()]); } async function startServer(): Promise { @@ -151,6 +159,20 @@ async function startServer(): Promise { return; } + if (scenario.kind === "exec-then-transport-error") { + const { responseFinished } = scenario; + stream.on("finish", () => responseFinished.resolve()); + // The exec request and the failure land in ONE chunk: the handler is + // dispatched fire-and-forget and is still running when the transport + // rejects. `turnEnded` is deliberately absent — this is the turn dying, + // not ending. + stream.write( + Buffer.concat([execRequestFrame(), connectEndErrorFrame("unavailable", "mid-exec transport failure")]), + ); + stream.end(); + return; + } + stream.write(Buffer.concat([textDeltaFrame("hello"), turnEndedFrame()])); if (scenario.kind === "connect-error-after-turn") { @@ -362,4 +384,66 @@ describe("Cursor terminal lifecycle after turnEnded", () => { expect(eventTypes).toContain("done"); expect(paired).toEqual(["call-final"]); }); + + it("waits for an in-flight exec handler before emitting the transport error", async () => { + // Same race as above, but the turn DIES instead of ending: the exec request + // and the transport failure arrive in one chunk. The Agent finalizes the + // synthesized call from the terminal error and clears its Cursor result + // buffer, so a handler still running would land its real result after + // `agent_end` and have it discarded — even though the tool may already + // have performed side effects. The error must not be pushed first. + const responseFinished = Promise.withResolvers(); + scenario = { kind: "exec-then-transport-error", 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]); + await Bun.sleep(0); + try { + expect(stream.resultSettled).toBe(false); + expect(paired).toEqual([]); + } finally { + handlerDone.resolve(); + } + })(); + + const eventTypes: string[] = []; + for await (const event of stream) { + // The handler's result must already exist by the time the terminal + // error is observed — that is the event the Agent drains on. + if (event.type === "error") expect(paired).toEqual(["call-final"]); + eventTypes.push(event.type); + } + await gate; + const result = await stream.result(); + + expect(eventTypes.at(-1)).toBe("error"); + expect(eventTypes).not.toContain("done"); + expect(result.errorMessage).toContain("mid-exec transport failure"); + expect(paired).toEqual(["call-final"]); + }); });