diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index e4fd848ab6e1..585cff588ed6 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -35,6 +35,7 @@ const emitGrokBackgroundTaskStarted = process.env.T3_ACP_EMIT_GROK_BACKGROUND_TA const emitForeignSessionUpdates = process.env.T3_ACP_EMIT_FOREIGN_SESSION_UPDATES === "1"; const waitForResumeRelease = process.env.T3_ACP_WAIT_FOR_RESUME_RELEASE === "1"; const completeFirstPromptOnCancel = process.env.T3_ACP_COMPLETE_FIRST_PROMPT_ON_CANCEL === "1"; +const failPromptOnCancel = process.env.T3_ACP_FAIL_PROMPT_ON_CANCEL === "1"; const floodStderr = process.env.T3_ACP_FLOOD_STDERR === "1"; const hangPromptForever = process.env.T3_ACP_HANG_PROMPT_FOREVER === "1"; const hangFirstPromptForever = process.env.T3_ACP_HANG_FIRST_PROMPT_FOREVER === "1"; @@ -590,7 +591,7 @@ const program = Effect.gen(function* () { Effect.gen(function* () { const cancelledSessionId = String(sessionId ?? "mock-session-1"); cancelledSessions.add(cancelledSessionId); - if (completeFirstPromptOnCancel) { + if (completeFirstPromptOnCancel || failPromptOnCancel) { yield* Deferred.succeed(nativeCancelRequested, undefined); yield* agent.client.sessionUpdate({ sessionId: cancelledSessionId, @@ -620,7 +621,7 @@ const program = Effect.gen(function* () { const requestedSessionId = String(request.sessionId ?? sessionId); promptCount += 1; - if (completeFirstPromptOnCancel && promptCount === 1) { + if ((completeFirstPromptOnCancel || failPromptOnCancel) && promptCount === 1) { yield* agent.client.sessionUpdate({ sessionId: requestedSessionId, update: { @@ -633,6 +634,12 @@ const program = Effect.gen(function* () { }); yield* Deferred.await(nativeCancelRequested); yield* Deferred.await(nativeCancelRelease); + if (failPromptOnCancel) { + return yield* new AcpError.AcpRequestError({ + code: -32000, + errorMessage: "context canceled: The request was canceled by the client.", + }); + } yield* agent.client.sessionUpdate({ sessionId: requestedSessionId, update: { @@ -1346,7 +1353,10 @@ const program = Effect.gen(function* () { return Deferred.succeed(resumeRelease, undefined).pipe(Effect.as({})); } if (method === "_test/finish-cancel") { - return Deferred.succeed(nativeCancelRelease, undefined).pipe(Effect.as({})); + return Effect.gen(function* () { + yield* Deferred.succeed(nativeCancelRequested, undefined); + yield* Deferred.succeed(nativeCancelRelease, undefined); + }).pipe(Effect.as({})); } if (method === "_test/startup-metadata") { return Effect.gen(function* () { diff --git a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts index 57323e2675e0..a8327cf94e23 100644 --- a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts +++ b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts @@ -31,6 +31,52 @@ const mockRuntimeOptions = { authMethodId: "test", } satisfies AcpSessionRuntime.AcpSessionRuntimeOptions; +describe("isPromptCancellationError", () => { + it("identifies cancel and abort error messages as cancellations", () => { + expect( + AcpSessionRuntime.isPromptCancellationError( + new EffectAcpErrors.AcpRequestError({ + code: -32000, + errorMessage: "context canceled: The request was canceled by the client.", + }), + ), + ).toBe(true); + expect( + AcpSessionRuntime.isPromptCancellationError( + new EffectAcpErrors.AcpRequestError({ + code: -32000, + errorMessage: "The operation was aborted", + }), + ), + ).toBe(true); + }); + + it("does not classify authRequired (-32000) or other non-cancellation errors as cancellation", () => { + expect( + AcpSessionRuntime.isPromptCancellationError( + EffectAcpErrors.AcpRequestError.authRequired("Authentication required"), + ), + ).toBe(false); + expect( + AcpSessionRuntime.isPromptCancellationError( + new EffectAcpErrors.AcpRequestError({ + code: -32000, + errorMessage: "Authentication required", + }), + ), + ).toBe(false); + expect( + AcpSessionRuntime.isPromptCancellationError( + new EffectAcpErrors.AcpRequestError({ + code: -32603, + errorMessage: "Internal server error", + }), + ), + ).toBe(false); + expect(AcpSessionRuntime.isPromptCancellationError(new Error("boom"))).toBe(false); + }); +}); + describe("AcpSessionRuntime", () => { for (const setupMethod of ["session/new", "session/resume"] as const) { it.effect(`buffers root metadata while ${setupMethod} startup is still pending`, () => @@ -221,6 +267,192 @@ describe("AcpSessionRuntime", () => { }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), ); + it.effect( + "drains active prompt successfully when session/cancel notification fails with transport error", + () => + Effect.gen(function* () { + const toolStarted = yield* Deferred.make(); + const cancelFailed = yield* Deferred.make(); + let promptRequests = 0; + const events: Array = []; + const runtime = yield* AcpSessionRuntime.make({ + ...mockRuntimeOptions, + spawn: { + ...mockRuntimeOptions.spawn, + env: { + T3_ACP_COMPLETE_FIRST_PROMPT_ON_CANCEL: "1", + }, + }, + protocolLogging: { + logOutgoing: true, + logger: (event) => { + if ( + event.direction === "outgoing" && + typeof event.payload === "object" && + event.payload !== null && + "_tag" in event.payload && + event.payload._tag === "Notification" && + "tag" in event.payload && + event.payload.tag === "session/cancel" + ) { + return Deferred.succeed(cancelFailed, undefined).pipe( + Effect.andThen( + Effect.fail( + new EffectAcpErrors.AcpTransportError({ + operation: "call-rpc", + method: "session/cancel", + detail: "Broken pipe", + cause: undefined, + }), + ), + ), + ) as unknown as Effect.Effect; + } + return Effect.void; + }, + }, + cancelBehavior: "wait-for-prompt", + requestLogger: (event) => + Effect.sync(() => { + if (event.method === "session/prompt" && event.status === "started") + promptRequests += 1; + }), + }); + yield* runtime.getEvents().pipe( + Stream.runForEach((event) => { + if (event._tag === "EventStreamBarrier") { + return Deferred.succeed(event.acknowledge, undefined); + } + events.push(event); + if (event._tag === "ToolCallUpdated" && event.toolCall.status === "inProgress") { + return Deferred.succeed(toolStarted, undefined); + } + return Effect.void; + }), + Effect.forkChild, + ); + yield* runtime.start(); + const prompt = yield* runtime + .prompt({ + prompt: [{ type: "text", text: "first" }], + }) + .pipe(Effect.forkChild); + yield* Deferred.await(toolStarted); + const cancellation = yield* runtime.cancel.pipe(Effect.forkChild); + yield* Deferred.await(cancelFailed); + const replacement = yield* runtime + .prompt({ + prompt: [{ type: "text", text: "second" }], + }) + .pipe(Effect.forkChild({ startImmediately: true })); + + expect(prompt.pollUnsafe()).toBeUndefined(); + expect(cancellation.pollUnsafe()).toBeUndefined(); + expect(promptRequests).toBe(1); + yield* runtime.request("_test/finish-cancel", {}); + yield* Fiber.join(cancellation); + + expect(yield* Fiber.join(prompt)).toEqual({ + stopReason: "cancelled", + _meta: { nativeCancel: true }, + }); + expect( + events.some( + (event) => + event._tag === "ToolCallUpdated" && + event.toolCall.status === "failed" && + event.toolCall.detail === "Cancelled.", + ), + ).toBe(true); + const cancelledDelta = events.find( + (event) => event._tag === "ContentDelta" && event.text === "Request cancelled.", + ); + expect(cancelledDelta?._tag).toBe("ContentDelta"); + if (cancelledDelta?._tag === "ContentDelta") { + expect( + events.filter( + (event) => + event._tag === "AssistantItemCompleted" && event.itemId === cancelledDelta.itemId, + ), + ).toHaveLength(1); + } + expect(yield* Fiber.join(replacement)).toMatchObject({ stopReason: "end_turn" }); + expect(promptRequests).toBe(2); + }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), + ); + + it.effect( + "drains active prompt successfully when agent cancels by failing in-flight prompt with context canceled", + () => + Effect.gen(function* () { + const toolStarted = yield* Deferred.make(); + let promptRequests = 0; + const events: Array = []; + const runtime = yield* AcpSessionRuntime.make({ + ...mockRuntimeOptions, + spawn: { + ...mockRuntimeOptions.spawn, + env: { + T3_ACP_FAIL_PROMPT_ON_CANCEL: "1", + }, + }, + cancelBehavior: "wait-for-prompt", + requestLogger: (event) => + Effect.sync(() => { + if (event.method === "session/prompt" && event.status === "started") + promptRequests += 1; + }), + }); + yield* runtime.getEvents().pipe( + Stream.runForEach((event) => { + if (event._tag === "EventStreamBarrier") { + return Deferred.succeed(event.acknowledge, undefined); + } + events.push(event); + if (event._tag === "ToolCallUpdated" && event.toolCall.status === "inProgress") { + return Deferred.succeed(toolStarted, undefined); + } + return Effect.void; + }), + Effect.forkChild, + ); + yield* runtime.start(); + const prompt = yield* runtime + .prompt({ + prompt: [{ type: "text", text: "first" }], + }) + .pipe(Effect.forkChild); + yield* Deferred.await(toolStarted); + const cancellation = yield* runtime.cancel.pipe(Effect.forkChild); + const replacement = yield* runtime + .prompt({ + prompt: [{ type: "text", text: "second" }], + }) + .pipe(Effect.forkChild({ startImmediately: true })); + + expect(prompt.pollUnsafe()).toBeUndefined(); + expect(cancellation.pollUnsafe()).toBeUndefined(); + expect(promptRequests).toBe(1); + yield* runtime.request("_test/finish-cancel", {}); + yield* Fiber.join(cancellation); + + expect(yield* Fiber.join(prompt)).toEqual({ + stopReason: "cancelled", + }); + expect( + events.some( + (event) => + event._tag === "ToolCallUpdated" && + event.toolCall.toolCallId === "native-cancel-tool" && + event.toolCall.status === "failed" && + event.toolCall.detail === "Cancelled.", + ), + ).toBe(true); + expect(yield* Fiber.join(replacement)).toMatchObject({ stopReason: "end_turn" }); + expect(promptRequests).toBe(2); + }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), + ); + it.effect("retires a process when native cancellation times out", () => Effect.gen(function* () { const toolStarted = yield* Deferred.make(); diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index b5894192eed9..b87233fa766c 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -317,6 +317,38 @@ interface EnsureActiveAssistantSegmentResult { interface AcpActivePrompt { readonly fiber: Fiber.Fiber; readonly completed: Deferred.Deferred; + readonly cancelled: Ref.Ref; +} + +/** + * Determines whether a given error represents an intentional cancellation or abort + * from the ACP provider or transport, ensuring other -32000 errors (such as authRequired) + * are not falsely classified as cancellation. + */ +export function isPromptCancellationError(error: unknown): boolean { + if ( + typeof error === "object" && + error !== null && + "_tag" in error && + error._tag === "AcpRequestError" + ) { + const reqErr = error as EffectAcpErrors.AcpRequestError; + const msg = (reqErr.errorMessage ?? "").toLowerCase(); + // Require a cancellation/abort keyword so other -32000 errors (e.g. authRequired) are not masked + return msg.includes("cancel") || msg.includes("abort"); + } + return false; +} + +/** + * Checks whether an Effect failure cause represents an interrupt or an + * underlying ACP prompt cancellation error. + */ +function isPromptCancellationCause(cause: Cause.Cause): boolean { + if (Cause.hasInterrupts(cause)) { + return true; + } + return isPromptCancellationError(Cause.squash(cause)); } export const make = ( @@ -919,6 +951,9 @@ export const make = ( const cancel = Effect.gen(function* () { const started = yield* getStartedState; const activePrompt = yield* Ref.get(activePromptRef); + if (Option.isSome(activePrompt)) { + yield* Ref.set(activePrompt.value.cancelled, true); + } if (options.cancelBehavior !== "wait-for-prompt") { if (Option.isSome(activePrompt)) { yield* Fiber.interrupt(activePrompt.value.fiber).pipe(Effect.ignore); @@ -928,7 +963,7 @@ export const make = ( return; } - yield* acp.agent.cancel({ sessionId: started.sessionId }); + yield* acp.agent.cancel({ sessionId: started.sessionId }).pipe(Effect.ignore); if (Option.isNone(activePrompt)) { return; } @@ -951,7 +986,9 @@ export const make = ( return yield* error; } if (Exit.isFailure(completed.value)) { - return yield* Effect.failCause(completed.value.cause); + if (!isPromptCancellationCause(completed.value.cause)) { + return yield* Effect.failCause(completed.value.cause); + } } }); @@ -989,12 +1026,13 @@ export const make = ( ...payload, } satisfies EffectAcpSchema.PromptRequest; const completed = yield* Deferred.make(); + const cancelled = yield* Ref.make(false); const fiber = yield* runLoggedRequest( "session/prompt", requestPayload, acp.agent.prompt(requestPayload), ).pipe(Effect.forkIn(runtimeScope)); - const active = { fiber, completed } satisfies AcpActivePrompt; + const active = { fiber, completed, cancelled } satisfies AcpActivePrompt; yield* Ref.set(activePromptRef, Option.some(active)); if (promptOptions?.dispatched) { yield* Deferred.succeed(promptOptions.dispatched, undefined); @@ -1005,11 +1043,23 @@ export const make = ( (activePrompt) => Fiber.join(activePrompt.fiber).pipe( Effect.catchCause((cause) => - options.cancelBehavior !== "wait-for-prompt" && Cause.hasInterruptsOnly(cause) - ? Effect.succeed({ + Effect.gen(function* () { + const isCancelled = yield* Ref.get(activePrompt.cancelled); + if ( + (options.cancelBehavior !== "wait-for-prompt" && + Cause.hasInterruptsOnly(cause)) || + (isCancelled && isPromptCancellationCause(cause)) + ) { + yield* finalizeActiveToolCallsOnCancellation({ + queue: eventQueue, + toolCallsRef, + }); + return { stopReason: "cancelled", - } satisfies EffectAcpSchema.PromptResponse) - : Effect.failCause(cause), + } satisfies EffectAcpSchema.PromptResponse; + } + return yield* Effect.failCause(cause); + }), ), Effect.tap(() => closeActiveAssistantSegment({ queue: eventQueue, assistantSegmentRef }), @@ -1271,6 +1321,35 @@ const ensureActiveAssistantSegment = ({ ), ); +/** + * Finalizes any active in-flight tool calls with a failed terminal cancellation status + * when an ACP prompt turn is cancelled. + */ +const finalizeActiveToolCallsOnCancellation = ({ + queue, + toolCallsRef, +}: { + readonly queue: Queue.Queue; + readonly toolCallsRef: Ref.Ref>; +}) => + Effect.gen(function* () { + const activeToolCalls = yield* Ref.modify(toolCallsRef, (current) => { + const active = Array.from(current.values()).map((tracked) => tracked.state); + return [active, new Map()] as const; + }); + for (const toolCall of activeToolCalls) { + yield* Queue.offer(queue, { + _tag: "ToolCallUpdated", + toolCall: { + ...toolCall, + status: "failed", + detail: toolCall.detail ?? "Cancelled.", + }, + rawPayload: undefined, + }); + } + }); + const closeActiveAssistantSegment = ({ queue, assistantSegmentRef, diff --git a/apps/server/src/provider/acp/AntigravityAcpSupport.ts b/apps/server/src/provider/acp/AntigravityAcpSupport.ts index f2f370068181..101f10620b5a 100644 --- a/apps/server/src/provider/acp/AntigravityAcpSupport.ts +++ b/apps/server/src/provider/acp/AntigravityAcpSupport.ts @@ -7,6 +7,7 @@ import { type RuntimeMode, } from "@t3tools/contracts"; import * as Crypto from "effect/Crypto"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; @@ -65,6 +66,7 @@ export const makeAntigravityAcpRuntime = Effect.fn("makeAntigravityAcpRuntime")( authMethodId: input.authMethod ?? "oauth-personal", resumeMethod: "resume", cancelBehavior: "wait-for-prompt", + cancelTimeout: Duration.seconds(60), clientCapabilities: { fs: { readTextFile: input.clientFileSystem === true,