feat(coding-agent): added AgentLifecycleManager for idle/parked subagents

Introduces AgentLifecycleManager: when the task executor adopts a finished subagent the manager arms a TTL timer on idle, parks the agent on expiry (disposes the live session, keeps the AgentRef + sessionFile), and revives it on demand via an injected reviver. Only this manager flips parked ↔ idle. AgentRegistry now annotates session: AgentRef["session"] as null exactly when parked/aborted, and sdk.ts wires the lifecycle dispose into main-session teardown, derives agentKind once, and gates ref unregistration on parking so a parked agent stays addressable (history://, revive).
This commit is contained in:
can1357
2026-06-10 17:46:02 +02:00
parent 5406edeed0
commit 1654c759ca
4 changed files with 497 additions and 14 deletions
@@ -0,0 +1,218 @@
/**
* AgentLifecycleManager - Owns the idle → parked → revived lifecycle of
* adopted subagents.
*
* The task executor hands a finished agent over via {@link AgentLifecycleManager.adopt};
* from then on the manager arms a TTL timer whenever the agent goes `idle`,
* parks it on expiry (disposes the live session, keeps the AgentRef +
* sessionFile), and revives it on demand through
* {@link AgentLifecycleManager.ensureLive}. Only this manager flips
* `parked` ↔ `idle`.
*/
import { logger } from "@oh-my-pi/pi-utils";
import type { AgentSession } from "../session/agent-session";
import { AgentRegistry, MAIN_AGENT_ID, type RegistryEvent } from "./agent-registry";
export type AgentReviver = () => Promise<AgentSession>;
export interface AdoptOptions {
/** TTL before an idle agent is parked. <= 0 disables parking. */
idleTtlMs: number;
/** Recreates a live AgentSession from the ref's sessionFile. Absent => not resumable after park (e.g. isolated runs). */
revive?: AgentReviver;
}
interface AdoptedAgent {
idleTtlMs: number;
revive?: AgentReviver;
timer?: NodeJS.Timeout;
}
export class AgentLifecycleManager {
static #global: AgentLifecycleManager | undefined;
static global(): AgentLifecycleManager {
if (!AgentLifecycleManager.#global) {
AgentLifecycleManager.#global = new AgentLifecycleManager();
}
return AgentLifecycleManager.#global;
}
/** Reset the global manager. Test-only. */
static resetGlobalForTests(): void {
const current = AgentLifecycleManager.#global;
if (current) {
current.#unsubscribe?.();
current.#unsubscribe = undefined;
for (const adopted of current.#adopted.values()) {
clearTimeout(adopted.timer);
}
current.#adopted.clear();
current.#revivals.clear();
current.#parking.clear();
}
AgentLifecycleManager.#global = undefined;
}
readonly #registry: AgentRegistry;
readonly #adopted = new Map<string, AdoptedAgent>();
/** Ids whose session is being disposed by {@link park} right now. */
readonly #parking = new Set<string>();
/** In-flight revives, so concurrent {@link ensureLive} calls coalesce. */
readonly #revivals = new Map<string, Promise<AgentSession>>();
#unsubscribe: (() => void) | undefined;
constructor(registry: AgentRegistry = AgentRegistry.global()) {
this.#registry = registry;
this.#unsubscribe = registry.onChange(event => this.#onRegistryEvent(event));
}
/**
* Take ownership of a finished subagent. Caller has already set registry
* status to "idle". Arms the TTL timer (idleTtlMs <= 0 adopts without one).
*/
adopt(id: string, opts: AdoptOptions): void {
if (id === MAIN_AGENT_ID) return;
if (!this.#registry.get(id)) {
logger.warn("AgentLifecycleManager.adopt: unknown agent id", { id });
return;
}
const existing = this.#adopted.get(id);
clearTimeout(existing?.timer);
const adopted: AdoptedAgent = { idleTtlMs: opts.idleTtlMs, revive: opts.revive };
this.#adopted.set(id, adopted);
this.#armTimer(id, adopted);
}
/** True if the id is adopted (parked or live). */
has(id: string): boolean {
return this.#adopted.has(id);
}
/** True while {@link park} is disposing this agent's session (lets dispose hooks distinguish park from teardown). */
isParking(id: string): boolean {
return this.#parking.has(id);
}
/**
* Dispose the live session, detach it from the registry, and mark the
* agent `parked`. No-op unless the id is adopted and live.
*/
async park(id: string): Promise<void> {
const adopted = this.#adopted.get(id);
if (!adopted) return;
const ref = this.#registry.get(id);
if (!ref?.session) return;
if (adopted.timer) {
clearTimeout(adopted.timer);
adopted.timer = undefined;
}
this.#parking.add(id);
try {
try {
await ref.session.dispose();
} catch (error) {
logger.warn("AgentLifecycleManager.park: session dispose failed", { id, error: String(error) });
}
this.#registry.detachSession(id);
this.#registry.setStatus(id, "parked");
} finally {
this.#parking.delete(id);
}
}
/**
* Return the live session, reviving from the sessionFile if parked.
* Throws a plain Error if the id is unknown or parked without a reviver.
* Concurrent calls share one in-flight revive.
*/
async ensureLive(id: string): Promise<AgentSession> {
const ref = this.#registry.get(id);
if (!ref) {
throw new Error(
`Unknown agent "${id}" — it was never registered or has been released. If a transcript exists, read history://${id}.`,
);
}
if (ref.session) return ref.session;
const inflight = this.#revivals.get(id);
if (inflight) return inflight;
const adopted = this.#adopted.get(id);
if (ref.status !== "parked" || !adopted?.revive) {
throw new Error(
`Agent "${id}" is ${ref.status} and cannot be revived${adopted?.revive ? "" : " (no reviver registered)"}. Its transcript remains readable at history://${id}.`,
);
}
const revival = this.#revive(id, adopted.revive, ref.sessionFile);
this.#revivals.set(id, revival);
try {
return await revival;
} finally {
this.#revivals.delete(id);
}
}
/** Hard removal: dispose if live, unregister from registry, drop timers. */
async release(id: string): Promise<void> {
const adopted = this.#adopted.get(id);
clearTimeout(adopted?.timer);
this.#adopted.delete(id);
const ref = this.#registry.get(id);
if (ref?.session) {
try {
await ref.session.dispose();
} catch (error) {
logger.warn("AgentLifecycleManager.release: session dispose failed", { id, error: String(error) });
}
}
this.#registry.unregister(id);
}
/** Teardown everything (process exit / main session dispose). */
async dispose(): Promise<void> {
this.#unsubscribe?.();
this.#unsubscribe = undefined;
const ids = [...this.#adopted.keys()];
await Promise.all(ids.map(id => this.release(id)));
this.#revivals.clear();
this.#parking.clear();
}
async #revive(id: string, revive: AgentReviver, sessionFile: string | null): Promise<AgentSession> {
const session = await revive();
this.#registry.attachSession(id, session, sessionFile);
// Emits status_changed → "idle", which re-arms the TTL timer below.
this.#registry.setStatus(id, "idle");
return session;
}
#armTimer(id: string, adopted: AdoptedAgent): void {
if (adopted.idleTtlMs <= 0) return;
clearTimeout(adopted.timer);
const timer = setTimeout(() => {
adopted.timer = undefined;
void this.park(id);
}, adopted.idleTtlMs);
timer.unref?.();
adopted.timer = timer;
}
#onRegistryEvent(event: RegistryEvent): void {
const adopted = this.#adopted.get(event.ref.id);
if (!adopted) return;
if (event.type === "removed") {
clearTimeout(adopted.timer);
this.#adopted.delete(event.ref.id);
return;
}
if (event.type !== "status_changed") return;
if (event.ref.status === "running") {
if (adopted.timer) {
clearTimeout(adopted.timer);
adopted.timer = undefined;
}
} else if (event.ref.status === "idle") {
this.#armTimer(event.ref.id, adopted);
}
}
}
@@ -1,16 +1,26 @@
/**
* AgentRegistry - Process-global registry of live AgentSession instances.
* AgentRegistry - Process-global registry of agents (the main session plus
* every subagent), keyed by stable id.
*
* Tracks every alive agent (the main session plus every subagent) so the
* `irc` tool can address peers by id. Sessions are registered explicitly at
* creation and removed when the owner releases them.
* Tracks each agent's status and (when live) its AgentSession so peers can be
* addressed by id (`irc`, `task resume`, `history://`). Sessions are
* registered explicitly at creation; finished agents stay registered as
* `idle` (live) or `parked` (session disposed, ref + sessionFile retained for
* revival) and are only removed on explicit release/teardown.
*/
import type { AgentSession } from "../session/agent-session";
export const MAIN_AGENT_ID = "Main";
export type AgentStatus = "running" | "idle" | "completed" | "aborted";
/**
* - `running`: a turn is in flight.
* - `idle`: live AgentSession in memory, awaiting work. Finished agents are
* `idle`, not removed.
* - `parked`: session disposed; AgentRef + sessionFile retained, revivable.
* - `aborted`: hard-killed, terminal.
*/
export type AgentStatus = "running" | "idle" | "parked" | "aborted";
export type AgentKind = "main" | "sub";
export interface AgentRef {
@@ -19,6 +29,7 @@ export interface AgentRef {
kind: AgentKind;
parentId?: string;
status: AgentStatus;
/** Null exactly when parked/aborted. */
session: AgentSession | null;
sessionFile: string | null;
createdAt: number;
+29 -9
View File
@@ -34,7 +34,7 @@ import {
Snowflake,
} from "@oh-my-pi/pi-utils";
import chalk from "chalk";
import { type AsyncJob, AsyncJobManager, isBackgroundJobSupportEnabled } from "./async";
import { type AsyncJob, AsyncJobManager } from "./async";
import { loadCapability } from "./capability";
import { type Rule, ruleCapability, setActiveRules } from "./capability/rule";
import { bucketRules } from "./capability/rule-buckets";
@@ -93,6 +93,7 @@ import { createSessionMemoryRuntimeContext, resolveMemoryBackend } from "./memor
import type { MnemopiSessionState } from "./mnemopi/state";
import asyncResultTemplate from "./prompts/tools/async-result.md" with { type: "text" };
import lateDiagnosticTemplate from "./prompts/tools/lsp-late-diagnostic.md" with { type: "text" };
import { AgentLifecycleManager } from "./registry/agent-lifecycle";
import { AgentRegistry, MAIN_AGENT_ID } from "./registry/agent-registry";
import {
collectEnvSecrets,
@@ -1293,7 +1294,6 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
let hasSession = false;
let hasRegistered = false;
const enableLsp = options.enableLsp ?? true;
const backgroundJobsEnabled = isBackgroundJobSupportEnabled(settings);
const asyncMaxJobs = Math.min(100, Math.max(1, settings.get("async.maxJobs") ?? 100));
const ASYNC_INLINE_RESULT_MAX_CHARS = 12_000;
const ASYNC_PREVIEW_MAX_CHARS = 4_000;
@@ -1326,7 +1326,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
// (issue #1923). The `instance()` guard means later sessions also skip
// constructing an orphaned manager that nothing would ever route to.
const asyncJobManager =
backgroundJobsEnabled && !options.parentTaskPrefix && !AsyncJobManager.instance()
!options.parentTaskPrefix && !AsyncJobManager.instance()
? new AsyncJobManager({
maxRunningJobs: asyncMaxJobs,
onJobComplete: async (jobId, result, job) => {
@@ -1351,6 +1351,17 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
const resolvedAgentId = options.agentId ?? options.parentTaskPrefix ?? MAIN_AGENT_ID;
const resolvedAgentDisplayName =
options.agentDisplayName ?? ((options.taskDepth ?? 0) > 0 || options.parentTaskPrefix ? "sub" : "main");
const agentKind = (options.taskDepth ?? 0) > 0 || options.parentTaskPrefix ? ("sub" as const) : ("main" as const);
/**
* Forget the agent ref on teardown — unless the agent is being parked (or is
* already parked). Parking disposes the session but keeps the ref addressable
* (history://, revive); only process teardown / explicit kill unregisters.
*/
const unregisterUnlessParked = (): void => {
if (agentRegistry.get(resolvedAgentId)?.status === "parked") return;
if (AgentLifecycleManager.global().isParking(resolvedAgentId)) return;
agentRegistry.unregister(resolvedAgentId);
};
const evalKernelOwnerId = `agent-session:${Snowflake.next()}`;
try {
@@ -1409,7 +1420,6 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
getTurnBudget: () => sessionManager.getTurnBudget(),
recordEvalSubagentUsage: output => sessionManager.recordEvalSubagentOutput(output),
getClientBridge: () => session?.clientBridge,
getCompactContext: () => session.formatCompactContext(),
queueDeferredDiagnostics: entry => session?.yieldQueue.enqueue(LSP_LATE_DIAGNOSTIC_MESSAGE_TYPE, entry),
bumpFileMutationVersion: path => {
const next = (fileMutationVersions.get(path) ?? 0) + 1;
@@ -2083,7 +2093,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
agentRegistry.register({
id: resolvedAgentId,
displayName: resolvedAgentDisplayName,
kind: (options.taskDepth ?? 0) > 0 || options.parentTaskPrefix ? "sub" : "main",
kind: agentKind,
parentId: options.parentTaskPrefix,
session: null,
sessionFile: sessionManager.getSessionFile() ?? null,
@@ -2320,7 +2330,6 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
ttsrManager,
obfuscator,
agentId: resolvedAgentId,
agentRegistry,
providerSessionId: options.providerSessionId,
parentEvalSessionId: options.parentEvalSessionId,
});
@@ -2341,15 +2350,26 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
// Attach the live session to the pre-registered ref so peers can route IRC
// messages here. Refresh sessionFile in case it was unavailable at pre-register
// time. The dispose wrapper below unregisters on teardown.
// time. The dispose wrapper below unregisters on teardown (unless parked).
agentRegistry.attachSession(resolvedAgentId, session, sessionManager.getSessionFile() ?? null);
{
const originalDispose = session.dispose.bind(session);
session.dispose = async () => {
try {
// Reject new session work (Python/eval starts) the moment disposal
// begins — the lifecycle await below opens an async gap before
// AgentSession.dispose() would otherwise set its guards.
session.beginDispose();
if (agentKind === "main") {
// Top-level teardown owns the global agent lifecycle: park timers,
// adopted subagent sessions, revivers. Tear it down while shared
// resources (kernels, MCP, LSP) are still live. Subagent disposal
// must NOT touch the global lifecycle.
await AgentLifecycleManager.global().dispose();
}
await originalDispose();
} finally {
agentRegistry.unregister(resolvedAgentId);
unregisterUnlessParked();
unsubscribeCredentialDisabled?.();
}
};
@@ -2502,7 +2522,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
if (hasSession) {
await session.dispose();
} else {
if (hasRegistered) agentRegistry.unregister(resolvedAgentId);
if (hasRegistered) unregisterUnlessParked();
if (asyncJobManager) {
if (AsyncJobManager.instance() === asyncJobManager) {
AsyncJobManager.setInstance(undefined);
@@ -0,0 +1,234 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
import { AgentLifecycleManager } from "@oh-my-pi/pi-coding-agent/registry/agent-lifecycle";
import { AgentRegistry, MAIN_AGENT_ID } from "@oh-my-pi/pi-coding-agent/registry/agent-registry";
import type { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
interface SessionStub {
session: AgentSession;
disposeCalls: () => number;
}
/** Minimal session: the lifecycle manager only ever calls dispose() on it. */
function makeSessionStub(dispose?: () => Promise<void>): SessionStub {
let calls = 0;
const stub = {
dispose: async () => {
calls++;
await dispose?.();
},
};
return { session: stub as unknown as AgentSession, disposeCalls: () => calls };
}
function deferred(): { promise: Promise<void>; resolve: () => void } {
let resolve!: () => void;
const promise = new Promise<void>(r => {
resolve = r;
});
return { promise, resolve };
}
/** Settle the async park chain (timer callback → park() → dispose → setStatus). */
async function flushAsync(): Promise<void> {
for (let i = 0; i < 5; i++) await Promise.resolve();
}
const TTL = 20;
describe("AgentLifecycleManager", () => {
let registry: AgentRegistry;
let lifecycle: AgentLifecycleManager;
beforeEach(() => {
AgentRegistry.resetGlobalForTests();
AgentLifecycleManager.resetGlobalForTests();
registry = AgentRegistry.global();
lifecycle = AgentLifecycleManager.global();
});
afterEach(() => {
vi.useRealTimers();
vi.restoreAllMocks();
AgentLifecycleManager.resetGlobalForTests();
AgentRegistry.resetGlobalForTests();
});
function registerIdleSub(id: string, session: AgentSession | null, sessionFile: string | null = `/tmp/${id}.jsonl`) {
return registry.register({ id, displayName: "task", kind: "sub", session, sessionFile, status: "idle" });
}
it("adopt arms the TTL: an idle agent is parked — session disposed, ref + sessionFile retained", async () => {
vi.useFakeTimers();
const stub = makeSessionStub();
registerIdleSub("1-Sub", stub.session, "/tmp/1-Sub.jsonl");
lifecycle.adopt("1-Sub", { idleTtlMs: TTL });
vi.advanceTimersByTime(TTL);
await flushAsync();
const ref = registry.get("1-Sub");
expect(stub.disposeCalls()).toBe(1);
expect(ref?.status).toBe("parked");
expect(ref?.session).toBeNull();
expect(ref?.sessionFile).toBe("/tmp/1-Sub.jsonl");
expect(lifecycle.has("1-Sub")).toBe(true);
});
it("running disarms the timer; returning to idle re-arms a fresh TTL", async () => {
vi.useFakeTimers();
const stub = makeSessionStub();
registerIdleSub("2-Sub", stub.session);
lifecycle.adopt("2-Sub", { idleTtlMs: TTL });
registry.setStatus("2-Sub", "running");
vi.advanceTimersByTime(TTL * 10);
await flushAsync();
expect(registry.get("2-Sub")?.status).toBe("running");
expect(registry.get("2-Sub")?.session).toBe(stub.session);
expect(stub.disposeCalls()).toBe(0);
registry.setStatus("2-Sub", "idle");
vi.advanceTimersByTime(TTL);
await flushAsync();
expect(registry.get("2-Sub")?.status).toBe("parked");
expect(stub.disposeCalls()).toBe(1);
});
it("ensureLive revives a parked agent through its reviver and flips it back to idle", async () => {
const revived = makeSessionStub();
registry.register({
id: "3-Sub",
displayName: "task",
kind: "sub",
session: null,
sessionFile: "/tmp/3-Sub.jsonl",
status: "parked",
});
lifecycle.adopt("3-Sub", { idleTtlMs: 0, revive: async () => revived.session });
const session = await lifecycle.ensureLive("3-Sub");
expect(session).toBe(revived.session);
const ref = registry.get("3-Sub");
expect(ref?.status).toBe("idle");
expect(ref?.session).toBe(revived.session);
expect(ref?.sessionFile).toBe("/tmp/3-Sub.jsonl");
});
it("concurrent ensureLive calls during a slow revive coalesce into one reviver run", async () => {
const gate = deferred();
const revived = makeSessionStub();
let reviverRuns = 0;
registry.register({
id: "4-Sub",
displayName: "task",
kind: "sub",
session: null,
sessionFile: "/tmp/4-Sub.jsonl",
status: "parked",
});
lifecycle.adopt("4-Sub", {
idleTtlMs: 0,
revive: async () => {
reviverRuns++;
await gate.promise;
return revived.session;
},
});
const first = lifecycle.ensureLive("4-Sub");
const second = lifecycle.ensureLive("4-Sub");
gate.resolve();
const [a, b] = await Promise.all([first, second]);
expect(reviverRuns).toBe(1);
expect(a).toBe(revived.session);
expect(b).toBe(revived.session);
});
it("ensureLive on an unknown id throws and points at history://", async () => {
await expect(lifecycle.ensureLive("9-Ghost")).rejects.toThrow(/history:\/\/9-Ghost/);
});
it("ensureLive on a parked agent without a reviver throws as not revivable", async () => {
registry.register({ id: "5-Sub", displayName: "task", kind: "sub", session: null, status: "parked" });
lifecycle.adopt("5-Sub", { idleTtlMs: 0 });
await expect(lifecycle.ensureLive("5-Sub")).rejects.toThrow(/cannot be revived.*no reviver registered/);
});
it("release disposes a live adopted agent, unregisters it, and leaves no pending park", async () => {
vi.useFakeTimers();
const stub = makeSessionStub();
registerIdleSub("6-Sub", stub.session);
lifecycle.adopt("6-Sub", { idleTtlMs: TTL });
await lifecycle.release("6-Sub");
expect(stub.disposeCalls()).toBe(1);
expect(registry.get("6-Sub")).toBeUndefined();
expect(lifecycle.has("6-Sub")).toBe(false);
// The disarmed timer must not fire a late park (which would double-dispose).
vi.advanceTimersByTime(TTL * 10);
await flushAsync();
expect(stub.disposeCalls()).toBe(1);
expect(registry.get("6-Sub")).toBeUndefined();
});
it("adopt(Main) is a no-op: Main is never adopted or parked", async () => {
vi.useFakeTimers();
const stub = makeSessionStub();
registry.register({
id: MAIN_AGENT_ID,
displayName: "main",
kind: "main",
session: stub.session,
status: "idle",
});
lifecycle.adopt(MAIN_AGENT_ID, { idleTtlMs: TTL });
expect(lifecycle.has(MAIN_AGENT_ID)).toBe(false);
vi.advanceTimersByTime(TTL * 10);
await flushAsync();
expect(registry.get(MAIN_AGENT_ID)?.status).toBe("idle");
expect(registry.get(MAIN_AGENT_ID)?.session).toBe(stub.session);
expect(stub.disposeCalls()).toBe(0);
});
it("isParking is true exactly while park's dispose is in flight; parked only after it completes", async () => {
const gate = deferred();
const stub = makeSessionStub(() => gate.promise);
registerIdleSub("7-Sub", stub.session);
lifecycle.adopt("7-Sub", { idleTtlMs: 0 });
// park() runs synchronously up to `await session.dispose()`, which we hold open.
const parking = lifecycle.park("7-Sub");
expect(stub.disposeCalls()).toBe(1);
expect(lifecycle.isParking("7-Sub")).toBe(true);
expect(registry.get("7-Sub")).toBeDefined();
expect(registry.get("7-Sub")?.status).toBe("idle"); // not yet flipped
gate.resolve();
await parking;
expect(lifecycle.isParking("7-Sub")).toBe(false);
expect(registry.get("7-Sub")?.status).toBe("parked");
expect(registry.get("7-Sub")?.session).toBeNull();
});
it("idleTtlMs <= 0 adopts without a timer: the agent never parks", async () => {
vi.useFakeTimers();
const stub = makeSessionStub();
registerIdleSub("8-Sub", stub.session);
lifecycle.adopt("8-Sub", { idleTtlMs: 0 });
vi.advanceTimersByTime(60_000);
await flushAsync();
const ref = registry.get("8-Sub");
expect(ref?.status).toBe("idle");
expect(ref?.session).toBe(stub.session);
expect(stub.disposeCalls()).toBe(0);
expect(lifecycle.has("8-Sub")).toBe(true);
});
});