diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index f4f440810..61e7000ba 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -98,6 +98,7 @@ - Fixed floating rejections from cmux browser guest JavaScript terminating the main process and every active session; attributable rejections now fail the browser run as tool errors while unrelated process rejections retain the fatal path ([#7365](https://github.com/can1357/oh-my-pi/issues/7365)). - Fixed the Windows bash tool silently taking down the whole omp process when a command blocked until its timeout: cancelling a timed-out run walked the spawned child's descendant tree from raw `th32ParentProcessID` links, and a recycled pid matching the harness's stale recorded parent pid could enumerate omp as a false descendant and `TerminateProcess` it, killing the session with no `session_exit` record. Run-cancellation sweeps now refuse to signal the harness or any process collected beneath it, while still reaping the timed-out target when it owns a recycled ancestor pid ([#7452](https://github.com/can1357/oh-my-pi/issues/7452)). - Fixed the unexpected-stop guard (`features.unexpectedStopDetection`) never firing for thinking-only stops: `isUnexpectedStopCandidate` only counted non-whitespace `text` blocks, so a `stopReason: "stop"` turn whose sole content was a signed `thinking` block (a trapped response or a truncated reasoning fragment from reasoning models) bypassed classification and silently ended the turn mid-task. Such stops are now candidates and are classified on their thinking text ([#7499](https://github.com/can1357/oh-my-pi/issues/7499)). +- Fixed Task cancellation hanging forever when a child ignored abort or stalled during cleanup ([#7483](https://github.com/can1357/oh-my-pi/issues/7483)). ## [17.2.5] - 2026-08-03 diff --git a/packages/coding-agent/src/async/job-manager.ts b/packages/coding-agent/src/async/job-manager.ts index 71d498c01..1fc9e6535 100644 --- a/packages/coding-agent/src/async/job-manager.ts +++ b/packages/coding-agent/src/async/job-manager.ts @@ -94,6 +94,12 @@ export interface AsyncJobDeliveryState { pendingJobIds: string[]; } +export interface AsyncJobReapResult { + settled: boolean; + pendingJobIds: string[]; + completion: Promise; +} + export interface AsyncJobRegisterOptions { id?: string; /** Registry id of the agent that owns this job; used to scope cancelAll. */ @@ -490,6 +496,26 @@ export class AsyncJobManager { } } + /** + * Cancel every job owned by `ownerId`, then wait only until `deadlineAt`. + * The returned completion keeps waiting for actual process settlement when + * the deadline expires, so callers can move that cleanup out of the + * user-visible Task wait without losing ownership of the live work. + */ + async cancelAndReapOwnerJobs(ownerId: string, deadlineAt: number): Promise { + this.cancelAll({ ownerId }); + const timeoutMs = Math.max(0, deadlineAt - Date.now()); + const settled = await this.waitForOwnerJobs(ownerId, { timeoutMs }); + if (settled) { + return { settled: true, pendingJobIds: [], completion: Promise.resolve() }; + } + const pendingJobIds = this.getAllJobs({ ownerId }) + .filter(job => job.status === "running" || job.status === "cancelled") + .map(job => job.id); + const completion = this.waitForOwnerJobs(ownerId).then(() => {}); + return { settled: false, pendingJobIds, completion }; + } + async #waitForAllUntil(deadline: number): Promise { const promises = Array.from(this.#jobs.values()).map(job => job.promise); if (promises.length === 0) return true; diff --git a/packages/coding-agent/src/registry/agent-lifecycle.ts b/packages/coding-agent/src/registry/agent-lifecycle.ts index 3864c60f1..525ee5116 100644 --- a/packages/coding-agent/src/registry/agent-lifecycle.ts +++ b/packages/coding-agent/src/registry/agent-lifecycle.ts @@ -20,8 +20,9 @@ * a superseded revive) can never clobber a newer same-id ref. */ -import { logger } from "@oh-my-pi/pi-utils"; +import { logger, untilAborted } from "@oh-my-pi/pi-utils"; import type { AgentSession } from "../session/agent-session"; +import { trackLateCleanup } from "../utils/late-cleanup"; import { type AgentRef, type AgentRefExpectation, @@ -32,6 +33,8 @@ import { export type AgentReviver = (expected: AgentRef) => Promise; +const AGENT_RELEASE_GRACE_MS = 5000; + /** * Builds a reviver for a `parked` ref restored from disk (Agent Hub scan, * collab mirror, resumed process) that carries a sessionFile but no in-memory @@ -402,11 +405,26 @@ export class AgentLifecycleManager { } /** Teardown everything (process exit / main session dispose). */ - async dispose(): Promise { + async dispose(deadlineAt: number = Date.now() + AGENT_RELEASE_GRACE_MS): Promise { this.#unsubscribe?.(); this.#unsubscribe = undefined; const ids = [...new Set([...this.#adopted.keys(), ...this.#parks.keys()])]; - await Promise.all(ids.map(id => this.release(id))); + await Promise.all( + ids.map(async id => { + const release = this.release(id).then(() => {}); + try { + await untilAborted(AbortSignal.timeout(Math.max(0, deadlineAt - Date.now())), () => release); + } catch (error) { + if (Date.now() >= deadlineAt) { + trackLateCleanup(release, { id, resource: "adopted-agent" }); + } + logger.warn("Agent cleanup exceeded its deadline", { + id, + error: error instanceof Error ? error.message : String(error), + }); + } + }), + ); this.#revivals.clear(); this.#parks.clear(); this.#persistedReviverFactory = undefined; diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index 5538130a0..610696357 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -54,6 +54,7 @@ import { normalizeSchema } from "../tools/jtd-to-json-schema"; import { buildOutputValidator, summarizeValidationFailure } from "../tools/output-schema-validator"; import { ToolAbortError } from "../tools/tool-errors"; import type { EventBus } from "../utils/event-bus"; +import { trackLateCleanup } from "../utils/late-cleanup"; import { buildNamedToolChoice } from "../utils/tool-choice"; import type { WorkspaceTree } from "../workspace-tree"; import { generateTaskLabel } from "./label"; @@ -79,6 +80,7 @@ import { arrayValuedLabels, assembleYieldResult } from "./yield-assembly"; export type { YieldItem } from "./types"; const MCP_CALL_TIMEOUT_MS = 60_000; +const TASK_ABORT_CLEANUP_GRACE_MS = 10_000; /** * Soft per-agent request budgets (assistant requests per run). Crossing the @@ -455,6 +457,8 @@ export interface ExecutorOptions { * set this false so disposal unregisters them instead of leaving idle peers. */ keepAlive?: boolean; + /** Internal ownership handoff for cleanup that outlives the visible Task result. */ + onCleanupDeferred?: (completion: Promise) => void; } function parseStringifiedJson(value: unknown): unknown { @@ -1964,9 +1968,7 @@ async function driveSessionToYield( // yield: the next iteration's ladder demands a fresh one. } - if (monitor.yieldCalled()) { - await session.waitForIdle(); - } else { + if (!monitor.yieldCalled()) { await awaitAbortable(session.waitForIdle()); } @@ -2140,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); @@ -2346,15 +2351,27 @@ export async function finalizeSubagentLifecycle(args: { isolated: boolean; agentIdleTtlMs: number; reviveSession: AgentReviver | null; + cleanupDeadlineAt?: number; + onCleanupDeferred?: (completion: Promise) => void; }): Promise { const registry = AgentRegistry.global(); const ref = registry.get(args.id); const ownsRef = Boolean(ref && ref.session === args.session); + const cleanupDeadlineAt = args.cleanupDeadlineAt ?? Date.now() + 5000; const disposeSession = async (): Promise => { + const disposal = args.session.dispose(); + const remainingMs = Math.max(0, cleanupDeadlineAt - Date.now()); try { - await untilAborted(AbortSignal.timeout(5000), () => args.session.dispose()); - } catch { - // Ignore cleanup errors + await untilAborted(AbortSignal.timeout(remainingMs), () => disposal); + } catch (error) { + if (Date.now() >= cleanupDeadlineAt) { + args.onCleanupDeferred?.(disposal); + return; + } + logger.warn("Subagent session cleanup failed", { + id: args.id, + error: error instanceof Error ? error.message : String(error), + }); } }; @@ -3186,6 +3203,20 @@ export async function runSubprocess(options: ExecutorOptions): Promise[] = []; + let deferredSessionShutdown: Promise | undefined; + const deferCleanup = (completion: Promise): void => { + lateCleanups.push(completion); + exitCode = 1; + aborted = true; + abortReasonText = `cleanup exceeded ${TASK_ABORT_CLEANUP_GRACE_MS} ms`; + error ??= `Task aborted. Cleanup did not finish within ${TASK_ABORT_CLEANUP_GRACE_MS} ms. ${cleanupChangeStatus}`; + }; if (abortSignal.aborted) { aborted = monitor.isAbortedRun(); if (aborted) { @@ -3194,10 +3225,21 @@ export async function runSubprocess(options: ExecutorOptions): Promise monitor.waitForActiveSessionAbort()); - } catch { - // Ignore abort cleanup timeouts/errors; terminal disposal below is still best-effort. + await untilAborted( + AbortSignal.timeout(Math.max(0, cleanupDeadlineAt - Date.now())), + () => activeSessionAbort, + ); + } catch (cleanupError) { + if (Date.now() >= cleanupDeadlineAt) { + deferCleanup(activeSessionAbort); + } else { + logger.warn("Subagent abort cleanup failed", { + id, + error: cleanupError instanceof Error ? cleanupError.message : String(cleanupError), + }); + } } if (unsubscribe) { try { @@ -3207,6 +3249,17 @@ export async function runSubprocess(options: ExecutorOptions): Promise { + deferredSessionShutdown = completion; + deferCleanup(completion); + }, }); } - // Structured-concurrency reap: cancel and await ALL surviving owner - // jobs (abort paths; suppressed/watched jobs the model left behind) - // so isolation capture/cleanup never races a live process writing - // into the worktree. This never proceeds while an owner process is - // live: cancellation SIGKILL-escalates, so settlement is expected - // within one interval — an unkillable process blocks here visibly - // (with periodic warnings) instead of silently racing teardown. - const jobManager = AsyncJobManager.instance(); if (jobManager) { - jobManager.cancelAll({ ownerId: id }); - while (!(await jobManager.waitForOwnerJobs(id, { timeoutMs: 10_000 }))) { - logger.warn("Subagent async jobs still settling; delaying teardown until process exit", { id }); + 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.allSettled(lateCleanups).then(() => {}); + trackLateCleanup(completion, { id, resource: "subagent" }); + options.onCleanupDeferred?.(completion); + } } // Launch-latency breakdown (subagent invocation → first chat dispatch). diff --git a/packages/coding-agent/src/task/isolation-runner.ts b/packages/coding-agent/src/task/isolation-runner.ts index c6e0a58f7..487d121ed 100644 --- a/packages/coding-agent/src/task/isolation-runner.ts +++ b/packages/coding-agent/src/task/isolation-runner.ts @@ -23,6 +23,7 @@ import type * as natives from "@oh-my-pi/pi-natives"; import type { ToolSession } from "../tools"; import { generateCommitMessage } from "../utils/commit-message-generator"; import * as git from "../utils/git"; +import { trackLateCleanup } from "../utils/late-cleanup"; import type { ExecutorOptions } from "./executor"; import { runSubprocess } from "./executor"; import type { SingleResult } from "./types"; @@ -146,6 +147,7 @@ async function writeIsolationPatch( */ export async function runIsolatedSubprocess(opts: IsolatedRunOptions): Promise { let handle: IsolationHandle | undefined; + let deferredCleanup: Promise | undefined; try { const taskBaseline = structuredClone(opts.context.baseline); handle = await ensureIsolation(opts.context.repoRoot, opts.agentId, opts.preferredBackend); @@ -155,7 +157,12 @@ export async function runIsolatedSubprocess(opts: IsolatedRunOptions): Promise { + deferredCleanup = completion; + opts.baseOptions.onCleanupDeferred?.(completion); + }, }); + if (deferredCleanup) return result; if (opts.mergeMode === "branch" && result.exitCode === 0) { try { const commitResult = await commitToBranch( @@ -213,7 +220,18 @@ export async function runIsolatedSubprocess(opts: IsolatedRunOptions): Promise cleanupIsolation(isolationHandle)), + { + agentId: opts.agentId, + resource: "isolation", + }, + ); + } else { + await cleanupIsolation(isolationHandle); + } } } } diff --git a/packages/coding-agent/src/task/structured-subagent.ts b/packages/coding-agent/src/task/structured-subagent.ts index 0726a102c..36680cbed 100644 --- a/packages/coding-agent/src/task/structured-subagent.ts +++ b/packages/coding-agent/src/task/structured-subagent.ts @@ -20,6 +20,7 @@ import type { TaskEffort } from "../thinking"; import type { ToolSession } from "../tools"; import { isIrcEnabled } from "../tools/hub"; import { buildOutputValidator } from "../tools/output-schema-validator"; +import { trackLateCleanup } from "../utils/late-cleanup"; import { type DiscoveryResult, discoverAgents, getAgent } from "./discovery"; import { type ExecutorOptions, runSubprocess } from "./executor"; import { @@ -540,12 +541,16 @@ export async function runStructuredSubagent(request: StructuredSubagentRequest): let mergeSummary = ""; let requiresRecoveryArtifacts = false; let completedSuccessfully = false; + let deferredCleanup: Promise | undefined; try { const id = await reserveStructuredSubagentId(request.session, { ...request.identity, label: request.identity?.label ?? (request.invocationKind === "eval" ? "EvalAgent" : undefined), }); const baseOptions = buildExecutorOptions(request, policy, lease, id); + baseOptions.onCleanupDeferred = completion => { + deferredCleanup = completion; + }; baseOptions.planReference = await loadPlanReference(request, policy); let isolationContext: IsolationContext | null = null; if (policy.isIsolated) { @@ -639,8 +644,18 @@ export async function runStructuredSubagent(request: StructuredSubagentRequest): (policy.isIsolated && (!policy.applyChanges || changesApplied === false || requiresRecoveryArtifacts)); const shouldCleanup = lease.temporary && !shouldRetainArtifacts; if (shouldCleanup) { - await fs.rm(lease.artifactsDir, { recursive: true, force: true }); - lease.unregister?.(); + const cleanupArtifacts = async (): Promise => { + await fs.rm(lease.artifactsDir, { recursive: true, force: true }); + lease.unregister?.(); + }; + if (deferredCleanup) { + trackLateCleanup(deferredCleanup.then(cleanupArtifacts), { + resource: "artifacts", + artifactsDir: lease.artifactsDir, + }); + } else { + await cleanupArtifacts(); + } } } } diff --git a/packages/coding-agent/src/utils/late-cleanup.ts b/packages/coding-agent/src/utils/late-cleanup.ts new file mode 100644 index 000000000..7d9ae2d8d --- /dev/null +++ b/packages/coding-agent/src/utils/late-cleanup.ts @@ -0,0 +1,17 @@ +import { logger } from "@oh-my-pi/pi-utils"; + +const pendingCleanups = new Set>(); + +/** Keep timed-out cleanup reachable until its resources really settle. */ +export function trackLateCleanup(work: Promise, context: Record): void { + let tracked: Promise; + tracked = work + .catch(error => { + logger.warn("Deferred cleanup failed", { + ...context, + error: error instanceof Error ? error.message : String(error), + }); + }) + .finally(() => pendingCleanups.delete(tracked)); + pendingCleanups.add(tracked); +} diff --git a/packages/coding-agent/test/async-job-manager.test.ts b/packages/coding-agent/test/async-job-manager.test.ts index 1a336456a..21fe7529a 100644 --- a/packages/coding-agent/test/async-job-manager.test.ts +++ b/packages/coding-agent/test/async-job-manager.test.ts @@ -113,6 +113,30 @@ describe("AsyncJobManager", () => { expect(completions).toHaveLength(0); }); + test("bounds owner-job reap while preserving late settlement", async () => { + const manager = new AsyncJobManager({ onJobComplete: async () => {} }); + const release = Promise.withResolvers(); + const jobId = manager.register( + "task", + "ignores abort", + async () => { + await release.promise; + return "late result"; + }, + { ownerId: "owner" }, + ); + + const reap = await manager.cancelAndReapOwnerJobs("owner", Date.now()); + + expect(reap.settled).toBe(false); + expect(reap.pendingJobIds).toEqual([jobId]); + expect(manager.getJob(jobId)?.status).toBe("cancelled"); + + release.resolve(); + await reap.completion; + expect(manager.getJob(jobId)?.resultText).toBe("late result"); + }); + test("enforces maxRunningJobs cap", () => { const manager = new AsyncJobManager({ maxRunningJobs: 1, diff --git a/packages/coding-agent/test/registry/agent-lifecycle.test.ts b/packages/coding-agent/test/registry/agent-lifecycle.test.ts index 6c3356f66..1de361c77 100644 --- a/packages/coding-agent/test/registry/agent-lifecycle.test.ts +++ b/packages/coding-agent/test/registry/agent-lifecycle.test.ts @@ -314,6 +314,23 @@ describe("AgentLifecycleManager", () => { expect(registry.get("6-Sub")).toBeUndefined(); }); + it("does not let one stuck adopted agent block sibling disposal", async () => { + const gate = deferred(); + const stuck = makeSessionStub(() => gate.promise); + const sibling = makeSessionStub(); + registerIdleSub("stuck-Sub", stuck.session); + registerIdleSub("sibling-Sub", sibling.session); + lifecycle.adopt("stuck-Sub", { idleTtlMs: TTL }); + lifecycle.adopt("sibling-Sub", { idleTtlMs: TTL }); + + await lifecycle.dispose(Date.now()); + + expect(stuck.disposeCalls()).toBe(1); + expect(sibling.disposeCalls()).toBe(1); + gate.resolve(); + await flushAsync(); + }); + it("a delayed release cannot remove or mutate a replacement ref with the same id", async () => { const gate = deferred(); const oldSession = makeSessionStub(() => gate.promise); 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 69e8e7dc8..ca95f5f5f 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 () => { @@ -250,4 +266,125 @@ describe("runSubprocess async quiescence fresh-yield contract", () => { expect(result.exitCode).toBe(0); expect(result.output).toContain("done"); }); + + it("does not wait on a second idle barrier after a terminal yield", async () => { + const harness = createAsyncSession(({ promptIndex, harness: h }) => { + if (promptIndex === 1) { + h.finishJob(); + h.emitTerminalYield({ report: "done" }); + } + }); + const idleStarted = Promise.withResolvers(); + const releaseIdle = Promise.withResolvers(); + let idleCalls = 0; + harness.session.waitForIdle = async () => { + idleCalls += 1; + idleStarted.resolve(); + await releaseIdle.promise; + }; + mockCreateAgentSession(harness.session); + + const run = runSubprocess({ + cwd: "/tmp", + agent: baseAgent, + task: "do the work", + index: 0, + id: "quiescence-no-second-idle", + }); + const outcome = await Promise.race([ + run.then(() => "completed" as const), + idleStarted.promise.then(() => "blocked" as const), + ]); + releaseIdle.resolve(); + const result = await run; + + expect(outcome).toBe("completed"); + expect(idleCalls).toBe(0); + 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).toBe( + "Task aborted. Cleanup did not finish within 10000 ms. This task was not isolated, so its changes may remain in the working directory.", + ); + 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); }); diff --git a/packages/coding-agent/test/task/isolation-runner.test.ts b/packages/coding-agent/test/task/isolation-runner.test.ts index 6a9a3c242..706196a18 100644 --- a/packages/coding-agent/test/task/isolation-runner.test.ts +++ b/packages/coding-agent/test/task/isolation-runner.test.ts @@ -139,6 +139,63 @@ describe("runIsolatedSubprocess", () => { expect(deleteSpy).toHaveBeenCalledWith(repoRoot, "omp/task/PreserveBranchFailure"); expect(cleanupSpy).toHaveBeenCalledTimes(1); }); + + it("keeps an isolated worktree until deferred child cleanup settles", async () => { + const cleanupGate = Promise.withResolvers(); + vi.spyOn(worktreeModule, "ensureIsolation").mockResolvedValue({ + mergedDir: "/repo/isolated", + backend: natives.IsoBackendKind.Rcopy, + fellBack: false, + fallbackReason: null, + }); + vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => { + options.onCleanupDeferred?.(cleanupGate.promise); + return result({ exitCode: 1, aborted: true, error: "cleanup exceeded its deadline" }); + }); + const cleanupSpy = vi.spyOn(worktreeModule, "cleanupIsolation").mockResolvedValue(); + + const outcome = await runIsolatedSubprocess({ + baseOptions: { + cwd: "/repo", + agent: { + name: "task", + description: "Task agent", + systemPrompt: "test", + source: "bundled", + }, + task: "Do work", + index: 0, + id: "DeferredCleanup", + }, + context: { + repoRoot: "/repo", + baseline: { + root: { + repoRoot: "/repo", + headCommit: "base", + staged: "", + unstaged: "", + untracked: [], + untrackedPatch: "", + }, + nested: [], + }, + }, + preferredBackend: undefined, + agentId: "DeferredCleanup", + mergeMode: "patch", + artifactsDir: "/artifacts", + buildFailureResult: error => result({ exitCode: 1, error: String(error) }), + }); + + expect(outcome.exitCode).toBe(1); + expect(cleanupSpy).not.toHaveBeenCalled(); + cleanupGate.resolve(); + await cleanupGate.promise; + await Promise.resolve(); + await Promise.resolve(); + expect(cleanupSpy).toHaveBeenCalledTimes(1); + }); }); describe("mergeIsolatedChanges", () => {