From d563e25cfe1bf53e1fcef70b4b8e6931ead11746 Mon Sep 17 00:00:00 2001 From: metaphorics <152830360+metaphorics@users.noreply.github.com> Date: Mon, 3 Aug 2026 21:39:33 +0900 Subject: [PATCH] fix(task): settle all late cleanup (#7488) --- packages/coding-agent/src/task/executor.ts | 39 +++++-- .../task/executor-async-quiescence.test.ts | 106 +++++++++++++++++- 2 files changed, 134 insertions(+), 11 deletions(-) diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index b6fe3540b..83c238377 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -2142,14 +2142,17 @@ async function finalizeRunResult(args: FinalizeRunArgs): Promise { exitCode = 1; } const wasAborted = - runtimeLimitExceeded || abortedViaYield || (!hasYield && (done.aborted || signal?.aborted || false)); + runtimeLimitExceeded || Boolean(done.aborted) || abortedViaYield || (!hasYield && Boolean(signal?.aborted)); const finalAbortReason = wasAborted ? runtimeLimitExceeded ? monitor.resolveAbortReasonText() - : abortedViaYield - ? yieldAbortReason - : (done.abortReason ?? - (signal?.aborted ? monitor.resolveSignalAbortReason() : monitor.resolveAbortReasonText())) + : done.aborted + ? (done.abortReason ?? monitor.resolveAbortReasonText()) + : abortedViaYield + ? yieldAbortReason + : signal?.aborted + ? monitor.resolveSignalAbortReason() + : monitor.resolveAbortReasonText() : undefined; progress.status = wasAborted ? "aborted" : exitCode === 0 ? "completed" : "failed"; monitor.scheduleProgress(true); @@ -3202,6 +3205,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise[] = []; + let deferredSessionShutdown: Promise | undefined; const deferCleanup = (completion: Promise): void => { lateCleanups.push(completion); exitCode = 1; @@ -3268,11 +3272,32 @@ export async function runSubprocess(options: ExecutorOptions): Promise { + deferredSessionShutdown = completion; + deferCleanup(completion); + }, }); } + if (jobManager) { + if (deferredSessionShutdown) { + const finalReap = Promise.allSettled([deferredSessionShutdown]).then(async () => { + const reap = await jobManager.cancelAndReapOwnerJobs(id, Date.now()); + await reap.completion; + }); + lateCleanups.push(finalReap); + } else { + const reap = await jobManager.cancelAndReapOwnerJobs(id, cleanupDeadlineAt); + if (!reap.settled) { + deferCleanup(reap.completion); + logger.warn("Subagent async job cleanup exceeded its deadline after session shutdown", { + id, + pendingJobIds: reap.pendingJobIds, + }); + } + } + } if (lateCleanups.length > 0) { - const completion = Promise.all(lateCleanups).then(() => {}); + const completion = Promise.allSettled(lateCleanups).then(() => {}); trackLateCleanup(completion, { id, resource: "subagent" }); options.onCleanupDeferred?.(completion); } diff --git a/packages/coding-agent/test/task/executor-async-quiescence.test.ts b/packages/coding-agent/test/task/executor-async-quiescence.test.ts index 29350000a..fbb893e47 100644 --- a/packages/coding-agent/test/task/executor-async-quiescence.test.ts +++ b/packages/coding-agent/test/task/executor-async-quiescence.test.ts @@ -8,6 +8,7 @@ */ import { afterEach, describe, expect, it, vi } from "bun:test"; import type { AssistantMessage } from "@oh-my-pi/pi-ai"; +import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async/job-manager"; import type { LoadExtensionsResult } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/types"; import type { CreateAgentSessionResult } from "@oh-my-pi/pi-coding-agent/sdk"; import * as sdkModule from "@oh-my-pi/pi-coding-agent/sdk"; @@ -18,7 +19,7 @@ import { EventBus } from "@oh-my-pi/pi-coding-agent/utils/event-bus"; const baseAgent: AgentDefinition = { name: "task", description: "test", systemPrompt: "test", source: "bundled" }; -function assistantStopMessage(text: string): AssistantMessage { +function assistantStopMessage(text: string, totalTokens = 0): AssistantMessage { return { role: "assistant", content: [{ type: "text", text }], @@ -27,10 +28,10 @@ function assistantStopMessage(text: string): AssistantMessage { model: "mock", usage: { input: 0, - output: 0, + output: totalTokens, cacheRead: 0, cacheWrite: 0, - totalTokens: 0, + totalTokens, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, }, stopReason: "stop", @@ -44,9 +45,15 @@ interface AsyncQuiescenceHarness { abortCalls: () => number; settleCalls: () => number; emitTerminalYield: (data: unknown) => void; + emitAssistant: (text: string, totalTokens?: number) => void; finishJob: () => void; } +interface AsyncSessionOptions { + abort?: () => Promise; + dispose?: () => Promise; +} + /** * Mock session with the owner-async surface the barrier drives: * `hasPendingAsyncWork` / `getAsyncJobSnapshot` / `settleAsyncWork`. The job @@ -56,6 +63,7 @@ interface AsyncQuiescenceHarness { */ function createAsyncSession( onPrompt: (params: { text: string; promptIndex: number; harness: AsyncQuiescenceHarness }) => void, + options: AsyncSessionOptions = {}, ): AsyncQuiescenceHarness { const listeners: Array<(event: AgentSessionEvent) => void> = []; const state = { messages: [] as AssistantMessage[] }; @@ -103,6 +111,11 @@ function createAsyncSession( state.messages.push(reaction); emit({ type: "message_end", message: reaction } as AgentSessionEvent); }; + const emitAssistant = (text: string, totalTokens = 0) => { + const message = assistantStopMessage(text, totalTokens); + state.messages.push(message); + emit({ type: "message_end", message } as AgentSessionEvent); + }; const harness: AsyncQuiescenceHarness = { session: undefined as unknown as AgentSession, @@ -110,6 +123,7 @@ function createAsyncSession( abortCalls: () => abortCount, settleCalls: () => settleCount, emitTerminalYield, + emitAssistant, finishJob, }; @@ -143,8 +157,9 @@ function createAsyncSession( }, abort: async () => { abortCount += 1; + await options.abort?.(); }, - dispose: async () => {}, + dispose: options.dispose ?? (async () => {}), setIrcWakeTurnObserver: () => {}, }; harness.session = session as unknown as AgentSession; @@ -163,6 +178,7 @@ function mockCreateAgentSession(session: AgentSession) { describe("runSubprocess async quiescence fresh-yield contract", () => { afterEach(() => { vi.restoreAllMocks(); + AsyncJobManager.resetForTests(); }); it("parks a pending yield, injects the result, and completes on the fresh yield", async () => { @@ -287,4 +303,86 @@ describe("runSubprocess async quiescence fresh-yield contract", () => { expect(result.exitCode).toBe(0); expect(result.output).toContain("done"); }); + + it("returns an aborted result after cleanup grace and waits for every late resource", async () => { + const abortStarted = Promise.withResolvers(); + const abortGate = Promise.withResolvers(); + const disposeGate = Promise.withResolvers(); + const lateJobGate = Promise.withResolvers(); + const manager = new AsyncJobManager({}); + AsyncJobManager.setInstance(manager); + let lateJobId: string | undefined; + let deferredCleanup: Promise | undefined; + const harness = createAsyncSession( + ({ promptIndex, harness: h }) => { + if (promptIndex !== 1) return; + h.finishJob(); + h.emitAssistant("captured before cleanup", 7); + h.emitTerminalYield({ report: "yielded output" }); + }, + { + abort: async () => { + abortStarted.resolve(); + await abortGate.promise; + }, + dispose: async () => { + lateJobId = manager.register( + "task", + "shutdown-time job", + async () => { + await lateJobGate.promise; + return "late result"; + }, + { ownerId: "cleanup-timeout" }, + ); + await disposeGate.promise; + }, + }, + ); + mockCreateAgentSession(harness.session); + + const run = runSubprocess({ + cwd: "/tmp", + agent: baseAgent, + task: "do the work", + index: 0, + id: "cleanup-timeout", + keepAlive: false, + onCleanupDeferred: completion => { + deferredCleanup = completion; + }, + }); + await abortStarted.promise; + + const result = await run; + expect(result.exitCode).toBe(1); + expect(result.aborted).toBe(true); + expect(result.abortReason).toBe("cleanup exceeded 10000 ms"); + expect(result.error).toContain("Cleanup did not finish within 10000 ms"); + expect(result.output).toContain("yielded output"); + expect(result.usage?.totalTokens).toBe(7); + expect(lateJobId).toBeDefined(); + expect(deferredCleanup).toBeDefined(); + + let cleanupSettled = false; + const cleanupOutcome = deferredCleanup?.then( + () => { + cleanupSettled = true; + }, + () => { + cleanupSettled = true; + }, + ); + abortGate.resolve(); + disposeGate.reject(new Error("dispose failed")); + for (let attempt = 0; attempt < 10 && manager.getJob(lateJobId ?? "")?.status === "running"; attempt += 1) { + await Promise.resolve(); + } + expect(cleanupSettled).toBe(false); + expect(manager.getJob(lateJobId ?? "")?.status).toBe("cancelled"); + + lateJobGate.resolve(); + await cleanupOutcome; + expect(cleanupSettled).toBe(true); + }, 15_000); });