From 13d1024519fec8b6b212f5302acc5d7791b999ca Mon Sep 17 00:00:00 2001 From: Will Blanchard Date: Sun, 13 Sep 2026 14:18:39 -0500 Subject: [PATCH 1/7] fix(antigravity): increase cancel timeout and ignore cancel rpc transport errors - In wait-for-prompt cancelBehavior, ignore transport errors on the session/cancel RPC itself (e.g. context canceled) so the runtime continues to await the active prompt completion instead of failing immediately. - Pass an explicit cancelTimeout of 60 seconds in makeAntigravityAcpRuntime to allow long-running tools and subagents sufficient time to settle before process termination. --- apps/server/src/provider/acp/AcpSessionRuntime.ts | 2 +- apps/server/src/provider/acp/AntigravityAcpSupport.ts | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index b5894192eed9..83b8176e6247 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -928,7 +928,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; } 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, From 90956a4679feae027b1eb3a18cdbfc07bfd7f5aa Mon Sep 17 00:00:00 2001 From: Will Blanchard Date: Sun, 13 Sep 2026 14:33:47 -0500 Subject: [PATCH 2/7] test(acp): add coverage for prompt drainage when cancel RPC fails with transport error --- apps/server/scripts/acp-mock-agent.ts | 7 ++ .../provider/acp/AcpJsonRpcConnection.test.ts | 90 +++++++++++++++++++ 2 files changed, 97 insertions(+) diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index e4fd848ab6e1..bbcbd8a665bc 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -51,6 +51,7 @@ const emitStaleXAiPromptCompleteBeforeSecondHang = const emitOverlappingXAiPromptCompleteOutOfOrder = process.env.T3_ACP_EMIT_OVERLAPPING_XAI_PROMPT_COMPLETE_OUT_OF_ORDER === "1"; const failPrompt = process.env.T3_ACP_FAIL_PROMPT === "1"; +const failCancel = process.env.T3_ACP_FAIL_CANCEL === "1"; const failSetConfigOption = process.env.T3_ACP_FAIL_SET_CONFIG_OPTION === "1"; const exitOnSetConfigOption = process.env.T3_ACP_EXIT_ON_SET_CONFIG_OPTION === "1"; const promptResponseText = process.env.T3_ACP_PROMPT_RESPONSE_TEXT; @@ -600,6 +601,12 @@ const program = Effect.gen(function* () { }, }); } + if (failCancel) { + return yield* new AcpError.AcpRequestError({ + code: -32000, + errorMessage: "context canceled", + }); + } if (emitLateUpdateAfterCancel) { yield* Effect.sleep("50 millis"); yield* Effect.sync(() => { diff --git a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts index 57323e2675e0..4a0f00d79c7b 100644 --- a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts +++ b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts @@ -221,6 +221,96 @@ 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 cancelReceived = 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", + T3_ACP_FAIL_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); + } + if (event._tag === "ThoughtDelta" && event.text === "native-cancel-received") { + return Deferred.succeed(cancelReceived, 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(cancelReceived); + 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("retires a process when native cancellation times out", () => Effect.gen(function* () { const toolStarted = yield* Deferred.make(); From 4865ad9210c1c431425bbafc237f226a19a35ecb Mon Sep 17 00:00:00 2001 From: Will Blanchard Date: Sun, 13 Sep 2026 14:45:36 -0500 Subject: [PATCH 3/7] test(acp): inject sender transport failure for cancel notification --- apps/server/scripts/acp-mock-agent.ts | 7 ---- .../provider/acp/AcpJsonRpcConnection.test.ts | 36 +++++++++++++++---- 2 files changed, 30 insertions(+), 13 deletions(-) diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index bbcbd8a665bc..e4fd848ab6e1 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -51,7 +51,6 @@ const emitStaleXAiPromptCompleteBeforeSecondHang = const emitOverlappingXAiPromptCompleteOutOfOrder = process.env.T3_ACP_EMIT_OVERLAPPING_XAI_PROMPT_COMPLETE_OUT_OF_ORDER === "1"; const failPrompt = process.env.T3_ACP_FAIL_PROMPT === "1"; -const failCancel = process.env.T3_ACP_FAIL_CANCEL === "1"; const failSetConfigOption = process.env.T3_ACP_FAIL_SET_CONFIG_OPTION === "1"; const exitOnSetConfigOption = process.env.T3_ACP_EXIT_ON_SET_CONFIG_OPTION === "1"; const promptResponseText = process.env.T3_ACP_PROMPT_RESPONSE_TEXT; @@ -601,12 +600,6 @@ const program = Effect.gen(function* () { }, }); } - if (failCancel) { - return yield* new AcpError.AcpRequestError({ - code: -32000, - errorMessage: "context canceled", - }); - } if (emitLateUpdateAfterCancel) { yield* Effect.sleep("50 millis"); yield* Effect.sync(() => { diff --git a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts index 4a0f00d79c7b..5eba05cca026 100644 --- a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts +++ b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts @@ -226,7 +226,7 @@ describe("AcpSessionRuntime", () => { () => Effect.gen(function* () { const toolStarted = yield* Deferred.make(); - const cancelReceived = yield* Deferred.make(); + const cancelFailed = yield* Deferred.make(); let promptRequests = 0; const events: Array = []; const runtime = yield* AcpSessionRuntime.make({ @@ -235,7 +235,34 @@ describe("AcpSessionRuntime", () => { ...mockRuntimeOptions.spawn, env: { T3_ACP_COMPLETE_FIRST_PROMPT_ON_CANCEL: "1", - T3_ACP_FAIL_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", @@ -254,9 +281,6 @@ describe("AcpSessionRuntime", () => { if (event._tag === "ToolCallUpdated" && event.toolCall.status === "inProgress") { return Deferred.succeed(toolStarted, undefined); } - if (event._tag === "ThoughtDelta" && event.text === "native-cancel-received") { - return Deferred.succeed(cancelReceived, undefined); - } return Effect.void; }), Effect.forkChild, @@ -269,7 +293,7 @@ describe("AcpSessionRuntime", () => { .pipe(Effect.forkChild); yield* Deferred.await(toolStarted); const cancellation = yield* runtime.cancel.pipe(Effect.forkChild); - yield* Deferred.await(cancelReceived); + yield* Deferred.await(cancelFailed); const replacement = yield* runtime .prompt({ prompt: [{ type: "text", text: "second" }], From bafa26c86e4ac727098a18146e5f0f48f22608ee Mon Sep 17 00:00:00 2001 From: Will Blanchard Date: Sun, 13 Sep 2026 14:48:50 -0500 Subject: [PATCH 4/7] test(acp): allow finish-cancel helper to resolve nativeCancelRequested --- apps/server/scripts/acp-mock-agent.ts | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index e4fd848ab6e1..e40adce41389 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -1346,7 +1346,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* () { From 3c761ea5814c07c9dc7e09cbc2208fb1455fb92b Mon Sep 17 00:00:00 2001 From: Will Blanchard Date: Sun, 13 Sep 2026 15:03:56 -0500 Subject: [PATCH 5/7] fix(acp): normalize active prompt cancellation errors during cancel drainage When Antigravity handles cancellation by rejecting the in-flight session/prompt RPC with code -32000 (context canceled), normalize the error as a cancelled prompt completion and suppress it in runtime.cancel so the turn completes cleanly without user-visible errors. Add deterministic test coverage via T3_ACP_FAIL_PROMPT_ON_CANCEL in acp-mock-agent. --- apps/server/scripts/acp-mock-agent.ts | 11 +++- .../provider/acp/AcpJsonRpcConnection.test.ts | 55 +------------------ .../src/provider/acp/AcpSessionRuntime.ts | 48 ++++++++++++++-- 3 files changed, 53 insertions(+), 61 deletions(-) diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index e40adce41389..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: { diff --git a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts index 5eba05cca026..da8696155611 100644 --- a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts +++ b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts @@ -222,11 +222,10 @@ describe("AcpSessionRuntime", () => { ); it.effect( - "drains active prompt successfully when session/cancel notification fails with transport error", + "drains active prompt successfully when agent cancels by failing in-flight prompt with context canceled", () => Effect.gen(function* () { const toolStarted = yield* Deferred.make(); - const cancelFailed = yield* Deferred.make(); let promptRequests = 0; const events: Array = []; const runtime = yield* AcpSessionRuntime.make({ @@ -234,35 +233,7 @@ describe("AcpSessionRuntime", () => { 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; + T3_ACP_FAIL_PROMPT_ON_CANCEL: "1", }, }, cancelBehavior: "wait-for-prompt", @@ -293,7 +264,6 @@ describe("AcpSessionRuntime", () => { .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" }], @@ -308,28 +278,7 @@ describe("AcpSessionRuntime", () => { 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)), diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index 83b8176e6247..f892f9ff1338 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -317,6 +317,28 @@ interface EnsureActiveAssistantSegmentResult { interface AcpActivePrompt { readonly fiber: Fiber.Fiber; readonly completed: Deferred.Deferred; + readonly cancelled: Ref.Ref; +} + +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(); + return reqErr.code === -32000 || msg.includes("cancel") || msg.includes("abort"); + } + return false; +} + +function isPromptCancellationCause(cause: Cause.Cause): boolean { + if (Cause.hasInterrupts(cause)) { + return true; + } + return isPromptCancellationError(Cause.squash(cause)); } export const make = ( @@ -919,6 +941,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); @@ -951,7 +976,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 +1016,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 +1033,19 @@ 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)) + ) { + return { stopReason: "cancelled", - } satisfies EffectAcpSchema.PromptResponse) - : Effect.failCause(cause), + } satisfies EffectAcpSchema.PromptResponse; + } + return yield* Effect.failCause(cause); + }), ), Effect.tap(() => closeActiveAssistantSegment({ queue: eventQueue, assistantSegmentRef }), From 40d459548f46e9280c0788fdf413c5dcc2337d1d Mon Sep 17 00:00:00 2001 From: Will Blanchard Date: Sun, 13 Sep 2026 15:17:34 -0500 Subject: [PATCH 6/7] fix(acp): refine cancellation detection and finalize in-flight tool calls on cancel --- .../provider/acp/AcpJsonRpcConnection.test.ts | 169 ++++++++++++++++++ .../src/provider/acp/AcpSessionRuntime.ts | 43 ++++- 2 files changed, 210 insertions(+), 2 deletions(-) diff --git a/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts b/apps/server/src/provider/acp/AcpJsonRpcConnection.test.ts index da8696155611..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,120 @@ 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", () => @@ -279,6 +439,15 @@ describe("AcpSessionRuntime", () => { 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)), diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index f892f9ff1338..2504d36ee423 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -320,7 +320,12 @@ interface AcpActivePrompt { readonly cancelled: Ref.Ref; } -function isPromptCancellationError(error: unknown): boolean { +/** + * 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 && @@ -329,7 +334,8 @@ function isPromptCancellationError(error: unknown): boolean { ) { const reqErr = error as EffectAcpErrors.AcpRequestError; const msg = (reqErr.errorMessage ?? "").toLowerCase(); - return reqErr.code === -32000 || msg.includes("cancel") || msg.includes("abort"); + // Require a cancellation/abort keyword so other -32000 errors (e.g. authRequired) are not masked + return msg.includes("cancel") || msg.includes("abort"); } return false; } @@ -1040,6 +1046,10 @@ export const make = ( Cause.hasInterruptsOnly(cause)) || (isCancelled && isPromptCancellationCause(cause)) ) { + yield* finalizeActiveToolCallsOnCancellation({ + queue: eventQueue, + toolCallsRef, + }); return { stopReason: "cancelled", } satisfies EffectAcpSchema.PromptResponse; @@ -1307,6 +1317,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, From bf301196c056a6b59d5707006a3c05cfc847727a Mon Sep 17 00:00:00 2001 From: Will Blanchard Date: Sun, 13 Sep 2026 15:24:06 -0500 Subject: [PATCH 7/7] docs(acp): add TSDoc to isPromptCancellationCause --- apps/server/src/provider/acp/AcpSessionRuntime.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index 2504d36ee423..b87233fa766c 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -340,6 +340,10 @@ export function isPromptCancellationError(error: unknown): boolean { 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;