diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index a4d66e8d2..75861cff9 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed handled OMP shutdown persisting running subagents as terminally aborted instead of restoring their transcripts as parked and revivable. ([#8216](https://github.com/can1357/oh-my-pi/issues/8216)) + ## [17.2.12] - 2026-08-08 ### Fixed diff --git a/packages/coding-agent/src/async/job-manager.ts b/packages/coding-agent/src/async/job-manager.ts index 1fc9e6535..7404a8e10 100644 --- a/packages/coding-agent/src/async/job-manager.ts +++ b/packages/coding-agent/src/async/job-manager.ts @@ -5,6 +5,8 @@ const DELIVERY_RETRY_MAX_MS = 30_000; const DELIVERY_RETRY_JITTER_MS = 200; const DEFAULT_RETENTION_MS = 5 * 60 * 1000; const DEFAULT_MAX_RUNNING_JOBS = 15; +/** Abort reason used only when the owning session shuts down the entire manager. */ +export const ASYNC_JOB_MANAGER_SHUTDOWN_REASON = Symbol("AsyncJobManager shutdown"); /** * Adaptive ("smart") `hub` poll-wait ladder (ms). A tight poll loop climbs @@ -416,9 +418,13 @@ export class AsyncJobManager { * (used by `dispose()` to nuke the manager's state). */ cancelAll(filter?: AsyncJobFilter): void { + this.#cancelJobs(filter); + } + + #cancelJobs(filter?: AsyncJobFilter, reason?: unknown): void { for (const job of this.getRunningJobs(filter)) { job.status = "cancelled"; - job.abortController.abort(); + job.abortController.abort(reason); this.#scheduleEviction(job.id); } } @@ -584,7 +590,7 @@ export class AsyncJobManager { async dispose(options?: { timeoutMs?: number }): Promise { this.#disposed = true; this.#clearEvictionTimers(); - this.cancelAll(); + this.#cancelJobs(undefined, ASYNC_JOB_MANAGER_SHUTDOWN_REASON); const timeoutMs = Math.max(options?.timeoutMs ?? 3_000, 0); const deadline = Date.now() + timeoutMs; const jobsSettled = await this.#waitForAllUntil(deadline); diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index a52793d63..37332477b 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -9,7 +9,7 @@ import type { AgentEvent, AgentIdentity, AgentMessage, AgentTelemetryConfig } fr import { recordHandoff, resolveTelemetry } from "@oh-my-pi/pi-agent-core"; import type { Api, Model, ServiceTierByFamily, Usage } from "@oh-my-pi/pi-ai"; import { logger, popLoopPhase, prompt, pushLoopPhase, untilAborted } from "@oh-my-pi/pi-utils"; -import { AsyncJobManager } from "../async"; +import { ASYNC_JOB_MANAGER_SHUTDOWN_REASON, AsyncJobManager } from "../async"; import type { Rule } from "../capability/rule"; import { ModelRegistry } from "../config/model-registry"; import { @@ -895,7 +895,7 @@ export function createSubagentSettings( ); } -export type AbortReason = "signal" | "terminate" | "timeout" | "budget"; +export type AbortReason = "signal" | "shutdown" | "terminate" | "timeout" | "budget"; const MAX_YIELD_TOOL_ERRORS = 6; @@ -1158,7 +1158,7 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor { signal.addEventListener( "abort", () => { - if (!resolved) requestAbort("signal"); + if (!resolved) requestAbort(signal.reason === ASYNC_JOB_MANAGER_SHUTDOWN_REASON ? "shutdown" : "signal"); }, { once: true, signal: listenerSignal }, ); @@ -1183,6 +1183,7 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor { } const resolveSignalAbortReason = (): string => { + if (signal?.reason === ASYNC_JOB_MANAGER_SHUTDOWN_REASON) return "Async job manager shutdown"; const reason = signal?.reason; if (reason instanceof Error) { const message = reason.message.trim(); @@ -1751,7 +1752,11 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor { runtimeLimitExceeded: () => runtimeLimitExceeded, terminalError: () => terminalError, hasExplicitAbortReason: () => - abortReason === "signal" || runtimeLimitExceeded || budgetLimitExceeded || budgetStopRequested, + abortReason === "signal" || + abortReason === "shutdown" || + runtimeLimitExceeded || + budgetLimitExceeded || + budgetStopRequested, budgetStopRequested: () => budgetStopRequested, waitForBudgetStop: () => budgetStopAbortPromise ?? Promise.resolve(), yieldInvalidatedByAsync: () => yieldInvalidatedByAsync, @@ -1777,7 +1782,11 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor { // the lifecycle can park the agent as resumable instead of killing it. abortKind: () => abortReason ?? (budgetStopRequested ? "budget" : undefined), isAbortedRun: () => - abortReason === "signal" || runtimeLimitExceeded || budgetLimitExceeded || abortReason === undefined, + abortReason === "signal" || + abortReason === "shutdown" || + runtimeLimitExceeded || + budgetLimitExceeded || + abortReason === undefined, requestAbort, failWithError, abortActiveSession, @@ -2416,23 +2425,37 @@ export async function finalizeSubagentLifecycle(args: { } }; - // A budget abort leaves a consistent session with its transcript on disk; - // caller signals, wall-clock timeouts (possible stream hang), and internal + // A budget abort leaves a consistent session with its transcript on disk. + // Manager shutdown also preserves the transcript, but disposes and unregisters + // the process-local session. Caller signals, wall-clock timeouts, and internal // terminations are genuine kills and stay terminal. const resumableAbort = args.abortKind === "budget" && args.keepAlive && !args.isolated && args.reviveSession !== null; if (args.aborted && !resumableAbort) { if (ref && ownsRef) { - // Route hard kills through the lifecycle owner so the terminal - // decision is durable and a restart cannot rediscover the transcript - // as a revivable parked agent. - try { - await AgentLifecycleManager.global().release(args.id, ref, { tombstone: true }); - } catch (error) { - logger.warn("runSubagent: failed to persist kill tombstone", { id: args.id, error: String(error) }); - registry.setStatus(args.id, "aborted", ref); - registry.detachSession(args.id, ref); - await disposeSession(); + if (args.abortKind === "shutdown") { + try { + await AgentLifecycleManager.global().release(args.id, ref); + } catch (error) { + logger.warn("runSubagent: failed to release session during manager shutdown", { + id: args.id, + error: String(error), + }); + await disposeSession(); + registry.unregister(args.id, ref); + } + } else { + // Route hard kills through the lifecycle owner so the terminal + // decision is durable and a restart cannot rediscover the transcript + // as a revivable parked agent. + try { + await AgentLifecycleManager.global().release(args.id, ref, { tombstone: true }); + } catch (error) { + logger.warn("runSubagent: failed to persist kill tombstone", { id: args.id, error: String(error) }); + registry.setStatus(args.id, "aborted", ref); + registry.detachSession(args.id, ref); + await disposeSession(); + } } } else { await disposeSession(); diff --git a/packages/coding-agent/test/agent-session-dispose-concurrent.test.ts b/packages/coding-agent/test/agent-session-dispose-concurrent.test.ts index 09cfb06e3..f2a557ee5 100644 --- a/packages/coding-agent/test/agent-session-dispose-concurrent.test.ts +++ b/packages/coding-agent/test/agent-session-dispose-concurrent.test.ts @@ -3,7 +3,7 @@ import * as path from "node:path"; import { Agent } from "@oh-my-pi/pi-agent-core"; import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock"; import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; -import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async"; +import { ASYNC_JOB_MANAGER_SHUTDOWN_REASON, AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async"; import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; import { HindsightSessionState } from "@oh-my-pi/pi-coding-agent/hindsight/state"; @@ -61,6 +61,33 @@ describe("AgentSession concurrent disposal", () => { return session; } + it("marks owned async job cancellation as manager shutdown", async () => { + const owned = new AsyncJobManager({ maxRunningJobs: 1 }); + const started = Promise.withResolvers(); + let abortReason: unknown; + owned.register("task", "running subagent", async ({ signal }) => { + const aborted = Promise.withResolvers(); + signal.addEventListener( + "abort", + () => { + abortReason = signal.reason; + aborted.resolve(); + }, + { once: true }, + ); + started.resolve(); + await aborted.promise; + return "stopped"; + }); + const current = createSession(owned); + + await started.promise; + await current.dispose(); + session = undefined; + + expect(abortReason).toBe(ASYNC_JOB_MANAGER_SHUTDOWN_REASON); + }); + it("starts independent writers together and closes persistence after their barrier", async () => { const owned = new AsyncJobManager({ maxRunningJobs: 1, retentionMs: 1_000, onJobComplete: () => {} }); const asyncGate = Promise.withResolvers(); diff --git a/packages/coding-agent/test/task/executor-soft-budget.test.ts b/packages/coding-agent/test/task/executor-soft-budget.test.ts index 7fbd28c57..d082e0e6a 100644 --- a/packages/coding-agent/test/task/executor-soft-budget.test.ts +++ b/packages/coding-agent/test/task/executor-soft-budget.test.ts @@ -1,4 +1,5 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; +import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async"; import type { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; import type { LoadExtensionsResult } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/types"; @@ -7,6 +8,7 @@ import { RpcSubagentRegistry } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-sub import type { RpcSubagentFrame } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-types"; import { AgentLifecycleManager } from "@oh-my-pi/pi-coding-agent/registry/agent-lifecycle"; import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry"; +import { registerPersistedSubagents } from "@oh-my-pi/pi-coding-agent/registry/persisted-agents"; import type { CreateAgentSessionResult } from "@oh-my-pi/pi-coding-agent/sdk"; import * as sdkModule from "@oh-my-pi/pi-coding-agent/sdk"; import type { AgentSession, AgentSessionEvent, PromptOptions } from "@oh-my-pi/pi-coding-agent/session/agent-session"; @@ -46,7 +48,8 @@ function createMockSession( promptIndex: number; emit: (event: AgentSessionEvent) => void; pushMessage: (message: unknown) => void; - }) => void, + }) => void | Promise, + onAbort?: () => void | Promise, ): MockSessionHandle { const listeners: Array<(event: AgentSessionEvent) => void> = []; const messages: unknown[] = []; @@ -81,7 +84,7 @@ function createMockSession( prompt: async (text: string, options?: PromptOptions) => { promptIndex += 1; prompts.push({ text, options }); - onPrompt({ promptIndex, emit, pushMessage: message => messages.push(message) }); + await onPrompt({ promptIndex, emit, pushMessage: message => messages.push(message) }); return true; }, waitForIdle: async () => {}, @@ -132,6 +135,7 @@ function createMockSession( }, abort: async () => { abortCount += 1; + await onAbort?.(); }, dispose: async () => { disposeCount += 1; @@ -169,12 +173,14 @@ describe("runSubprocess soft request budget", () => { beforeEach(() => { AgentRegistry.resetGlobalForTests(); AgentLifecycleManager.resetGlobalForTests(); + AsyncJobManager.resetForTests(); tempDir = TempDir.createSync("@pi-soft-budget-"); }); afterEach(() => { vi.restoreAllMocks(); AgentLifecycleManager.resetGlobalForTests(); AgentRegistry.resetGlobalForTests(); + AsyncJobManager.resetForTests(); tempDir[Symbol.dispose](); }); @@ -193,13 +199,13 @@ describe("runSubprocess soft request budget", () => { }; } - function registerRunning(id: string, session: AgentSession) { + function registerRunning(id: string, session: AgentSession, sessionFile: string | null = null) { AgentRegistry.global().register({ id, displayName: id, kind: "sub", session, - sessionFile: null, + sessionFile, status: "running", }); } @@ -338,6 +344,47 @@ describe("runSubprocess soft request budget", () => { rpcRegistry.dispose(); }); + it("manager shutdown restores a running kept-alive agent as parked without a tombstone", async () => { + const id = "ShutdownScout"; + const rootSessionFile = `${tempDir.path()}/main.jsonl`; + const workerSessionFile = `${tempDir.path()}/main/${id}.jsonl`; + await Bun.write(rootSessionFile, ""); + await Bun.write(workerSessionFile, ""); + const promptStarted = Promise.withResolvers(); + const promptStopped = Promise.withResolvers(); + const handle = createMockSession( + async ({ promptIndex }) => { + if (promptIndex !== 1) return; + promptStarted.resolve(); + await promptStopped.promise; + }, + () => promptStopped.resolve(), + ); + mockCreateAgentSession(handle.session); + registerRunning(id, handle.session, workerSessionFile); + const manager = new AsyncJobManager({ maxRunningJobs: 1 }); + AsyncJobManager.setInstance(manager); + manager.register( + "task", + "shutdown regression", + async ({ signal }) => { + const result = await runSubprocess({ ...baseOptions(id), signal }); + return result.output; + }, + { ownerId: "Main", agentId: id }, + ); + + await promptStarted.promise; + await manager.dispose({ timeoutMs: 1_000 }); + AsyncJobManager.setInstance(undefined); + + expect(await Bun.file(`${workerSessionFile}.tombstone`).exists()).toBe(false); + expect(AgentRegistry.global().get(id)).toBeUndefined(); + const restoredRegistry = new AgentRegistry(); + await registerPersistedSubagents(restoredRegistry, rootSessionFile); + expect(restoredRegistry.get(id)?.status).toBe("parked"); + }); + it("a caller-signal abort stays terminal and irc names the aborted agent precisely", async () => { const id = "CancelledScout"; const controller = new AbortController();