From ea72ce9e9820451bcf6b41deacd00342fef2c4ae Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 19 Aug 2026 06:02:23 -0700 Subject: [PATCH] fix(agents): reject aborted steer before acceptance (#126314) --- .../run/attempt-queue-message.ts | 6 ++- .../run/attempt.queue-message.test.ts | 43 +++++++++++++++++++ 2 files changed, 47 insertions(+), 2 deletions(-) diff --git a/src/agents/embedded-agent-runner/run/attempt-queue-message.ts b/src/agents/embedded-agent-runner/run/attempt-queue-message.ts index a2135e9f71a2..abf8c8535120 100644 --- a/src/agents/embedded-agent-runner/run/attempt-queue-message.ts +++ b/src/agents/embedded-agent-runner/run/attempt-queue-message.ts @@ -318,9 +318,11 @@ async function steerAndWaitForTranscriptCommit( ); function onAbort() { abortRequested = true; - if (accepted) { - rejectAfterCancellation("queued steering message was cancelled before delivery"); + if (!accepted) { + rejectBeforeAcceptance("queued steering message was cancelled before acceptance"); + return; } + rejectAfterCancellation("queued steering message was cancelled before delivery"); } abortSignal?.addEventListener("abort", onAbort, { once: true }); }); diff --git a/src/agents/embedded-agent-runner/run/attempt.queue-message.test.ts b/src/agents/embedded-agent-runner/run/attempt.queue-message.test.ts index 837cb86d0622..58513496d532 100644 --- a/src/agents/embedded-agent-runner/run/attempt.queue-message.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt.queue-message.test.ts @@ -348,6 +348,49 @@ describe("embedded OpenClaw queued steering cancellation", () => { } }); + it("fences an aborted steer before delayed preparation can enqueue it", async () => { + let releasePreparation!: () => void; + const preparation = new Promise((resolve) => { + releasePreparation = resolve; + }); + let preparationStarted!: () => void; + const started = new Promise((resolve) => { + preparationStarted = resolve; + }); + let enqueued = false; + const onQueueAccepted = vi.fn(); + const activeSession: EmbeddedAgentActiveSessionSteerTarget = { + steer: async (_text, _images, _recorder, _media, _imageOrder, _identity, canInject) => { + preparationStarted(); + await preparation; + if (canInject && !canInject()) { + throw new Error("active session is finalizing"); + } + enqueued = true; + }, + subscribe: () => () => {}, + }; + const controller = new AbortController(); + const wait = steerActiveSessionWithOptionalDeliveryWait(activeSession, "delayed steer", { + abortSignal: controller.signal, + deliveryTimeoutMs: 10_000, + onQueueAccepted, + waitForTranscriptCommit: true, + }); + const rejection = expect(wait).rejects.toThrow( + "queued steering message was cancelled before acceptance", + ); + + await started; + controller.abort(); + releasePreparation(); + + await rejection; + expect(enqueued).toBe(false); + expect(onQueueAccepted).toHaveBeenCalledOnce(); + expect(onQueueAccepted).toHaveBeenCalledWith(false); + }); + it("matches identical steering text by stable queue identity", async () => { let emit!: (event: unknown) => void; const first = {