From b26ee4fb35005cd69a97f4e2d0c3f9fde9bcfd1f Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 23 Aug 2026 06:47:23 -0700 Subject: [PATCH] fix(agents): honor cancellation across compaction planning (#128215) --- src/agents/compaction-planning-worker.test.ts | 94 ++++++++++++++++++- src/agents/compaction-planning-worker.ts | 3 + .../compaction.summarize-fallback.test.ts | 10 +- 3 files changed, 98 insertions(+), 9 deletions(-) diff --git a/src/agents/compaction-planning-worker.test.ts b/src/agents/compaction-planning-worker.test.ts index 4ff95bd3857d..c9654f53d8f0 100644 --- a/src/agents/compaction-planning-worker.test.ts +++ b/src/agents/compaction-planning-worker.test.ts @@ -2,10 +2,16 @@ import { MAX_TIMER_TIMEOUT_MS } from "@openclaw/normalization-core/number-coercion"; import { serializeConversation } from "openclaw/plugin-sdk/agent-core"; import { beforeAll, describe, expect, it, vi } from "vitest"; -import { runCompactionPlanningWorker } from "./compaction-planning-worker-runtime.js"; +import * as compactionPlanningWorkerRuntime from "./compaction-planning-worker-runtime.js"; +import { + CompactionPlanningWorkerError, + runCompactionPlanningWorker, +} from "./compaction-planning-worker-runtime.js"; import { buildOversizedFallbackPlanWithWorker, + buildStageSplitPlanWithWorker, buildSummaryChunksWithWorker, + computeAdaptiveChunkRatioWithWorker, } from "./compaction-planning-worker.js"; import { estimateMessagesTokens } from "./compaction-planning.js"; import { runCompactionPlanningWorkerInput } from "./compaction-planning.worker.js"; @@ -26,6 +32,29 @@ function createSyntheticWorkerUrl(source: string): URL { return new URL(`data:text/javascript,${encodeURIComponent(source)}`); } +const cancellablePlanningOperations = [ + { + operation: "summary chunks", + run: (messages: AgentMessage[], signal: AbortSignal) => + buildSummaryChunksWithWorker({ messages, maxChunkTokens: 1_200, signal }), + }, + { + operation: "oversized fallback", + run: (messages: AgentMessage[], signal: AbortSignal) => + buildOversizedFallbackPlanWithWorker({ messages, contextWindow: 1_200, signal }), + }, + { + operation: "stage splitting", + run: (messages: AgentMessage[], signal: AbortSignal) => + buildStageSplitPlanWithWorker({ messages, maxChunkTokens: 1_200, signal }), + }, + { + operation: "adaptive chunk sizing", + run: (messages: AgentMessage[], signal: AbortSignal) => + computeAdaptiveChunkRatioWithWorker({ messages, contextWindow: 1_200, signal }), + }, +]; + describe("compaction planning worker", () => { let packagedSummaryChunks: Awaited>; @@ -59,6 +88,69 @@ describe("compaction planning worker", () => { } }); + it.each( + cancellablePlanningOperations.flatMap(({ operation, run }) => + [63, 64].map((messageCount) => ({ operation, run, messageCount })), + ), + )( + "honors cancellation for $operation with $messageCount messages", + async ({ run, messageCount }) => { + const reason = new Error("operator cancelled compaction"); + const signal = AbortSignal.abort(reason); + const messages = Array.from({ length: messageCount }, (_, index) => + makeMessage(index + 1, "active user request"), + ); + + await expect(run(messages, signal)).rejects.toBe(reason); + }, + ); + + it("does not resume cancelled compaction when its worker becomes unavailable", async () => { + const controller = new AbortController(); + const reason = new Error("operator cancelled compaction"); + const worker = vi + .spyOn(compactionPlanningWorkerRuntime, "runCompactionPlanningWorker") + .mockImplementationOnce(async () => { + controller.abort(reason); + throw new CompactionPlanningWorkerError("worker disappeared", "unavailable"); + }); + + try { + await expect( + buildSummaryChunksWithWorker({ + messages: Array.from({ length: 64 }, (_, index) => makeMessage(index + 1, "request")), + maxChunkTokens: 1_200, + signal: controller.signal, + }), + ).rejects.toBe(reason); + } finally { + worker.mockRestore(); + } + }); + + it("does not restore a worker plan after its compaction has been cancelled", async () => { + const controller = new AbortController(); + const reason = new Error("operator cancelled compaction"); + const worker = vi + .spyOn(compactionPlanningWorkerRuntime, "runCompactionPlanningWorker") + .mockImplementationOnce(async () => { + controller.abort(reason); + return { kind: "summaryChunks", chunkIndexes: [[0]] }; + }); + + try { + await expect( + buildSummaryChunksWithWorker({ + messages: Array.from({ length: 64 }, (_, index) => makeMessage(index + 1, "request")), + maxChunkTokens: 1_200, + signal: controller.signal, + }), + ).rejects.toBe(reason); + } finally { + worker.mockRestore(); + } + }); + it("plans summary chunks in the packaged worker", () => { expect(packagedSummaryChunks.kind).toBe("summaryChunks"); if (packagedSummaryChunks.kind !== "summaryChunks") { diff --git a/src/agents/compaction-planning-worker.ts b/src/agents/compaction-planning-worker.ts index 8aea43bd9ba6..e8d32d352daf 100644 --- a/src/agents/compaction-planning-worker.ts +++ b/src/agents/compaction-planning-worker.ts @@ -48,6 +48,7 @@ async function runCompactionPlan TResult; }): Promise { + params.signal?.throwIfAborted(); const messages = sanitizeCompactionMessages(params.input.messages); if (messages.length < COMPACTION_PLANNING_WORKER_MIN_MESSAGES) { return params.fallback(params.input.messages); @@ -61,6 +62,7 @@ async function runCompactionPlan { expect(agentSessionMocks.generateSummary).toHaveBeenCalledTimes(2); }); - it("does not retry and propagates AbortError immediately when caller signal is already aborted", async () => { + it("does not contact the provider when the caller signal is already aborted", async () => { const controller = new AbortController(); controller.abort(); - const providerAbortErr = Object.assign(new Error("This operation was aborted"), { - name: "AbortError", - }); - agentSessionMocks.generateSummary.mockRejectedValueOnce(providerAbortErr); - await expect( summarizeWithFallback({ messages: [ @@ -164,8 +159,7 @@ describe("summarizeWithFallback", () => { }), ).rejects.toMatchObject({ name: "AbortError" }); - // Caller abort is terminal — no retry, no fallback to placeholder. - expect(agentSessionMocks.generateSummary).toHaveBeenCalledTimes(1); + expect(agentSessionMocks.generateSummary).not.toHaveBeenCalled(); }); it("stops retry backoff promptly when the caller aborts mid-sleep", async () => {