From e143495d6c0dce8f03576ed276aba766196ea1d6 Mon Sep 17 00:00:00 2001 From: alexion Date: Sat, 1 Aug 2026 16:47:34 -0400 Subject: [PATCH] feat(pi): expire recent subagent statuses --- .../pi/extensions/subagents/config.test.ts | 25 ++++++- .../agents/pi/extensions/subagents/config.ts | 6 ++ .../agents/pi/extensions/subagents/index.ts | 1 + .../extensions/subagents/supervisor.test.ts | 70 +++++++++++++++++++ .../pi/extensions/subagents/supervisor.ts | 38 ++++++++++ 5 files changed, 138 insertions(+), 2 deletions(-) diff --git a/modules/agents/pi/extensions/subagents/config.test.ts b/modules/agents/pi/extensions/subagents/config.test.ts index e292df3..679cc91 100644 --- a/modules/agents/pi/extensions/subagents/config.test.ts +++ b/modules/agents/pi/extensions/subagents/config.test.ts @@ -28,6 +28,7 @@ test("missing config files and agent directories are normal", () => { assert.equal(config.defaultContext, "independent"); assert.equal(config.defaultTools, "read-only"); + assert.equal(config.recentTerminalTtlMs, 300000); assert.equal(agents.size, 0); assert.deepEqual(diag.warnings, []); }); @@ -35,16 +36,36 @@ test("missing config files and agent directories are normal", () => { test("global and trusted project config merge in order", () => { const { cwd, agentDir } = fixture(); mkdirSync(join(cwd, ".pi"), { recursive: true }); - writeFileSync(join(agentDir, "subagents.json"), JSON.stringify({ defaultTools: "global-profile", toolProfiles: { "global-profile": { activeTools: ["read"] } } })); - writeFileSync(join(cwd, ".pi", "subagents.json"), JSON.stringify({ defaultTools: "project-profile", toolProfiles: { "project-profile": { activeTools: ["ls"] } } })); + writeFileSync(join(agentDir, "subagents.json"), JSON.stringify({ defaultTools: "global-profile", recentTerminalTtlMs: 1000, toolProfiles: { "global-profile": { activeTools: ["read"] } } })); + writeFileSync(join(cwd, ".pi", "subagents.json"), JSON.stringify({ defaultTools: "project-profile", recentTerminalTtlMs: 2000, toolProfiles: { "project-profile": { activeTools: ["ls"] } } })); const config = loadConfig(cwd, true, diagnostics(), agentDir); assert.equal(config.defaultTools, "project-profile"); + assert.equal(config.recentTerminalTtlMs, 2000); assert.deepEqual(config.toolProfiles["global-profile"].activeTools, ["read"]); assert.deepEqual(config.toolProfiles["project-profile"].activeTools, ["ls"]); }); +test("recent terminal ttl preserves zero and rejects invalid values", () => { + const { cwd, agentDir } = fixture(); + writeFileSync(join(agentDir, "subagents.json"), JSON.stringify({ recentTerminalTtlMs: 0 })); + const zeroDiag = diagnostics(); + + const zeroConfig = loadConfig(cwd, true, zeroDiag, agentDir); + + assert.equal(zeroConfig.recentTerminalTtlMs, 0); + assert.deepEqual(zeroDiag.warnings, []); + + writeFileSync(join(agentDir, "subagents.json"), JSON.stringify({ recentTerminalTtlMs: -1 })); + const invalidDiag = diagnostics(); + + const invalidConfig = loadConfig(cwd, true, invalidDiag, agentDir); + + assert.equal(invalidConfig.recentTerminalTtlMs, 300000); + assert.ok(invalidDiag.warnings.some((warning) => warning.includes("Invalid global recentTerminalTtlMs ignored"))); +}); + test("project config is ignored when project is untrusted", () => { const { cwd, agentDir } = fixture(); mkdirSync(join(cwd, ".pi"), { recursive: true }); diff --git a/modules/agents/pi/extensions/subagents/config.ts b/modules/agents/pi/extensions/subagents/config.ts index 3fc6cf1..36d2785 100644 --- a/modules/agents/pi/extensions/subagents/config.ts +++ b/modules/agents/pi/extensions/subagents/config.ts @@ -12,6 +12,7 @@ export interface SubagentsConfig { defaultContext: ContextMode; defaultTools: string; maxConcurrent: number; + recentTerminalTtlMs: number; ui: { enabled: boolean; defaultExpanded: boolean; @@ -38,6 +39,7 @@ const DEFAULT_CONFIG: SubagentsConfig = { defaultContext: "independent", defaultTools: "read-only", maxConcurrent: 3, + recentTerminalTtlMs: 5 * 60 * 1000, ui: { enabled: true, defaultExpanded: false }, toolProfiles: { ...BUILT_IN_TOOL_PROFILES }, }; @@ -108,6 +110,9 @@ function normalizeConfig(raw: unknown, diagnostics: Diagnostics, label: string): else if (input.defaultTools !== undefined) diagnostics.warnings.push(`Invalid ${label} defaultTools ignored`); if (typeof input.maxConcurrent === "number" && Number.isInteger(input.maxConcurrent) && input.maxConcurrent > 0) config.maxConcurrent = input.maxConcurrent; else if (input.maxConcurrent !== undefined) diagnostics.warnings.push(`Invalid ${label} maxConcurrent ignored`); + if (typeof input.recentTerminalTtlMs === "number" && Number.isInteger(input.recentTerminalTtlMs) && input.recentTerminalTtlMs >= 0) { + config.recentTerminalTtlMs = input.recentTerminalTtlMs; + } else if (input.recentTerminalTtlMs !== undefined) diagnostics.warnings.push(`Invalid ${label} recentTerminalTtlMs ignored`); if (input.ui !== undefined) config.ui = normalizeUi(input.ui, diagnostics, label); if (input.toolProfiles !== undefined) config.toolProfiles = normalizeProfiles(input.toolProfiles, diagnostics, label); return config; @@ -159,6 +164,7 @@ function mergeConfig(base: SubagentsConfig, override: Partial | if (override.defaultContext) merged.defaultContext = override.defaultContext; if (override.defaultTools) merged.defaultTools = override.defaultTools; if (override.maxConcurrent) merged.maxConcurrent = override.maxConcurrent; + if (override.recentTerminalTtlMs !== undefined) merged.recentTerminalTtlMs = override.recentTerminalTtlMs; if (override.ui) merged.ui = { ...merged.ui, ...override.ui }; if (override.toolProfiles) merged.toolProfiles = { ...merged.toolProfiles, ...override.toolProfiles }; for (const key of Object.keys(merged.toolProfiles)) { diff --git a/modules/agents/pi/extensions/subagents/index.ts b/modules/agents/pi/extensions/subagents/index.ts index 5ea2958..bb194e1 100644 --- a/modules/agents/pi/extensions/subagents/index.ts +++ b/modules/agents/pi/extensions/subagents/index.ts @@ -22,6 +22,7 @@ export default function subagents(pi: ExtensionAPI) { uiExpanded = config.ui.defaultExpanded; supervisor = new Supervisor(new SubprocessRpcRunner(), cwd, { maxConcurrent: config.maxConcurrent, + recentTerminalTtlMs: config.recentTerminalTtlMs, onMilestone: (status, event) => pi.appendEntry("subagent_milestone", { event, status }), onChange: (statuses) => { lastStatuses = statuses; diff --git a/modules/agents/pi/extensions/subagents/supervisor.test.ts b/modules/agents/pi/extensions/subagents/supervisor.test.ts index 7cff49e..edaca0c 100644 --- a/modules/agents/pi/extensions/subagents/supervisor.test.ts +++ b/modules/agents/pi/extensions/subagents/supervisor.test.ts @@ -110,6 +110,24 @@ test("completed children ignore later cancel", async () => { assert.equal(runner.starts[0].handle.cancelCalls, 0); }); +test("shutdown clears recent terminal expiry timer", async () => { + const runner = new FakeRunner(); + let changes = 0; + const supervisor = new Supervisor(runner, "/tmp", { + recentTerminalTtlMs: 5, + onChange: () => { + changes += 1; + }, + }); + await spawnStarted(supervisor); + + await supervisor.shutdown(); + const afterShutdown = changes; + await sleep(15); + + assert.equal(changes, afterShutdown); +}); + test("batch spawn returns accepted ids and per-entry failures", async () => { const runner = new FakeRunner(); const supervisor = new Supervisor(runner, "/tmp"); @@ -140,6 +158,58 @@ test("maxConcurrent preserves queued records", async () => { assert.equal(runner.starts.length, 2); }); +test("recent terminal statuses expire from list by ttl", async () => { + const runner = new FakeRunner(); + const supervisor = new Supervisor(runner, "/tmp", { recentTerminalTtlMs: 5 }); + const accepted = await spawnStarted(supervisor); + + runner.starts[0].events.completed("done", "agent_settled"); + assert.equal(supervisor.list().some((status) => status.id === accepted.id), true); + + await sleep(10); + + assert.equal(supervisor.list().some((status) => status.id === accepted.id), false); + assert.equal(supervisor.result(accepted.id).result, "done"); +}); + +test("recent terminal ttl does not hide active statuses", async () => { + const runner = new FakeRunner(); + const supervisor = new Supervisor(runner, "/tmp", { recentTerminalTtlMs: 0 }); + const accepted = await spawnStarted(supervisor); + + assert.equal(supervisor.list().some((status) => status.id === accepted.id), true); +}); + +test("zero recent terminal ttl hides terminal statuses immediately", async () => { + const runner = new FakeRunner(); + const supervisor = new Supervisor(runner, "/tmp", { recentTerminalTtlMs: 0 }); + const accepted = await spawnStarted(supervisor); + + runner.starts[0].events.completed("done", "agent_settled"); + + assert.equal(supervisor.list().some((status) => status.id === accepted.id), false); + assert.equal(supervisor.result(accepted.id).result, "done"); +}); + +test("recent terminal ttl emits a change when an entry expires", async () => { + const runner = new FakeRunner(); + let changes = 0; + const supervisor = new Supervisor(runner, "/tmp", { + recentTerminalTtlMs: 5, + onChange: () => { + changes += 1; + }, + }); + await spawnStarted(supervisor); + const beforeComplete = changes; + + runner.starts[0].events.completed("done", "agent_settled"); + await sleep(15); + + assert.ok(changes > beforeComplete + 1); + assert.equal(supervisor.list().length, 0); +}); + test("wait blocks until multiple subagents are terminal", async () => { const runner = new FakeRunner(); const supervisor = new Supervisor(runner, "/tmp"); diff --git a/modules/agents/pi/extensions/subagents/supervisor.ts b/modules/agents/pi/extensions/subagents/supervisor.ts index 743bfde..702aa46 100644 --- a/modules/agents/pi/extensions/subagents/supervisor.ts +++ b/modules/agents/pi/extensions/subagents/supervisor.ts @@ -24,6 +24,7 @@ interface RunningChild { interface SupervisorOptions { maxConcurrent?: number; recentTerminalLimit?: number; + recentTerminalTtlMs?: number; timeouts?: { startMs?: number; runMs?: number; @@ -47,6 +48,7 @@ export class Supervisor { private readonly children = new Map(); private readonly queue: RunningChild[] = []; private readonly waiters = new Set<() => void>(); + private recentTerminalTimer?: ReturnType; constructor( private readonly runner: ChildRunner, @@ -76,6 +78,7 @@ export class Supervisor { const active = statuses.filter((status) => !isTerminal(status.state)); const terminal = statuses .filter((status) => isTerminal(status.state)) + .filter((status) => this.isRecentTerminal(status)) .sort((a, b) => Date.parse(b.completedAt ?? b.startedAt) - Date.parse(a.completedAt ?? a.startedAt)) .slice(0, this.options.recentTerminalLimit ?? 10); return [...active, ...terminal]; @@ -134,6 +137,7 @@ export class Supervisor { } async shutdown(): Promise { + this.clearRecentTerminalTimer(); await Promise.allSettled( [...this.children.values()].map(async (child) => { if (!isTerminal(child.record.status.state)) { @@ -142,6 +146,7 @@ export class Supervisor { } }), ); + this.clearRecentTerminalTimer(); } private createChild(request: SpawnRequest): SpawnAccepted { @@ -340,6 +345,39 @@ export class Supervisor { private emitChange() { this.options.onChange?.(this.list()); for (const waiter of this.waiters) waiter(); + this.scheduleRecentTerminalExpiry(); + } + + private scheduleRecentTerminalExpiry() { + this.clearRecentTerminalTimer(); + const ttl = this.options.recentTerminalTtlMs; + if (ttl === undefined || ttl <= 0) return; + const now = Date.now(); + const nextExpiryMs = [...this.children.values()] + .map((child) => child.record.status) + .filter((status) => isTerminal(status.state)) + .map((status) => Date.parse(status.completedAt ?? status.startedAt)) + .filter((completed) => Number.isFinite(completed)) + .map((completed) => completed + ttl - now) + .filter((remaining) => remaining > 0) + .sort((a, b) => a - b)[0]; + if (nextExpiryMs === undefined) return; + this.recentTerminalTimer = setTimeout(() => this.emitChange(), nextExpiryMs + 1); + } + + private clearRecentTerminalTimer() { + if (!this.recentTerminalTimer) return; + clearTimeout(this.recentTerminalTimer); + this.recentTerminalTimer = undefined; + } + + private isRecentTerminal(status: SubagentStatus): boolean { + const ttl = this.options.recentTerminalTtlMs; + if (ttl === undefined) return true; + if (ttl <= 0) return false; + const completed = Date.parse(status.completedAt ?? status.startedAt); + if (!Number.isFinite(completed)) return true; + return Date.now() - completed <= ttl; } private waitReady(ids: string[], mode: SubagentWaitMode): boolean {