diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index ad0fbffd..1aee83ad 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -6,6 +6,7 @@ - Added `/login ` support with provider autocomplete. - Added public SDK exports for CLI-equivalent model and scoped-model resolution ([#6201](https://github.com/earendil-works/pi/issues/6201)). +- Added extension and RPC `agent_settled` events plus session-level idle waiting for fully settled agent runs ([#6363](https://github.com/earendil-works/pi/issues/6363)). - Added extension entry renderers for persisted display-only session entries that are rendered in interactive mode without being sent to the model context. ### Fixed diff --git a/packages/coding-agent/docs/extensions.md b/packages/coding-agent/docs/extensions.md index ddfbec04..d14c1685 100644 --- a/packages/coding-agent/docs/extensions.md +++ b/packages/coding-agent/docs/extensions.md @@ -307,7 +307,8 @@ user sends prompt ──────────────────── │ │ │ │ │ └─► turn_end │ │ │ │ - └─► agent_end │ + ├─► agent_end │ + └─► agent_settled (no retry/compaction/follow-up left) │ │ user sends another prompt ◄────────────────────────────────┘ @@ -546,15 +547,19 @@ The `systemPromptOptions` field gives extensions access to the same structured d Inside `before_agent_start`, `event.systemPrompt` and `ctx.getSystemPrompt()` both reflect the chained system prompt as of the current handler. Later `before_agent_start` handlers can still modify it again. -#### agent_start / agent_end +#### agent_start / agent_end / agent_settled -Fired once per user prompt. +`agent_start` fires when a low-level agent run begins. `agent_end` fires when that run ends, but Pi may still auto-retry, auto-compact and retry, or continue with queued follow-up messages. Use `agent_settled` for status integrations that need to know Pi will not continue running automatically. ```typescript pi.on("agent_start", async (_event, ctx) => {}); pi.on("agent_end", async (event, ctx) => { - // event.messages - messages from this prompt + // event.messages - messages from this low-level run +}); + +pi.on("agent_settled", async (_event, ctx) => { + // ctx.isIdle() is true here unless another extension started a new run. }); ``` @@ -1000,7 +1005,7 @@ pi.on("tool_result", async (event, ctx) => { ### ctx.isIdle() / ctx.abort() / ctx.hasPendingMessages() -Control flow helpers. +Control flow helpers. `ctx.isIdle()` is false while Pi is processing an agent run, automatic retry, auto-compaction retry, or queued continuation. ### ctx.shutdown() @@ -1082,7 +1087,7 @@ This reports the current base prompt inputs. It does not include per-turn `befor ### ctx.waitForIdle() -Wait for the agent to finish streaming: +Wait for the agent to fully settle, including automatic retries, auto-compaction retries, and queued continuations: ```typescript pi.registerCommand("my-cmd", { diff --git a/packages/coding-agent/docs/rpc.md b/packages/coding-agent/docs/rpc.md index bf6cf130..89c1c72f 100644 --- a/packages/coding-agent/docs/rpc.md +++ b/packages/coding-agent/docs/rpc.md @@ -808,7 +808,8 @@ Events are streamed to stdout as JSON lines during agent operation. Events do NO | Event | Description | |-------|-------------| | `agent_start` | Agent begins processing | -| `agent_end` | Agent completes (includes all generated messages) | +| `agent_end` | One low-level agent run completes (may still be followed by retry, compaction, or queued continuations) | +| `agent_settled` | Agent run is fully settled; no automatic retry, compaction retry, or queued continuation remains | | `turn_start` | New turn begins | | `turn_end` | Turn completes (includes assistant message and tool results) | | `message_start` | Message begins | @@ -834,15 +835,24 @@ Emitted when the agent begins processing a prompt. ### agent_end -Emitted when the agent completes. Contains all messages generated during this run. +Emitted when one low-level agent run completes. Contains all messages generated during this run. If `willRetry` is true, an automatic retry will follow. ```json { "type": "agent_end", - "messages": [...] + "messages": [...], + "willRetry": false } ``` +### agent_settled + +Emitted after the full session-level run settles. At this point Pi will not continue automatically through retry, compaction retry, or queued follow-up messages. + +```json +{"type": "agent_settled"} +``` + ### turn_start / turn_end A turn consists of one assistant response plus any resulting tool calls and results. diff --git a/packages/coding-agent/src/core/agent-session.ts b/packages/coding-agent/src/core/agent-session.ts index d0fbed7c..9a3a5c5f 100644 --- a/packages/coding-agent/src/core/agent-session.ts +++ b/packages/coding-agent/src/core/agent-session.ts @@ -131,6 +131,7 @@ export type AgentSessionEvent = messages: AgentMessage[]; willRetry: boolean; } + | { type: "agent_settled" } | { type: "queue_update"; steering: readonly string[]; @@ -275,6 +276,9 @@ export class AgentSession { // Event subscription state private _unsubscribeAgent?: () => void; private _eventListeners: AgentSessionEventListener[] = []; + private _isAgentRunActive = false; + private _idleWaitPromise: Promise | undefined; + private _resolveIdleWait: (() => void) | undefined; /** Tracks pending steering messages for UI display. Removed when delivered. */ private _steeringMessages: string[] = []; @@ -508,6 +512,35 @@ export class AgentSession { }); } + private _getIdleWaitPromise(): Promise { + if (!this._idleWaitPromise) { + this._idleWaitPromise = new Promise((resolve) => { + this._resolveIdleWait = resolve; + }); + } + return this._idleWaitPromise; + } + + private _resolveIdleWaitIfIdle(): void { + if (this._isAgentRunActive || !this._resolveIdleWait) { + return; + } + const resolve = this._resolveIdleWait; + this._idleWaitPromise = undefined; + this._resolveIdleWait = undefined; + resolve(); + } + + private async _emitAgentSettled(): Promise { + this._isAgentRunActive = false; + try { + await this._extensionRunner.emit({ type: "agent_settled" }); + this._emit({ type: "agent_settled" }); + } finally { + this._resolveIdleWaitIfIdle(); + } + } + // Track last assistant message for auto-compaction check private _lastAssistantMessage: AssistantMessage | undefined = undefined; @@ -801,9 +834,14 @@ export class AgentSession { return this.agent.state.thinkingLevel; } - /** Whether agent is currently streaming a response */ + /** Whether the session is currently processing an agent run or post-run continuation. */ get isStreaming(): boolean { - return this.agent.state.isStreaming; + return this._isAgentRunActive; + } + + /** Whether the session has no active agent run, retry, auto-compaction, or queued continuation. */ + get isIdle(): boolean { + return !this._isAgentRunActive; } /** Current effective system prompt (includes any per-turn extension modifications) */ @@ -983,6 +1021,7 @@ export class AgentSession { // ========================================================================= private async _runAgentPrompt(messages: AgentMessage | AgentMessage[]): Promise { + this._isAgentRunActive = true; try { await this.agent.prompt(messages); while (await this._handlePostAgentRun()) { @@ -991,6 +1030,7 @@ export class AgentSession { } finally { this._systemPromptOverride = undefined; this._flushPendingBashMessages(); + await this._emitAgentSettled(); } } @@ -1461,7 +1501,14 @@ export class AgentSession { async abort(): Promise { this.abortRetry(); this.agent.abort(); - await this.agent.waitForIdle(); + await this.waitForIdle(); + } + + async waitForIdle(): Promise { + if (this.isIdle) { + return; + } + await this._getIdleWaitPromise(); } // ========================================================================= @@ -2305,7 +2352,7 @@ export class AgentSession { }, { getModel: () => this.model, - isIdle: () => !this.isStreaming, + isIdle: () => this.isIdle, isProjectTrusted: () => this.settingsManager.isProjectTrusted(), getSignal: () => this.agent.signal, abort: () => { diff --git a/packages/coding-agent/src/core/extensions/index.ts b/packages/coding-agent/src/core/extensions/index.ts index cbe772f1..29847b63 100644 --- a/packages/coding-agent/src/core/extensions/index.ts +++ b/packages/coding-agent/src/core/extensions/index.ts @@ -22,6 +22,7 @@ export { ExtensionRunner } from "./runner.ts"; export type { AfterProviderResponseEvent, AgentEndEvent, + AgentSettledEvent, AgentStartEvent, // Re-exports AgentToolResult, diff --git a/packages/coding-agent/src/core/extensions/types.ts b/packages/coding-agent/src/core/extensions/types.ts index 4a04c261..9a37d8c3 100644 --- a/packages/coding-agent/src/core/extensions/types.ts +++ b/packages/coding-agent/src/core/extensions/types.ts @@ -705,6 +705,11 @@ export interface AgentEndEvent { messages: AgentMessage[]; } +/** Fired after an agent run has fully settled and no automatic retry, compaction, or queued continuation will run. */ +export interface AgentSettledEvent { + type: "agent_settled"; +} + /** Fired at the start of each turn */ export interface TurnStartEvent { type: "turn_start"; @@ -1021,6 +1026,7 @@ export type ExtensionEvent = | BeforeAgentStartEvent | AgentStartEvent | AgentEndEvent + | AgentSettledEvent | TurnStartEvent | TurnEndEvent | MessageStartEvent @@ -1188,6 +1194,7 @@ export interface ExtensionAPI { on(event: "before_agent_start", handler: ExtensionHandler): void; on(event: "agent_start", handler: ExtensionHandler): void; on(event: "agent_end", handler: ExtensionHandler): void; + on(event: "agent_settled", handler: ExtensionHandler): void; on(event: "turn_start", handler: ExtensionHandler): void; on(event: "turn_end", handler: ExtensionHandler): void; on(event: "message_start", handler: ExtensionHandler): void; diff --git a/packages/coding-agent/src/core/index.ts b/packages/coding-agent/src/core/index.ts index a63a471f..924932e6 100644 --- a/packages/coding-agent/src/core/index.ts +++ b/packages/coding-agent/src/core/index.ts @@ -32,6 +32,7 @@ export { areExperimentalFeaturesEnabled } from "./experimental.ts"; // Extensions system export { type AgentEndEvent, + type AgentSettledEvent, type AgentStartEvent, type AgentToolResult, type AgentToolUpdateCallback, diff --git a/packages/coding-agent/src/index.ts b/packages/coding-agent/src/index.ts index 29c82175..7c6de94b 100644 --- a/packages/coding-agent/src/index.ts +++ b/packages/coding-agent/src/index.ts @@ -61,6 +61,7 @@ export { createEventBus, type EventBus, type EventBusController } from "./core/e // Extension system export type { AgentEndEvent, + AgentSettledEvent, AgentStartEvent, AgentToolResult, AgentToolUpdateCallback, diff --git a/packages/coding-agent/src/modes/interactive/interactive-mode.ts b/packages/coding-agent/src/modes/interactive/interactive-mode.ts index d1fe0d84..8773cc0c 100644 --- a/packages/coding-agent/src/modes/interactive/interactive-mode.ts +++ b/packages/coding-agent/src/modes/interactive/interactive-mode.ts @@ -1624,7 +1624,7 @@ export class InteractiveMode { this.restoreQueuedMessagesToEditor({ abort: true }); }, commandContextActions: { - waitForIdle: () => this.session.agent.waitForIdle(), + waitForIdle: () => this.session.waitForIdle(), newSession: async (options) => { this.clearStatusIndicator(); try { @@ -1674,7 +1674,7 @@ export class InteractiveMode { }, shutdownHandler: () => { this.shutdownRequested = true; - if (!this.session.isStreaming) { + if (this.session.isIdle) { void this.shutdown(); } }, @@ -1774,7 +1774,7 @@ export class InteractiveMode { sessionManager: this.sessionManager, modelRegistry: this.session.modelRegistry, model: this.session.model, - isIdle: () => !this.session.isStreaming, + isIdle: () => this.session.isIdle, isProjectTrusted: () => this.settingsManager.isProjectTrusted(), signal: this.session.agent.signal, abort: () => { @@ -3024,11 +3024,13 @@ export class InteractiveMode { } this.pendingTools.clear(); - await this.checkShutdownRequested(); - this.ui.requestRender(); break; + case "agent_settled": + await this.checkShutdownRequested(); + break; + case "compaction_start": { if (this.settingsManager.getShowTerminalProgress()) { this.ui.terminal.setProgress(true); diff --git a/packages/coding-agent/src/modes/print-mode.ts b/packages/coding-agent/src/modes/print-mode.ts index a5ee0235..ca63157b 100644 --- a/packages/coding-agent/src/modes/print-mode.ts +++ b/packages/coding-agent/src/modes/print-mode.ts @@ -73,7 +73,7 @@ export async function runPrintMode(runtimeHost: AgentSessionRuntime, options: Pr await session.bindExtensions({ mode: mode === "json" ? "json" : "print", commandContextActions: { - waitForIdle: () => session.agent.waitForIdle(), + waitForIdle: () => session.waitForIdle(), newSession: async (newSessionOptions) => runtimeHost.newSession(newSessionOptions), fork: async (entryId, forkOptions) => { const result = await runtimeHost.fork(entryId, forkOptions); diff --git a/packages/coding-agent/src/modes/rpc/rpc-client.ts b/packages/coding-agent/src/modes/rpc/rpc-client.ts index eca52af7..934c6f1f 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-client.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-client.ts @@ -5,9 +5,9 @@ */ import { type ChildProcess, spawn } from "node:child_process"; -import type { AgentEvent, AgentMessage, ThinkingLevel } from "@earendil-works/pi-agent-core"; +import type { AgentMessage, ThinkingLevel } from "@earendil-works/pi-agent-core"; import type { ImageContent } from "@earendil-works/pi-ai"; -import type { SessionStats } from "../../core/agent-session.ts"; +import type { AgentSessionEvent, SessionStats } from "../../core/agent-session.ts"; import type { BashResult } from "../../core/bash-executor.ts"; import type { CompactionResult } from "../../core/compaction/index.ts"; import type { SessionEntry, SessionTreeNode } from "../../core/session-manager.ts"; @@ -46,7 +46,7 @@ export interface ModelInfo { reasoning: boolean; } -export type RpcEventListener = (event: AgentEvent) => void; +export type RpcEventListener = (event: AgentSessionEvent) => void; // ============================================================================ // RPC Client @@ -442,7 +442,7 @@ export class RpcClient { /** * Wait for agent to become idle (no streaming). - * Resolves when agent_end event is received. + * Resolves when agent_settled event is received. */ waitForIdle(timeout = 60000): Promise { return new Promise((resolve, reject) => { @@ -452,7 +452,7 @@ export class RpcClient { }, timeout); const unsubscribe = this.onEvent((event) => { - if (event.type === "agent_end") { + if (event.type === "agent_settled") { clearTimeout(timer); unsubscribe(); resolve(); @@ -464,9 +464,9 @@ export class RpcClient { /** * Collect events until agent becomes idle. */ - collectEvents(timeout = 60000): Promise { + collectEvents(timeout = 60000): Promise { return new Promise((resolve, reject) => { - const events: AgentEvent[] = []; + const events: AgentSessionEvent[] = []; const timer = setTimeout(() => { unsubscribe(); reject(new Error(`Timeout collecting events. Stderr: ${this.stderr}`)); @@ -474,7 +474,7 @@ export class RpcClient { const unsubscribe = this.onEvent((event) => { events.push(event); - if (event.type === "agent_end") { + if (event.type === "agent_settled") { clearTimeout(timer); unsubscribe(); resolve(events); @@ -486,7 +486,7 @@ export class RpcClient { /** * Send prompt and wait for completion, returning all events. */ - async promptAndWait(message: string, images?: ImageContent[], timeout = 60000): Promise { + async promptAndWait(message: string, images?: ImageContent[], timeout = 60000): Promise { const eventsPromise = this.collectEvents(timeout); await this.prompt(message, images); return eventsPromise; @@ -510,7 +510,7 @@ export class RpcClient { // Otherwise it's an event for (const listener of this.eventListeners) { - listener(data as AgentEvent); + listener(data as AgentSessionEvent); } } catch { // Ignore non-JSON lines diff --git a/packages/coding-agent/src/modes/rpc/rpc-mode.ts b/packages/coding-agent/src/modes/rpc/rpc-mode.ts index 38cac641..b099f874 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-mode.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-mode.ts @@ -319,7 +319,7 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise session.agent.waitForIdle(), + waitForIdle: () => session.waitForIdle(), newSession: async (options) => runtimeHost.newSession(options), fork: async (entryId, forkOptions) => { const result = await runtimeHost.fork(entryId, forkOptions); @@ -353,6 +353,9 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise { output(event); + if (event.type === "agent_settled") { + void checkShutdownRequested(); + } }); unsubscribeBackpressure = session.agent.subscribe(async () => { await waitForRawStdoutBackpressure(); diff --git a/packages/coding-agent/test/suite/agent-session-retry-events.test.ts b/packages/coding-agent/test/suite/agent-session-retry-events.test.ts index 946f4da9..0e07f82f 100644 --- a/packages/coding-agent/test/suite/agent-session-retry-events.test.ts +++ b/packages/coding-agent/test/suite/agent-session-retry-events.test.ts @@ -253,6 +253,7 @@ describe("AgentSession retry and event characterization", () => { "message_end:assistant", "turn_end", "agent_end", + "agent_settled", ]); }); @@ -298,6 +299,7 @@ describe("AgentSession retry and event characterization", () => { "message_end:assistant", "turn_end", "agent_end", + "agent_settled", ]); }); @@ -328,7 +330,8 @@ describe("AgentSession retry and event characterization", () => { await harness.session.prompt("hi"); - expect(harness.events[harness.events.length - 1]?.type).toBe("agent_end"); + expect(harness.eventsOfType("agent_end")).toHaveLength(1); + expect(harness.events[harness.events.length - 1]?.type).toBe("agent_settled"); }); it("emits agent_end for aborted runs and persists the aborted assistant message", async () => { @@ -350,7 +353,8 @@ describe("AgentSession retry and event characterization", () => { await harness.session.abort(); await promptPromise; - expect(harness.events[harness.events.length - 1]?.type).toBe("agent_end"); + expect(harness.eventsOfType("agent_end")).toHaveLength(1); + expect(harness.events[harness.events.length - 1]?.type).toBe("agent_settled"); const lastMessage = harness.session.messages[harness.session.messages.length - 1]; expect(lastMessage?.role).toBe("assistant"); if (lastMessage?.role === "assistant") { diff --git a/packages/coding-agent/test/suite/regressions/6363-agent-settled-event.test.ts b/packages/coding-agent/test/suite/regressions/6363-agent-settled-event.test.ts new file mode 100644 index 00000000..d3101047 --- /dev/null +++ b/packages/coding-agent/test/suite/regressions/6363-agent-settled-event.test.ts @@ -0,0 +1,158 @@ +import type { AgentTool } from "@earendil-works/pi-agent-core"; +import { fauxAssistantMessage, fauxToolCall } from "@earendil-works/pi-ai"; +import { Type } from "typebox"; +import { afterEach, describe, expect, it } from "vitest"; +import { createHarness, getUserTexts, type Harness } from "../harness.ts"; + +function createWaitTool(released: Promise): AgentTool { + return { + name: "wait", + label: "Wait", + description: "Wait until released", + parameters: Type.Object({}), + execute: async () => { + await released; + return { content: [{ type: "text", text: "released" }], details: {} }; + }, + }; +} + +describe("regression #6363: agent settled event and idle waiting", () => { + const harnesses: Harness[] = []; + + afterEach(() => { + while (harnesses.length > 0) { + harnesses.pop()?.cleanup(); + } + }); + + it("emits one agent_settled event after automatic retry finishes", async () => { + const extensionEvents: string[] = []; + const publicEvents: string[] = []; + const harness = await createHarness({ + settings: { retry: { enabled: true, maxRetries: 3, baseDelayMs: 1 } }, + extensionFactories: [ + (pi) => { + pi.on("agent_end", () => { + extensionEvents.push("agent_end"); + }); + pi.on("agent_settled", (_event, ctx) => { + extensionEvents.push(`agent_settled:${ctx.isIdle()}`); + }); + }, + ], + }); + harnesses.push(harness); + harness.session.subscribe((event) => { + if (event.type === "agent_settled") { + publicEvents.push("agent_settled"); + } + }); + harness.setResponses([ + fauxAssistantMessage("", { stopReason: "error", errorMessage: "overloaded_error" }), + fauxAssistantMessage("recovered"), + ]); + + await harness.session.prompt("test"); + + expect(harness.eventsOfType("agent_end").map((event) => event.willRetry)).toEqual([true, false]); + expect(harness.eventsOfType("agent_settled")).toHaveLength(1); + expect(extensionEvents).toEqual(["agent_end", "agent_end", "agent_settled:true"]); + expect(publicEvents).toEqual(["agent_settled"]); + }); + + it("settles only after follow-ups queued by agent_end handlers run", async () => { + let queuedFollowUp = false; + const settledIdleStates: boolean[] = []; + const harness = await createHarness({ + extensionFactories: [ + (pi) => { + pi.on("agent_end", () => { + if (queuedFollowUp) return; + queuedFollowUp = true; + pi.sendUserMessage("status follow-up", { deliverAs: "followUp" }); + }); + pi.on("agent_settled", (_event, ctx) => { + settledIdleStates.push(ctx.isIdle()); + }); + }, + ], + }); + harnesses.push(harness); + harness.setResponses([fauxAssistantMessage("first"), fauxAssistantMessage("second")]); + + await harness.session.prompt("hello"); + + expect(getUserTexts(harness)).toEqual(["hello", "status follow-up"]); + expect(harness.eventsOfType("agent_end")).toHaveLength(2); + expect(harness.eventsOfType("agent_settled")).toHaveLength(1); + expect(settledIdleStates).toEqual([true]); + }); + + it("extension command waitForIdle waits for session-level settlement", async () => { + let releaseTool = () => {}; + const released = new Promise((resolve) => { + releaseTool = resolve; + }); + let markCommandStarted = () => {}; + const commandStarted = new Promise((resolve) => { + markCommandStarted = resolve; + }); + const commandResults: boolean[] = []; + const harness = await createHarness({ + tools: [createWaitTool(released)], + extensionFactories: [ + (pi) => { + pi.registerCommand("after-idle", { + description: "Wait for idle", + handler: async (_args, ctx) => { + markCommandStarted(); + await ctx.waitForIdle(); + commandResults.push(ctx.isIdle()); + }, + }); + }, + ], + }); + harnesses.push(harness); + await harness.session.bindExtensions({ + commandContextActions: { + waitForIdle: () => harness.session.waitForIdle(), + newSession: async () => ({ cancelled: false }), + fork: async () => ({ cancelled: false }), + navigateTree: async () => ({ cancelled: false }), + switchSession: async () => ({ cancelled: false }), + reload: async () => {}, + }, + }); + const toolStarted = new Promise((resolve) => { + const unsubscribe = harness.session.subscribe((event) => { + if (event.type === "tool_execution_start" && event.toolName === "wait") { + unsubscribe(); + resolve(); + } + }); + }); + harness.setResponses([ + fauxAssistantMessage(fauxToolCall("wait", {}), { stopReason: "toolUse" }), + fauxAssistantMessage("done"), + ]); + + const promptPromise = harness.session.prompt("start"); + await toolStarted; + const commandPromise = harness.session.prompt("/after-idle"); + await commandStarted; + let commandFinished = false; + void commandPromise.then(() => { + commandFinished = true; + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(commandFinished).toBe(false); + + releaseTool(); + await Promise.all([promptPromise, commandPromise]); + + expect(commandResults).toEqual([true]); + expect(harness.eventsOfType("agent_settled")).toHaveLength(1); + }); +});