@@ -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<void> | 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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
this.abortRetry();
|
||||
this.agent.abort();
|
||||
await this.agent.waitForIdle();
|
||||
await this.waitForIdle();
|
||||
}
|
||||
|
||||
async waitForIdle(): Promise<void> {
|
||||
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: () => {
|
||||
|
||||
@@ -22,6 +22,7 @@ export { ExtensionRunner } from "./runner.ts";
|
||||
export type {
|
||||
AfterProviderResponseEvent,
|
||||
AgentEndEvent,
|
||||
AgentSettledEvent,
|
||||
AgentStartEvent,
|
||||
// Re-exports
|
||||
AgentToolResult,
|
||||
|
||||
@@ -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<BeforeAgentStartEvent, BeforeAgentStartEventResult>): void;
|
||||
on(event: "agent_start", handler: ExtensionHandler<AgentStartEvent>): void;
|
||||
on(event: "agent_end", handler: ExtensionHandler<AgentEndEvent>): void;
|
||||
on(event: "agent_settled", handler: ExtensionHandler<AgentSettledEvent>): void;
|
||||
on(event: "turn_start", handler: ExtensionHandler<TurnStartEvent>): void;
|
||||
on(event: "turn_end", handler: ExtensionHandler<TurnEndEvent>): void;
|
||||
on(event: "message_start", handler: ExtensionHandler<MessageStartEvent>): void;
|
||||
|
||||
@@ -32,6 +32,7 @@ export { areExperimentalFeaturesEnabled } from "./experimental.ts";
|
||||
// Extensions system
|
||||
export {
|
||||
type AgentEndEvent,
|
||||
type AgentSettledEvent,
|
||||
type AgentStartEvent,
|
||||
type AgentToolResult,
|
||||
type AgentToolUpdateCallback,
|
||||
|
||||
@@ -61,6 +61,7 @@ export { createEventBus, type EventBus, type EventBusController } from "./core/e
|
||||
// Extension system
|
||||
export type {
|
||||
AgentEndEvent,
|
||||
AgentSettledEvent,
|
||||
AgentStartEvent,
|
||||
AgentToolResult,
|
||||
AgentToolUpdateCallback,
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<void> {
|
||||
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<AgentEvent[]> {
|
||||
collectEvents(timeout = 60000): Promise<AgentSessionEvent[]> {
|
||||
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<AgentEvent[]> {
|
||||
async promptAndWait(message: string, images?: ImageContent[], timeout = 60000): Promise<AgentSessionEvent[]> {
|
||||
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
|
||||
|
||||
@@ -319,7 +319,7 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
||||
uiContext: createExtensionUIContext(),
|
||||
mode: "rpc",
|
||||
commandContextActions: {
|
||||
waitForIdle: () => 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<neve
|
||||
unsubscribeBackpressure?.();
|
||||
unsubscribe = session.subscribe((event) => {
|
||||
output(event);
|
||||
if (event.type === "agent_settled") {
|
||||
void checkShutdownRequested();
|
||||
}
|
||||
});
|
||||
unsubscribeBackpressure = session.agent.subscribe(async () => {
|
||||
await waitForRawStdoutBackpressure();
|
||||
|
||||
Reference in New Issue
Block a user