fix(rpc): restored frames for IRC-revived subagents
- Monitored autonomous IRC wake turns with the task executor lifecycle and progress channels. - Preserved monitoring after idle-TTL parking and session revival. - Covered RPC subscriptions for both idle and parked keep-alive agents. Fixes #7105
This commit is contained in:
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed RPC hosts receiving no subagent lifecycle or progress frames when an IRC message revives an idle or parked keep-alive subagent ([#7105](https://github.com/can1357/oh-my-pi/issues/7105)).
|
||||
|
||||
## [17.2.1] - 2026-07-30
|
||||
|
||||
### Added
|
||||
|
||||
@@ -502,6 +502,7 @@ export class AgentSession {
|
||||
#asyncDeliveryEpoch = 0;
|
||||
|
||||
readonly #irc: IrcBridge;
|
||||
#ircWakeTurnObserver: ((records: CustomMessage[]) => ((error?: unknown) => void) | undefined) | undefined;
|
||||
// Agent identity (registry id) used for IRC routing and job ownership.
|
||||
#agentId: string | undefined;
|
||||
#agentKind: "main" | "sub" = "main";
|
||||
@@ -726,14 +727,27 @@ export class AgentSession {
|
||||
if (parkedFollowUps.length > 0) {
|
||||
this.agent.replaceQueues([...this.agent.peekSteeringQueue()], []);
|
||||
}
|
||||
let finishObservation: ((error?: unknown) => void) | undefined;
|
||||
try {
|
||||
finishObservation = this.#ircWakeTurnObserver?.(records);
|
||||
} catch (error) {
|
||||
logger.warn("IRC wake turn observer failed to start", { error: String(error) });
|
||||
}
|
||||
this.#resetPromptMaintenanceState();
|
||||
this.#beginInFlight();
|
||||
let turnError: unknown;
|
||||
void this.agent
|
||||
.prompt(records)
|
||||
.catch(error => {
|
||||
turnError = error;
|
||||
logger.warn("IRC wake turn failed", { error: String(error) });
|
||||
})
|
||||
.finally(() => {
|
||||
try {
|
||||
finishObservation?.(turnError);
|
||||
} catch (error) {
|
||||
logger.warn("IRC wake turn observer failed to finish", { error: String(error) });
|
||||
}
|
||||
if (parkedFollowUps.length > 0) {
|
||||
this.agent.replaceQueues(
|
||||
[...this.agent.peekSteeringQueue()],
|
||||
@@ -6868,6 +6882,13 @@ export class AgentSession {
|
||||
return this.#irc.deliver(msg, opts);
|
||||
}
|
||||
|
||||
/** Installs task-executor monitoring around autonomous IRC wake turns. */
|
||||
setIrcWakeTurnObserver(
|
||||
observer: ((records: CustomMessage[]) => ((error?: unknown) => void) | undefined) | undefined,
|
||||
): void {
|
||||
this.#ircWakeTurnObserver = observer;
|
||||
}
|
||||
|
||||
/** Emits an IRC relay observation for UI rendering without persisting it. */
|
||||
emitIrcRelayObservation(record: CustomMessage): void {
|
||||
this.#irc.emitRelayObservation(record);
|
||||
|
||||
@@ -2544,6 +2544,102 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
}
|
||||
});
|
||||
};
|
||||
const installIrcWakeTurnMonitor = (target: AgentSession): void => {
|
||||
target.setIrcWakeTurnObserver(records => {
|
||||
const ircTask =
|
||||
records
|
||||
.map(record => {
|
||||
const body =
|
||||
record.details && typeof record.details === "object"
|
||||
? Reflect.get(record.details, "message")
|
||||
: undefined;
|
||||
return typeof body === "string" ? body : record.content;
|
||||
})
|
||||
.filter(Boolean)
|
||||
.join("\n\n") || "IRC follow-up";
|
||||
const turnStartTime = Date.now();
|
||||
const sessionFile = AgentRegistry.global().get(id)?.sessionFile ?? subtaskSessionFile ?? undefined;
|
||||
const turnMonitor = createSubagentRunMonitor({
|
||||
index,
|
||||
id,
|
||||
agent,
|
||||
task: ircTask,
|
||||
description: options.description,
|
||||
modelOverride,
|
||||
eventBus: options.eventBus,
|
||||
parentToolCallId: options.parentToolCallId,
|
||||
detached: true,
|
||||
sessionFile,
|
||||
softRequestBudget: 0,
|
||||
softRequestBudgetNotice: false,
|
||||
maxRuntimeMs,
|
||||
});
|
||||
|
||||
if (options.eventBus) {
|
||||
options.eventBus.emit(TASK_SUBAGENT_LIFECYCLE_CHANNEL, {
|
||||
id,
|
||||
agent: agent.name,
|
||||
parentToolCallId: options.parentToolCallId,
|
||||
detached: true,
|
||||
agentSource: agent.source,
|
||||
description: options.description,
|
||||
status: "started",
|
||||
sessionFile,
|
||||
index,
|
||||
});
|
||||
}
|
||||
|
||||
turnMonitor.setActiveSession(target);
|
||||
const unsubscribeTurn = turnMonitor.attach(target);
|
||||
return turnError => {
|
||||
unsubscribeTurn();
|
||||
const activeSession = turnMonitor.takeActiveSession();
|
||||
if (activeSession) turnMonitor.captureSalvage(activeSession);
|
||||
const lastAssistant = target.getLastAssistantMessage();
|
||||
const yielded = turnMonitor.yieldCalled();
|
||||
const runtimeLimitExceeded = turnMonitor.runtimeLimitExceeded();
|
||||
const aborted = runtimeLimitExceeded || (lastAssistant?.stopReason === "aborted" && !yielded);
|
||||
const error =
|
||||
lastAssistant?.stopReason === "error"
|
||||
? lastAssistant.errorMessage || "Subagent failed"
|
||||
: turnError !== undefined && !yielded
|
||||
? turnError instanceof Error
|
||||
? turnError.stack || turnError.message
|
||||
: String(turnError)
|
||||
: undefined;
|
||||
turnMonitor.finish();
|
||||
void finalizeRunResult({
|
||||
monitor: turnMonitor,
|
||||
done: {
|
||||
exitCode: aborted || error ? 1 : 0,
|
||||
error,
|
||||
aborted,
|
||||
abortReason: aborted ? turnMonitor.resolveAbortReasonText() : undefined,
|
||||
durationMs: Date.now() - turnStartTime,
|
||||
},
|
||||
index,
|
||||
id,
|
||||
agent,
|
||||
task: ircTask,
|
||||
modelOverride,
|
||||
outputSchema,
|
||||
outputSchemaMode: options.outputSchemaMode,
|
||||
outputSchemaSource: options.outputSchemaSource,
|
||||
artifactsDir: options.artifactsDir,
|
||||
eventBus: options.eventBus,
|
||||
parentToolCallId: options.parentToolCallId,
|
||||
detached: true,
|
||||
sessionFile,
|
||||
startTime: turnStartTime,
|
||||
}).catch(error => {
|
||||
logger.warn("IRC subagent turn finalization failed", {
|
||||
id,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
});
|
||||
};
|
||||
});
|
||||
};
|
||||
|
||||
const runSubagent = async (): Promise<{
|
||||
exitCode: number;
|
||||
@@ -2882,6 +2978,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
buildSubagentSessionOptions(reopened, expectedAgentRef),
|
||||
);
|
||||
installRegistryStatusSync(revived);
|
||||
installIrcWakeTurnMonitor(revived);
|
||||
return revived;
|
||||
};
|
||||
}
|
||||
@@ -3053,6 +3150,9 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
const session = monitor.takeActiveSession();
|
||||
if (session) {
|
||||
monitor.captureSalvage(session);
|
||||
if (options.keepAlive !== false && worktree === undefined) {
|
||||
installIrcWakeTurnMonitor(session);
|
||||
}
|
||||
await finalizeSubagentLifecycle({
|
||||
id,
|
||||
session,
|
||||
|
||||
@@ -3,11 +3,14 @@ import type { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-regis
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import type { LoadExtensionsResult } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/types";
|
||||
import { IrcBus } from "@oh-my-pi/pi-coding-agent/irc/bus";
|
||||
import { RpcSubagentRegistry } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-subagents";
|
||||
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 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";
|
||||
import type { CustomMessage } from "@oh-my-pi/pi-coding-agent/session/messages";
|
||||
import { resolveSoftRequestBudget, runSubprocess } from "@oh-my-pi/pi-coding-agent/task/executor";
|
||||
import type { AgentDefinition } from "@oh-my-pi/pi-coding-agent/task/types";
|
||||
import { EventBus } from "@oh-my-pi/pi-coding-agent/utils/event-bus";
|
||||
@@ -21,7 +24,8 @@ import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
* forced final `yield`, so the run finishes as a normal completion with a
|
||||
* partial report — not as an abort with no output.
|
||||
* 2. If the agent still refuses to yield (grace exhausted → hard abort), a
|
||||
* kept-alive agent stays adopted (`idle`), so `irc` can message/resume it.
|
||||
* kept-alive agent stays adopted (`idle`), so `irc` can message/resume it
|
||||
* with the same RPC lifecycle/progress frames as the original run.
|
||||
* 3. Caller-signal aborts remain terminal, and the irc bus names the aborted
|
||||
* agent precisely instead of claiming it is unknown.
|
||||
*/
|
||||
@@ -50,6 +54,7 @@ function createMockSession(
|
||||
let abortCount = 0;
|
||||
let disposeCount = 0;
|
||||
let promptIndex = 0;
|
||||
let ircWakeTurnObserver: ((records: CustomMessage[]) => ((error?: unknown) => void) | undefined) | undefined;
|
||||
|
||||
const emit = (event: AgentSessionEvent) => {
|
||||
for (const listener of [...listeners]) listener(event);
|
||||
@@ -80,7 +85,49 @@ function createMockSession(
|
||||
waitForIdle: async () => {},
|
||||
getLastAssistantMessage: () => messages[messages.length - 1] as never,
|
||||
sendUserMessage: async () => {},
|
||||
deliverIrcMessage: async () => "woken",
|
||||
setIrcWakeTurnObserver: observer => {
|
||||
ircWakeTurnObserver = observer;
|
||||
},
|
||||
deliverIrcMessage: async msg => {
|
||||
const record: CustomMessage = {
|
||||
role: "custom",
|
||||
customType: "irc:incoming",
|
||||
content: msg.body,
|
||||
display: true,
|
||||
details: { id: msg.id, from: msg.from, message: msg.body },
|
||||
attribution: "agent",
|
||||
timestamp: msg.ts,
|
||||
};
|
||||
const finishObservation = ircWakeTurnObserver?.([record]);
|
||||
const yieldMessage = {
|
||||
role: "assistant" as const,
|
||||
content: [
|
||||
{
|
||||
type: "toolCall" as const,
|
||||
id: "tool-irc-yield",
|
||||
name: "yield",
|
||||
arguments: { result: { data: { report: "resumed findings" } } },
|
||||
},
|
||||
],
|
||||
stopReason: "toolUse" as const,
|
||||
};
|
||||
messages.push(yieldMessage);
|
||||
emit({ type: "agent_start" } as AgentSessionEvent);
|
||||
emit({ type: "message_end", message: yieldMessage } as unknown as AgentSessionEvent);
|
||||
emit({
|
||||
type: "tool_execution_end",
|
||||
toolCallId: "tool-irc-yield",
|
||||
toolName: "yield",
|
||||
result: {
|
||||
content: [{ type: "text", text: "Result submitted." }],
|
||||
details: { status: "success", data: { report: "resumed findings" } },
|
||||
},
|
||||
isError: false,
|
||||
} as AgentSessionEvent);
|
||||
emit({ type: "agent_end", messages: [yieldMessage] } as unknown as AgentSessionEvent);
|
||||
finishObservation?.();
|
||||
return "woken";
|
||||
},
|
||||
abort: async () => {
|
||||
abortCount += 1;
|
||||
},
|
||||
@@ -129,7 +176,7 @@ describe("runSubprocess soft request budget", () => {
|
||||
tempDir[Symbol.dispose]();
|
||||
});
|
||||
|
||||
function baseOptions(id: string) {
|
||||
function baseOptions(id: string, eventBus?: EventBus) {
|
||||
return {
|
||||
cwd: "/tmp",
|
||||
agent: baseAgent,
|
||||
@@ -140,6 +187,7 @@ describe("runSubprocess soft request budget", () => {
|
||||
modelRegistry: { refresh: async () => {} } as unknown as ModelRegistry,
|
||||
enableLsp: false,
|
||||
artifactsDir: tempDir.path(),
|
||||
eventBus,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -221,6 +269,22 @@ describe("runSubprocess soft request budget", () => {
|
||||
|
||||
it("a budget hard-abort keeps the kept-alive agent adopted and messageable via irc", async () => {
|
||||
const id = "StubbornScout";
|
||||
const eventBus = new EventBus();
|
||||
const frames: RpcSubagentFrame[] = [];
|
||||
let resolveFollowUpTerminal: (() => void) | undefined;
|
||||
const waitForFollowUpTerminal = (): Promise<void> => {
|
||||
const deferred = Promise.withResolvers<void>();
|
||||
resolveFollowUpTerminal = deferred.resolve;
|
||||
return deferred.promise;
|
||||
};
|
||||
const rpcRegistry = new RpcSubagentRegistry(eventBus, frame => {
|
||||
frames.push(frame);
|
||||
if (frame.type !== "subagent_lifecycle" || frame.payload.status === "started") return;
|
||||
const resolve = resolveFollowUpTerminal;
|
||||
resolveFollowUpTerminal = undefined;
|
||||
resolve?.();
|
||||
});
|
||||
rpcRegistry.setSubscriptionLevel("progress");
|
||||
const handle = createMockSession(({ promptIndex, emit, pushMessage }) => {
|
||||
if (promptIndex !== 1) return;
|
||||
// Never yields: budget 2 → stop at 3, grace exhausted at 3 + 5 = 8.
|
||||
@@ -233,7 +297,7 @@ describe("runSubprocess soft request budget", () => {
|
||||
mockCreateAgentSession(handle.session);
|
||||
registerRunning(id, handle.session);
|
||||
|
||||
const result = await runSubprocess(baseOptions(id));
|
||||
const result = await runSubprocess(baseOptions(id, eventBus));
|
||||
|
||||
expect(result.aborted).toBe(true);
|
||||
expect(result.abortReason).toMatch(/Soft request budget exceeded/);
|
||||
@@ -242,9 +306,34 @@ describe("runSubprocess soft request budget", () => {
|
||||
expect(AgentLifecycleManager.global().has(id)).toBe(true);
|
||||
expect(handle.disposeCalls()).toBe(0);
|
||||
|
||||
// The whole point: irc can reach the stopped agent to resume it.
|
||||
const receipt = await new IrcBus().send({ from: "Main", to: id, body: "resume your inventory" });
|
||||
expect(receipt.outcome).toBe("woken");
|
||||
const expectRpcTurn = (): void => {
|
||||
expect(frames[0]).toMatchObject({
|
||||
type: "subagent_lifecycle",
|
||||
payload: { id, status: "started" },
|
||||
});
|
||||
expect(frames.some(frame => frame.type === "subagent_progress")).toBe(true);
|
||||
expect(frames.at(-1)).toMatchObject({
|
||||
type: "subagent_lifecycle",
|
||||
payload: { id, status: "completed" },
|
||||
});
|
||||
};
|
||||
|
||||
frames.length = 0;
|
||||
const idleTerminal = waitForFollowUpTerminal();
|
||||
const idleReceipt = await new IrcBus().send({ from: "Main", to: id, body: "resume your inventory" });
|
||||
expect(idleReceipt.outcome).toBe("woken");
|
||||
await idleTerminal;
|
||||
expectRpcTurn();
|
||||
|
||||
await AgentLifecycleManager.global().park(id);
|
||||
expect(AgentRegistry.global().get(id)?.status).toBe("parked");
|
||||
frames.length = 0;
|
||||
const revivedTerminal = waitForFollowUpTerminal();
|
||||
const revivedReceipt = await new IrcBus().send({ from: "Main", to: id, body: "resume after parking" });
|
||||
expect(revivedReceipt.outcome).toBe("revived");
|
||||
await revivedTerminal;
|
||||
expectRpcTurn();
|
||||
rpcRegistry.dispose();
|
||||
});
|
||||
|
||||
it("a caller-signal abort stays terminal and irc names the aborted agent precisely", async () => {
|
||||
|
||||
Reference in New Issue
Block a user