Files
oh-my-pi/packages/ai/test/cursor-terminal-error.test.ts
T
Diogo Soares Rodrigues 90b89bf3c4 fix(cursor): drain exec dispatches on transport error and preserve raw failure state
Two P2 findings on the exec-handler path.

1. The success path waits for in-flight exec dispatches before pushing
   done, but the catch path did not. When the HTTP/2 completion rejects
   while a handler decoded from the last chunk is still running, the
   Agent finalizes the synthesized call from the terminal error and
   clears its Cursor result buffer; the handler then lands its real
   result after agent_end and it is discarded at the next run - even
   though the tool may already have performed side effects. The barrier
   is now a shared helper and inFlightDispatches is hoisted out of the
   try so both exits drain it.

2. resolveExecHandler always synthesized "Tool produced no transcript
   result" with isError: false for the TResult-only handler form. A
   rejected or error protocol result therefore reached Cursor as a
   failure while the rebuilt transcript recorded the same call as
   successful. Every exec result is a proto oneof whose only non-failure
   variant is success, so describeExecResult() derives the state from
   the variant and reuses its own error/reason text - the same string
   the server received.
2026-07-26 09:29:13 -03:00

450 lines
14 KiB
TypeScript

import { afterEach, describe, expect, it } from "bun:test";
import * as http2 from "node:http2";
import { create, toBinary } from "@bufbuild/protobuf";
import { streamCursor } from "@oh-my-pi/pi-ai/providers/cursor";
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";
const CONNECT_END_STREAM_FLAG = 0b00000010;
type Scenario =
| { kind: "success" }
| { kind: "connect-error-after-turn" }
| { kind: "grpc-trailer-after-turn" }
| { kind: "end-before-turn" }
| { kind: "hang-after-turn" }
| { kind: "exec-in-final-chunk"; responseFinished: PromiseWithResolvers<void> }
| { kind: "exec-then-transport-error"; responseFinished: PromiseWithResolvers<void> };
let server: http2.Http2Server | undefined;
const sessions = new Set<http2.Http2Session>();
let scenario: Scenario = { kind: "success" };
function frameConnectMessage(data: Uint8Array, flags = 0): Buffer {
const frame = Buffer.alloc(5 + data.length);
frame[0] = flags;
frame.writeUInt32BE(data.length, 1);
frame.set(data, 5);
return frame;
}
function textDeltaFrame(text: string): Buffer {
const message = create(AgentServerMessageSchema, {
message: {
case: "interactionUpdate",
value: create(InteractionUpdateSchema, {
message: {
case: "textDelta",
value: create(TextDeltaUpdateSchema, { text }),
},
}),
},
});
return frameConnectMessage(toBinary(AgentServerMessageSchema, message));
}
function turnEndedFrame(): Buffer {
const message = create(AgentServerMessageSchema, {
message: {
case: "interactionUpdate",
value: create(InteractionUpdateSchema, {
message: {
case: "turnEnded",
value: create(TurnEndedUpdateSchema, {}),
},
}),
},
});
return frameConnectMessage(toBinary(AgentServerMessageSchema, message));
}
function connectEndErrorFrame(code: string, message: string): Buffer {
const payload = Buffer.from(JSON.stringify({ error: { code, message } }), "utf8");
return frameConnectMessage(payload, CONNECT_END_STREAM_FLAG);
}
/**
* 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 execRequestFrame(): 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 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<string> {
server = http2.createServer();
server.on("session", session => {
sessions.add(session);
session.on("close", () => sessions.delete(session));
});
server.on("stream", (stream: http2.ServerHttp2Stream, headers: http2.IncomingHttpHeaders) => {
stream.on("data", () => {});
if (headers[":path"] !== "/agent.v1.AgentService/Run") {
stream.respond({ ":status": 404 });
stream.end();
return;
}
if (scenario.kind === "grpc-trailer-after-turn") {
stream.respond(
{
":status": 200,
"content-type": "application/connect+proto",
},
{ waitForTrailers: true },
);
stream.on("wantTrailers", () => {
stream.sendTrailers({
"grpc-status": "13",
"grpc-message": encodeURIComponent("post-turn trailer failure"),
});
});
stream.write(textDeltaFrame("hello"));
stream.write(turnEndedFrame());
stream.end();
return;
}
stream.respond({
":status": 200,
"content-type": "application/connect+proto",
});
if (scenario.kind === "end-before-turn") {
stream.write(textDeltaFrame("partial"));
stream.end();
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;
}
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") {
stream.write(connectEndErrorFrame("unavailable", "post-turn connect failure"));
stream.end();
return;
}
if (scenario.kind === "hang-after-turn") {
return;
}
stream.end();
});
const listening = Promise.withResolvers<void>();
server.once("error", listening.reject);
server.listen(0, "127.0.0.1", listening.resolve);
await listening.promise;
const address = server.address();
if (!address || typeof address === "string") {
throw new Error("expected http2 fixture server to bind a tcp port");
}
return `http://127.0.0.1:${address.port}`;
}
function makeModel(baseUrl: string): Model<"cursor-agent"> {
return buildModel({
id: "cursor-terminal-fixture",
name: "Cursor terminal fixture",
api: "cursor-agent",
provider: "cursor",
baseUrl,
reasoning: false,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 1,
maxTokens: 1,
});
}
const context: Context = {
messages: [{ role: "user", content: "terminal lifecycle", timestamp: 1 }],
};
async function collectStream(model: Model<"cursor-agent">, options?: { signal?: AbortSignal }) {
const stream = streamCursor(model, context, { apiKey: "test-token", signal: options?.signal });
const eventTypes: string[] = [];
for await (const event of stream) {
eventTypes.push(event.type);
}
const result = await stream.result();
return { eventTypes, result };
}
async function stopServer(): Promise<void> {
for (const session of sessions) {
session.destroy();
}
sessions.clear();
if (!server) return;
const closing = server;
server = undefined;
const closed = Promise.withResolvers<void>();
closing.close(error => {
if (error) {
closed.reject(error);
} else {
closed.resolve();
}
});
await closed.promise;
}
afterEach(async () => {
scenario = { kind: "success" };
await stopServer();
});
describe("Cursor terminal lifecycle after turnEnded", () => {
it("emits done only after turnEnded and a clean protocol end", async () => {
scenario = { kind: "success" };
const baseUrl = await startServer();
const { eventTypes, result } = await collectStream(makeModel(baseUrl));
expect(eventTypes).toEqual(["start", "text_start", "text_delta", "text_end", "done"]);
expect(result.stopReason).toBe("stop");
expect(result.errorMessage).toBeUndefined();
});
it("surfaces CONNECT end-stream errors that arrive after turnEnded", async () => {
scenario = { kind: "connect-error-after-turn" };
const baseUrl = await startServer();
const { eventTypes, result } = await collectStream(makeModel(baseUrl));
expect(eventTypes[0]).toBe("start");
expect(eventTypes.at(-1)).toBe("error");
expect(eventTypes).not.toContain("done");
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toContain("Connect error unavailable: post-turn connect failure");
});
it("surfaces nonzero gRPC trailers that arrive after turnEnded", async () => {
scenario = { kind: "grpc-trailer-after-turn" };
const baseUrl = await startServer();
const { eventTypes, result } = await collectStream(makeModel(baseUrl));
expect(eventTypes[0]).toBe("start");
expect(eventTypes.at(-1)).toBe("error");
expect(eventTypes).not.toContain("done");
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toContain("gRPC error 13: post-turn trailer failure");
});
it("rejects when the stream ends before turnEnded", async () => {
scenario = { kind: "end-before-turn" };
const baseUrl = await startServer();
const { eventTypes, result } = await collectStream(makeModel(baseUrl));
expect(eventTypes[0]).toBe("start");
expect(eventTypes.at(-1)).toBe("error");
expect(eventTypes).not.toContain("done");
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toContain("Cursor stream ended before turnEnded");
});
it("aborts without emitting done when the signal fires", async () => {
scenario = { kind: "hang-after-turn" };
const baseUrl = await startServer();
const controller = new AbortController();
const stream = streamCursor(makeModel(baseUrl), context, {
apiKey: "test-token",
signal: controller.signal,
});
const eventTypes: string[] = [];
for await (const event of stream) {
eventTypes.push(event.type);
if (event.type === "text_delta") controller.abort();
}
const result = await stream.result();
expect(eventTypes[0]).toBe("start");
expect(eventTypes.at(-1)).toBe("error");
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<void>();
scenario = { kind: "exec-in-final-chunk", responseFinished };
const baseUrl = await startServer();
const paired: string[] = [];
const handlerStarted = Promise.withResolvers<void>();
const handlerDone = Promise.withResolvers<void>();
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);
try {
expect(stream.resultSettled).toBe(false);
expect(paired).toEqual([]);
} finally {
// Always release: a failing assertion here must surface as that
// failure, not as a hung `for await` that waits for a handler
// nobody will ever unblock.
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"]);
});
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<void>();
scenario = { kind: "exec-then-transport-error", responseFinished };
const baseUrl = await startServer();
const paired: string[] = [];
const handlerStarted = Promise.withResolvers<void>();
const handlerDone = Promise.withResolvers<void>();
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"]);
});
});