Add RPC subagent subscriptions

This commit is contained in:
Kormákur
2026-06-09 23:42:26 +00:00
committed by can1357
parent 96defff9a5
commit 23cce69f4b
11 changed files with 730 additions and 16 deletions
+1
View File
@@ -145,6 +145,7 @@
### Fixed
- Fixed a turn-ending provider error (e.g. a 502 whose body is the proxy's full HTML page) flooding the transcript: `AnthropicApiError` folds the entire response body into `errorMessage`, and the inline transcript render reprinted it verbatim — every embedded blank line included — leaving a tall mostly-empty block ending in `</html>`. The inline error now drops blank lines, clamps to 8 lines, and width-truncates each line via `getPreviewLines`, matching the pinned error banner.
- Added RPC subagent subscription frames, snapshots, and transcript catch-up APIs for desktop clients embedding `omp --mode rpc`.
## [15.10.7] - 2026-06-08
+2 -1
View File
@@ -85,6 +85,7 @@ type RunPrintMode = (session: AgentSession, options: PrintModeOptions) => Promis
type RunRpcMode = (
session: AgentSession,
setToolUIContext?: (uiContext: ExtensionUIContext, hasUI: boolean) => void,
eventBus?: EventBus,
) => Promise<never>;
function maybeShowStartupSplash(options: {
@@ -1298,7 +1299,7 @@ export async function runRootCommand(
// Branch-only protocol runner: keep RPC host code out of normal interactive startup.
const runRpcMode: RunRpcMode = (await import("./modes/rpc/rpc-mode")).runRpcMode;
stopStartupWatchdog();
await runRpcMode(session, mode === "rpc-ui" ? setToolUIContext : undefined);
await runRpcMode(session, mode === "rpc-ui" ? setToolUIContext : undefined, eventBus);
} else if (isInteractive) {
const versionCheckPromise = checkForNewVersion(VERSION).catch(() => undefined);
const changelogMarkdown = await logger.time("main:getChangelogForDisplay", getChangelogForDisplay, parsedArgs);
+12
View File
@@ -18,6 +18,10 @@ export {
type RpcClientToolContext,
type RpcClientToolResult,
type RpcEventListener,
type RpcSessionEventListener,
type RpcSubagentEventListener,
type RpcSubagentLifecycleListener,
type RpcSubagentProgressListener,
} from "./rpc/rpc-client";
export type {
RpcCommand,
@@ -27,7 +31,15 @@ export type {
RpcHostToolResult,
RpcHostToolUpdate,
RpcResponse,
RpcSessionEventFrame,
RpcSessionState,
RpcSubagentEventFrame,
RpcSubagentFrame,
RpcSubagentLifecycleFrame,
RpcSubagentMessagesResult,
RpcSubagentProgressFrame,
RpcSubagentSnapshot,
RpcSubagentSubscriptionLevel,
} from "./rpc/rpc-types";
postmortem.register("terminal-restore", () => {
@@ -10,7 +10,7 @@ import type { CompactionResult } from "@oh-my-pi/pi-agent-core/compaction";
import type { ImageContent, Model } from "@oh-my-pi/pi-ai";
import { isRecord, ptree, readJsonl } from "@oh-my-pi/pi-utils";
import type { BashResult } from "../../exec/bash-executor";
import type { SessionStats } from "../../session/agent-session";
import type { AgentSessionEvent, SessionStats } from "../../session/agent-session";
import type {
RpcCommand,
RpcExtensionUIRequest,
@@ -22,6 +22,12 @@ import type {
RpcHostToolUpdate,
RpcResponse,
RpcSessionState,
RpcSubagentEventFrame,
RpcSubagentLifecycleFrame,
RpcSubagentMessagesResult,
RpcSubagentProgressFrame,
RpcSubagentSnapshot,
RpcSubagentSubscriptionLevel,
} from "./rpc-types";
/** Distributive Omit that works with union types */
@@ -52,6 +58,10 @@ export interface RpcClientOptions {
export type ModelInfo = Pick<Model, "provider" | "id" | "contextWindow" | "reasoning" | "thinking">;
export type RpcEventListener = (event: AgentEvent) => void;
export type RpcSessionEventListener = (event: AgentSessionEvent) => void;
export type RpcSubagentLifecycleListener = (payload: RpcSubagentLifecycleFrame["payload"]) => void;
export type RpcSubagentProgressListener = (payload: RpcSubagentProgressFrame["payload"]) => void;
export type RpcSubagentEventListener = (payload: RpcSubagentEventFrame["payload"]) => void;
export interface RpcClientToolContext<TDetails = unknown> {
toolCallId: string;
@@ -92,6 +102,23 @@ const agentEventTypes = new Set<AgentEvent["type"]>([
"tool_execution_end",
]);
const sessionEventTypes = new Set<AgentSessionEvent["type"]>([
...agentEventTypes,
"auto_compaction_start",
"auto_compaction_end",
"auto_retry_start",
"auto_retry_end",
"retry_fallback_applied",
"retry_fallback_succeeded",
"ttsr_triggered",
"todo_reminder",
"todo_auto_clear",
"irc_message",
"notice",
"thinking_level_changed",
"goal_updated",
]);
function isRpcResponse(value: unknown): value is RpcResponse {
if (!isRecord(value)) return false;
if (value.type !== "response") return false;
@@ -111,6 +138,28 @@ function isAgentEvent(value: unknown): value is AgentEvent {
return agentEventTypes.has(type as AgentEvent["type"]);
}
function isAgentSessionEvent(value: unknown): value is AgentSessionEvent {
if (!isRecord(value)) return false;
const type = value.type;
if (typeof type !== "string") return false;
return sessionEventTypes.has(type as AgentSessionEvent["type"]);
}
function isRpcSubagentLifecycleFrame(value: unknown): value is RpcSubagentLifecycleFrame {
if (!isRecord(value)) return false;
return value.type === "subagent_lifecycle" && isRecord(value.payload);
}
function isRpcSubagentProgressFrame(value: unknown): value is RpcSubagentProgressFrame {
if (!isRecord(value)) return false;
return value.type === "subagent_progress" && isRecord(value.payload);
}
function isRpcSubagentEventFrame(value: unknown): value is RpcSubagentEventFrame {
if (!isRecord(value)) return false;
return value.type === "subagent_event" && isRecord(value.payload);
}
function isRpcHostToolCallRequest(value: unknown): value is RpcHostToolCallRequest {
if (!isRecord(value)) return false;
return (
@@ -148,6 +197,10 @@ function normalizeToolResult<TDetails>(result: RpcClientToolResult<TDetails>): A
export class RpcClient {
#process: ptree.ChildProcess | null = null;
#eventListeners: RpcEventListener[] = [];
#sessionEventListeners: RpcSessionEventListener[] = [];
#subagentLifecycleListeners = new Set<RpcSubagentLifecycleListener>();
#subagentProgressListeners = new Set<RpcSubagentProgressListener>();
#subagentEventListeners = new Set<RpcSubagentEventListener>();
#pendingRequests: Map<string, { resolve: (response: RpcResponse) => void; reject: (error: Error) => void }> =
new Map();
#customTools: RpcClientCustomTool[] = [];
@@ -286,6 +339,43 @@ export class RpcClient {
};
}
/**
* Subscribe to all top-level session events, including non-core session state events.
*/
onSessionEvent(listener: RpcSessionEventListener): () => void {
this.#sessionEventListeners.push(listener);
return () => {
const index = this.#sessionEventListeners.indexOf(listener);
if (index !== -1) {
this.#sessionEventListeners.splice(index, 1);
}
};
}
/**
* Subscribe to subagent lifecycle frames emitted by the task tool.
*/
onSubagentLifecycle(listener: RpcSubagentLifecycleListener): () => void {
this.#subagentLifecycleListeners.add(listener);
return () => this.#subagentLifecycleListeners.delete(listener);
}
/**
* Subscribe to aggregated subagent progress frames emitted by the task tool.
*/
onSubagentProgress(listener: RpcSubagentProgressListener): () => void {
this.#subagentProgressListeners.add(listener);
return () => this.#subagentProgressListeners.delete(listener);
}
/**
* Subscribe to raw subagent session events. Call setSubagentSubscription(\"events\") to enable them server-side.
*/
onSubagentEvent(listener: RpcSubagentEventListener): () => void {
this.#subagentEventListeners.add(listener);
return () => this.#subagentEventListeners.delete(listener);
}
/**
* Get collected stderr output (useful for debugging).
*/
@@ -358,6 +448,40 @@ export class RpcClient {
return this.#getData(response);
}
/**
* Configure subagent frames emitted by the RPC server.
* Progress emits lifecycle/progress frames; events additionally emits raw subagent session events.
*/
async setSubagentSubscription(level: RpcSubagentSubscriptionLevel): Promise<RpcSubagentSubscriptionLevel> {
const response = await this.#send({ type: "set_subagent_subscription", level });
return this.#getData<{ level: RpcSubagentSubscriptionLevel }>(response).level;
}
/**
* Return the RPC server's current subagent snapshot.
*/
async getSubagents(): Promise<RpcSubagentSnapshot[]> {
const response = await this.#send({ type: "get_subagents" });
return this.#getData<{ subagents: RpcSubagentSnapshot[] }>(response).subagents;
}
/**
* Read persisted transcript entries for a tracked subagent session.
*/
async getSubagentMessages(selector: {
subagentId?: string;
sessionFile?: string;
fromByte?: number;
}): Promise<RpcSubagentMessagesResult> {
const response = await this.#send({
type: "get_subagent_messages",
subagentId: selector.subagentId,
sessionFile: selector.sessionFile,
fromByte: selector.fromByte,
});
return this.#getData<RpcSubagentMessagesResult>(response);
}
/**
* Set model by provider and ID.
*/
@@ -679,9 +803,35 @@ export class RpcClient {
return;
}
if (isRpcSubagentLifecycleFrame(data)) {
for (const listener of this.#subagentLifecycleListeners) {
listener(data.payload);
}
return;
}
if (isRpcSubagentProgressFrame(data)) {
for (const listener of this.#subagentProgressListeners) {
listener(data.payload);
}
return;
}
if (isRpcSubagentEventFrame(data)) {
for (const listener of this.#subagentEventListeners) {
listener(data.payload);
}
return;
}
if (!isAgentSessionEvent(data)) return;
for (const listener of this.#sessionEventListeners) {
listener(data);
}
if (!isAgentEvent(data)) return;
// Otherwise it's an event
for (const listener of this.#eventListeners) {
listener(data);
}
@@ -21,9 +21,11 @@ import {
} from "../../extensibility/extensions";
import { type Theme, theme } from "../../modes/theme/theme";
import type { AgentSession } from "../../session/agent-session";
import type { EventBus } from "../../utils/event-bus";
import { initializeExtensions } from "../runtime-init";
import { isRpcHostToolResult, isRpcHostToolUpdate, RpcHostToolBridge } from "./host-tools";
import { isRpcHostUriResult, RpcHostUriBridge } from "./host-uris";
import { RpcSubagentRegistry, readRpcSubagentTranscript } from "./rpc-subagents";
import type {
RpcCommand,
RpcExtensionUIRequest,
@@ -35,6 +37,7 @@ import type {
RpcHostUriRequest,
RpcResponse,
RpcSessionState,
RpcSubagentSubscriptionLevel,
} from "./rpc-types";
// Re-export types for consumers
@@ -99,6 +102,10 @@ function shouldEmitRpcTitles(): boolean {
return normalized === "1" || normalized === "true" || normalized === "yes" || normalized === "on";
}
function isSubagentSubscriptionLevel(value: unknown): value is RpcSubagentSubscriptionLevel {
return value === "off" || value === "progress" || value === "events";
}
export function requestRpcEditor(
pendingRequests: Map<string, PendingExtensionRequest>,
output: RpcOutput,
@@ -169,6 +176,7 @@ export function requestRpcEditor(
export async function runRpcMode(
session: AgentSession,
setToolUIContext?: (uiContext: ExtensionUIContext, hasUI: boolean) => void,
eventBus?: EventBus,
): Promise<never> {
// Signal to RPC clients that the server is ready to accept commands
// Suppress terminal notifications: they write \x07 (BEL) or OSC sequences directly to
@@ -201,6 +209,7 @@ export async function runRpcMode(
const pendingExtensionRequests = new Map<string, PendingExtensionRequest>();
const hostToolBridge = new RpcHostToolBridge(output);
const hostUriBridge = new RpcHostUriBridge(output);
const subagentRegistry = eventBus ? new RpcSubagentRegistry(eventBus, output) : undefined;
// Shutdown request flag (wrapped in object to allow mutation with const)
const shutdownState = { requested: false };
@@ -564,6 +573,41 @@ export async function runRpcMode(
}
}
case "set_subagent_subscription": {
if (!subagentRegistry) {
return error(id, "set_subagent_subscription", "Subagent event bus is unavailable");
}
if (!isSubagentSubscriptionLevel(command.level)) {
return error(
id,
"set_subagent_subscription",
`Invalid subagent subscription level: ${String(command.level)}`,
);
}
subagentRegistry.setSubscriptionLevel(command.level);
return success(id, "set_subagent_subscription", { level: subagentRegistry.getSubscriptionLevel() });
}
case "get_subagents": {
return success(id, "get_subagents", { subagents: subagentRegistry?.getSubagents() ?? [] });
}
case "get_subagent_messages": {
if (!subagentRegistry) {
return error(id, "get_subagent_messages", "Subagent event bus is unavailable");
}
try {
if (command.fromByte !== undefined && !Number.isFinite(command.fromByte)) {
return error(id, "get_subagent_messages", "fromByte must be a finite number");
}
const sessionFile = subagentRegistry.resolveSessionFile(command);
const transcript = await readRpcSubagentTranscript(sessionFile, command.fromByte);
return success(id, "get_subagent_messages", transcript);
} catch (err) {
return error(id, "get_subagent_messages", err instanceof Error ? err.message : String(err));
}
}
// =================================================================
// Model
// =================================================================
@@ -858,5 +902,6 @@ export async function runRpcMode(
// stdin closed — RPC client is gone, exit cleanly
hostToolBridge.rejectAllPending("RPC client disconnected before host tool execution completed");
hostUriBridge.clear("RPC client disconnected before host URI request completed");
subagentRegistry?.dispose();
process.exit(0);
}
@@ -0,0 +1,170 @@
import * as fs from "node:fs/promises";
import type { FileEntry, SessionMessageEntry } from "../../session/session-manager";
import { parseSessionEntries } from "../../session/session-manager";
import {
type AgentProgress,
type SubagentEventPayload,
type SubagentLifecyclePayload,
type SubagentProgressPayload,
TASK_SUBAGENT_EVENT_CHANNEL,
TASK_SUBAGENT_LIFECYCLE_CHANNEL,
TASK_SUBAGENT_PROGRESS_CHANNEL,
} from "../../task";
import type { EventBus } from "../../utils/event-bus";
import type {
RpcSubagentEventFrame,
RpcSubagentFrame,
RpcSubagentMessagesResult,
RpcSubagentSnapshot,
RpcSubagentSubscriptionLevel,
} from "./rpc-types";
export interface RpcSubagentTranscriptSelector {
subagentId?: string;
sessionFile?: string;
fromByte?: number;
}
type RpcSubagentOutput = (frame: RpcSubagentFrame) => void;
function isSessionMessageEntry(entry: FileEntry): entry is SessionMessageEntry {
return entry.type === "message";
}
function statusFromLifecycle(status: SubagentLifecyclePayload["status"]): AgentProgress["status"] {
return status === "started" ? "running" : status;
}
export async function readRpcSubagentTranscript(sessionFile: string, fromByte = 0): Promise<RpcSubagentMessagesResult> {
let startByte = Number.isFinite(fromByte) ? Math.max(0, Math.trunc(fromByte)) : 0;
const file = Bun.file(sessionFile);
const { size } = await fs.stat(sessionFile);
let reset = false;
if (startByte > size) {
startByte = 0;
reset = true;
}
const text = startByte >= size ? "" : await file.slice(startByte).text();
const lastNewline = text.lastIndexOf("\n");
const completeText = lastNewline >= 0 ? text.slice(0, lastNewline + 1) : "";
const entries = completeText.length > 0 ? parseSessionEntries(completeText) : [];
const nextByte = startByte + Buffer.byteLength(completeText, "utf8");
return {
sessionFile,
fromByte: startByte,
nextByte,
reset,
entries,
messages: entries.filter(isSessionMessageEntry).map(entry => entry.message),
};
}
export class RpcSubagentRegistry {
#subagents = new Map<string, RpcSubagentSnapshot>();
#unsubscribers: Array<() => void> = [];
#output: RpcSubagentOutput;
#subscriptionLevel: RpcSubagentSubscriptionLevel = "progress";
constructor(eventBus: EventBus, output: RpcSubagentOutput) {
this.#output = output;
this.#unsubscribers.push(
eventBus.on(TASK_SUBAGENT_LIFECYCLE_CHANNEL, data => {
this.handleLifecycle(data as SubagentLifecyclePayload);
}),
eventBus.on(TASK_SUBAGENT_PROGRESS_CHANNEL, data => {
this.handleProgress(data as SubagentProgressPayload);
}),
eventBus.on(TASK_SUBAGENT_EVENT_CHANNEL, data => {
this.handleEvent(data as SubagentEventPayload);
}),
);
}
dispose(): void {
for (const unsubscribe of this.#unsubscribers) unsubscribe();
this.#unsubscribers = [];
this.#subagents.clear();
}
setSubscriptionLevel(level: RpcSubagentSubscriptionLevel): void {
this.#subscriptionLevel = level;
}
getSubscriptionLevel(): RpcSubagentSubscriptionLevel {
return this.#subscriptionLevel;
}
getSubagents(): RpcSubagentSnapshot[] {
return [...this.#subagents.values()].sort((a, b) => a.index - b.index || a.id.localeCompare(b.id));
}
handleLifecycle(payload: SubagentLifecyclePayload): void {
const existing = this.#subagents.get(payload.id);
const snapshot: RpcSubagentSnapshot = {
id: payload.id,
index: payload.index,
agent: payload.agent,
agentSource: payload.agentSource,
description: payload.description ?? existing?.description,
status: statusFromLifecycle(payload.status),
task: existing?.task,
assignment: existing?.assignment,
sessionFile: payload.sessionFile ?? existing?.sessionFile,
parentToolCallId: payload.parentToolCallId ?? existing?.parentToolCallId,
lastUpdate: Date.now(),
progress: existing?.progress,
};
this.#subagents.set(payload.id, snapshot);
if (this.#subscriptionLevel !== "off") {
this.#output({ type: "subagent_lifecycle", payload });
}
}
handleProgress(payload: SubagentProgressPayload): void {
const progress = payload.progress;
const existing = this.#subagents.get(progress.id);
this.#subagents.set(progress.id, {
id: progress.id,
index: payload.index,
agent: payload.agent,
agentSource: payload.agentSource,
description: progress.description ?? existing?.description,
status: progress.status,
task: payload.task,
assignment: payload.assignment,
sessionFile: payload.sessionFile ?? existing?.sessionFile,
lastUpdate: Date.now(),
parentToolCallId: payload.parentToolCallId ?? existing?.parentToolCallId,
progress,
});
if (this.#subscriptionLevel !== "off") {
this.#output({ type: "subagent_progress", payload });
}
}
handleEvent(payload: SubagentEventPayload): void {
if (this.#subscriptionLevel !== "events") return;
this.#output({ type: "subagent_event", payload } satisfies RpcSubagentEventFrame);
}
resolveSessionFile(selector: RpcSubagentTranscriptSelector): string {
if (selector.subagentId) {
const snapshot = this.#subagents.get(selector.subagentId);
if (!snapshot?.sessionFile) {
throw new Error(`Unknown subagent or session file unavailable: ${selector.subagentId}`);
}
return snapshot.sessionFile;
}
if (selector.sessionFile) {
for (const snapshot of this.#subagents.values()) {
if (snapshot.sessionFile === selector.sessionFile) return selector.sessionFile;
}
throw new Error("Unknown subagent session file");
}
throw new Error("get_subagent_messages requires subagentId or sessionFile");
}
}
@@ -9,7 +9,14 @@ import type { CompactionResult } from "@oh-my-pi/pi-agent-core/compaction";
import type { Effort, ImageContent, Model } from "@oh-my-pi/pi-ai";
import type { BashResult } from "../../exec/bash-executor";
import type { ContextUsage } from "../../extensibility/extensions/types";
import type { SessionStats } from "../../session/agent-session";
import type { AgentSessionEvent, SessionStats } from "../../session/agent-session";
import type { FileEntry } from "../../session/session-manager";
import type {
AgentProgress,
SubagentEventPayload,
SubagentLifecyclePayload,
SubagentProgressPayload,
} from "../../task";
import type { TodoPhase } from "../../tools/todo";
// ============================================================================
@@ -30,6 +37,9 @@ export type RpcCommand =
| { id?: string; type: "set_todos"; phases: TodoPhase[] }
| { id?: string; type: "set_host_tools"; tools: RpcHostToolDefinition[] }
| { id?: string; type: "set_host_uri_schemes"; schemes: RpcHostUriSchemeDefinition[] }
| { id?: string; type: "set_subagent_subscription"; level: RpcSubagentSubscriptionLevel }
| { id?: string; type: "get_subagents" }
| { id?: string; type: "get_subagent_messages"; subagentId?: string; sessionFile?: string; fromByte?: number }
// Model
| { id?: string; type: "set_model"; provider: string; modelId: string }
@@ -104,6 +114,32 @@ export interface RpcHandoffResult {
savedPath?: string;
}
export type RpcSubagentSubscriptionLevel = "off" | "progress" | "events";
export interface RpcSubagentSnapshot {
id: string;
index: number;
agent: string;
agentSource: AgentProgress["agentSource"];
description?: string;
status: AgentProgress["status"];
task?: string;
assignment?: string;
sessionFile?: string;
lastUpdate: number;
progress?: AgentProgress;
parentToolCallId?: string;
}
export interface RpcSubagentMessagesResult {
sessionFile: string;
fromByte: number;
nextByte: number;
reset: boolean;
entries: FileEntry[];
messages: AgentMessage[];
}
// ============================================================================
// RPC Responses (stdout)
// ============================================================================
@@ -123,6 +159,27 @@ export type RpcResponse =
| { id?: string; type: "response"; command: "set_todos"; success: true; data: { todoPhases: TodoPhase[] } }
| { id?: string; type: "response"; command: "set_host_tools"; success: true; data: { toolNames: string[] } }
| { id?: string; type: "response"; command: "set_host_uri_schemes"; success: true; data: { schemes: string[] } }
| {
id?: string;
type: "response";
command: "set_subagent_subscription";
success: true;
data: { level: RpcSubagentSubscriptionLevel };
}
| {
id?: string;
type: "response";
command: "get_subagents";
success: true;
data: { subagents: RpcSubagentSnapshot[] };
}
| {
id?: string;
type: "response";
command: "get_subagent_messages";
success: true;
data: RpcSubagentMessagesResult;
}
// Model
| {
@@ -212,6 +269,29 @@ export type RpcResponse =
// Error response (any command can fail)
| { id?: string; type: "response"; command: string; success: false; error: string };
// ============================================================================
// Subagent Events (stdout)
// ============================================================================
export interface RpcSubagentLifecycleFrame {
type: "subagent_lifecycle";
payload: SubagentLifecyclePayload;
}
export interface RpcSubagentProgressFrame {
type: "subagent_progress";
payload: SubagentProgressPayload;
}
export interface RpcSubagentEventFrame {
type: "subagent_event";
payload: SubagentEventPayload;
}
export type RpcSubagentFrame = RpcSubagentLifecycleFrame | RpcSubagentProgressFrame | RpcSubagentEventFrame;
export type RpcSessionEventFrame = AgentSessionEvent | RpcSubagentFrame;
// ============================================================================
// Extension UI Events (stdout)
// ============================================================================
+20 -12
View File
@@ -162,6 +162,7 @@ export interface ExecutorOptions {
description?: string;
index: number;
id: string;
parentToolCallId?: string;
modelOverride?: string | string[];
/**
* Active model selector of the parent session, used as an auth-aware fallback
@@ -840,6 +841,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
agent: agent.name,
agentSource: agent.source,
task,
parentToolCallId: options.parentToolCallId,
assignment,
progress: { ...progress },
sessionFile: subtaskSessionFile,
@@ -922,20 +924,23 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
progress.recentOutput = [];
};
const emitSubagentEvent = (event: AgentSessionEvent) => {
if (!options.eventBus) return;
options.eventBus.emit(TASK_SUBAGENT_EVENT_CHANNEL, {
id,
index,
agent: agent.name,
agentSource: agent.source,
task,
assignment,
sessionFile: subtaskSessionFile,
event,
parentToolCallId: options.parentToolCallId,
});
};
const processEvent = (event: AgentEvent) => {
if (resolved) return;
if (options.eventBus) {
options.eventBus.emit(TASK_SUBAGENT_EVENT_CHANNEL, {
index,
agent: agent.name,
agentSource: agent.source,
task,
assignment,
event,
});
}
const now = Date.now();
let flushProgress = false;
@@ -1354,6 +1359,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
options.eventBus.emit(TASK_SUBAGENT_LIFECYCLE_CHANNEL, {
id,
agent: agent.name,
parentToolCallId: options.parentToolCallId,
agentSource: agent.source,
description: options.description,
status: "started",
@@ -1452,6 +1458,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
const MAX_YIELD_RETRIES = 3;
unsubscribe = session.subscribe(event => {
emitSubagentEvent(event);
if (event.type === "auto_retry_start") {
progress.retryState = {
attempt: event.attempt,
@@ -1704,6 +1711,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
options.eventBus.emit(TASK_SUBAGENT_LIFECYCLE_CHANNEL, {
id,
agent: agent.name,
parentToolCallId: options.parentToolCallId,
agentSource: agent.source,
description: options.description,
status: progress.status as "completed" | "failed" | "aborted",
+3
View File
@@ -119,6 +119,7 @@ export type {
AgentDefinition,
AgentProgress,
SingleResult,
SubagentEventPayload,
SubagentLifecyclePayload,
SubagentProgressPayload,
TaskParams,
@@ -1035,6 +1036,7 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
planReference,
description: task.description,
index,
parentToolCallId: _toolCallId,
id: task.id,
taskDepth,
modelOverride,
@@ -1097,6 +1099,7 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
planReference,
description: task.description,
index,
parentToolCallId: _toolCallId,
id: task.id,
taskDepth,
modelOverride,
+16
View File
@@ -2,6 +2,7 @@ import type { ThinkingLevel } from "@oh-my-pi/pi-agent-core";
import type { Usage } from "@oh-my-pi/pi-ai";
import { $env } from "@oh-my-pi/pi-utils";
import * as z from "zod/v4";
import type { AgentSessionEvent } from "../session/agent-session";
import { getTaskSimpleModeCapabilities, type TaskSimpleMode } from "./simple-mode";
import type { NestedRepoPatch } from "./worktree";
@@ -41,11 +42,25 @@ export interface SubagentProgressPayload {
agent: string;
agentSource: AgentSource;
task: string;
parentToolCallId?: string;
assignment?: string;
progress: AgentProgress;
sessionFile?: string;
}
/** Payload emitted on TASK_SUBAGENT_EVENT_CHANNEL */
export interface SubagentEventPayload {
id: string;
index: number;
agent: string;
agentSource: AgentSource;
task: string;
parentToolCallId?: string;
assignment?: string;
sessionFile?: string;
event: AgentSessionEvent;
}
/** Payload emitted on TASK_SUBAGENT_LIFECYCLE_CHANNEL */
export interface SubagentLifecyclePayload {
id: string;
@@ -54,6 +69,7 @@ export interface SubagentLifecyclePayload {
description?: string;
status: "started" | "completed" | "failed" | "aborted";
sessionFile?: string;
parentToolCallId?: string;
index: number;
}
@@ -0,0 +1,228 @@
import { afterEach, describe, expect, test } from "bun:test";
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import { RpcClient } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-client";
import { RpcSubagentRegistry, readRpcSubagentTranscript } 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 {
type AgentProgress,
type SubagentEventPayload,
type SubagentLifecyclePayload,
type SubagentProgressPayload,
TASK_SUBAGENT_EVENT_CHANNEL,
TASK_SUBAGENT_LIFECYCLE_CHANNEL,
TASK_SUBAGENT_PROGRESS_CHANNEL,
} from "@oh-my-pi/pi-coding-agent/task";
import { EventBus } from "@oh-my-pi/pi-coding-agent/utils/event-bus";
const tempPaths: string[] = [];
afterEach(() => {
for (const tempPath of tempPaths.splice(0)) {
fs.rmSync(tempPath, { recursive: true, force: true });
}
});
function createProgress(overrides: Partial<AgentProgress> = {}): AgentProgress {
return {
index: 0,
id: "SubagentA",
agent: "task",
agentSource: "bundled",
status: "running",
task: "Do work",
assignment: "Implement work",
description: "Worker",
recentTools: [],
recentOutput: [],
toolCount: 0,
tokens: 0,
cost: 0,
durationMs: 0,
...overrides,
};
}
describe("RPC subagent registry", () => {
test("emits progress frames and snapshots tracked subagents", () => {
const eventBus = new EventBus();
const frames: RpcSubagentFrame[] = [];
const registry = new RpcSubagentRegistry(eventBus, frame => frames.push(frame));
const lifecycle: SubagentLifecyclePayload = {
id: "SubagentA",
index: 0,
agent: "task",
agentSource: "bundled",
description: "Worker",
status: "started",
sessionFile: "/tmp/subagent.jsonl",
parentToolCallId: "toolu_parent",
};
const progressPayload: SubagentProgressPayload = {
index: 0,
agent: "task",
agentSource: "bundled",
task: "Do work",
assignment: "Implement work",
parentToolCallId: "toolu_parent",
sessionFile: "/tmp/subagent.jsonl",
progress: createProgress(),
};
eventBus.emit(TASK_SUBAGENT_LIFECYCLE_CHANNEL, lifecycle);
eventBus.emit(TASK_SUBAGENT_PROGRESS_CHANNEL, progressPayload);
expect(frames.map(frame => frame.type)).toEqual(["subagent_lifecycle", "subagent_progress"]);
expect(registry.getSubagents()).toMatchObject([
{
id: "SubagentA",
status: "running",
task: "Do work",
assignment: "Implement work",
sessionFile: "/tmp/subagent.jsonl",
parentToolCallId: "toolu_parent",
},
]);
registry.dispose();
});
test("gates raw subagent events behind the events subscription level", () => {
const eventBus = new EventBus();
const frames: RpcSubagentFrame[] = [];
const registry = new RpcSubagentRegistry(eventBus, frame => frames.push(frame));
const eventPayload: SubagentEventPayload = {
id: "SubagentA",
index: 0,
agent: "task",
agentSource: "bundled",
task: "Do work",
assignment: "Implement work",
parentToolCallId: "toolu_parent",
sessionFile: "/tmp/subagent.jsonl",
event: { type: "agent_start" },
};
eventBus.emit(TASK_SUBAGENT_EVENT_CHANNEL, eventPayload);
expect(frames).toHaveLength(0);
registry.setSubscriptionLevel("events");
eventBus.emit(TASK_SUBAGENT_EVENT_CHANNEL, eventPayload);
expect(frames).toHaveLength(1);
expect(frames[0]).toMatchObject({ type: "subagent_event", payload: { id: "SubagentA" } });
registry.dispose();
});
});
describe("readRpcSubagentTranscript", () => {
test("returns complete JSONL entries and byte cursor", async () => {
const dir = fs.mkdtempSync(path.join(os.tmpdir(), "omp-rpc-subagent-transcript-"));
tempPaths.push(dir);
const sessionFile = path.join(dir, "session.jsonl");
const headerLine = `${JSON.stringify({ type: "session", id: "s1", timestamp: "2026-06-09T00:00:00.000Z", cwd: dir })}\n`;
const messageLine = `${JSON.stringify({
type: "message",
id: "m1",
parentId: null,
timestamp: "2026-06-09T00:00:00.000Z",
message: { role: "user", content: [{ type: "text", text: "hello" }] },
})}\n`;
await Bun.write(sessionFile, `${headerLine}${messageLine}{"type":"message"`);
const result = await readRpcSubagentTranscript(sessionFile);
expect(result.entries).toHaveLength(2);
expect(result.messages).toHaveLength(1);
expect(result.nextByte).toBe(Buffer.byteLength(`${headerLine}${messageLine}`, "utf8"));
expect(result.reset).toBe(false);
});
});
describe("RpcClient subagent frames", () => {
test("dispatches subagent frames and session-specific events", async () => {
const scriptPath = path.join(os.tmpdir(), `omp-rpc-subagent-client-${Date.now()}.js`);
tempPaths.push(scriptPath);
await Bun.write(
scriptPath,
`
let buffer = "";
function write(frame) {
process.stdout.write(JSON.stringify(frame) + "\\n");
}
const progress = {
index: 0,
id: "SubagentA",
agent: "task",
agentSource: "bundled",
status: "running",
task: "Do work",
assignment: "Implement work",
recentTools: [],
recentOutput: [],
toolCount: 0,
tokens: 0,
cost: 0,
durationMs: 0
};
write({ type: "ready" });
process.stdin.on("data", chunk => {
buffer += chunk.toString("utf8");
let index = buffer.indexOf("\\n");
while (index !== -1) {
const line = buffer.slice(0, index).trim();
buffer = buffer.slice(index + 1);
if (line) handle(JSON.parse(line));
index = buffer.indexOf("\\n");
}
});
function handle(frame) {
if (frame.type === "set_subagent_subscription") {
write({ id: frame.id, type: "response", command: "set_subagent_subscription", success: true, data: { level: frame.level } });
return;
}
if (frame.type === "get_subagents") {
write({ id: frame.id, type: "response", command: "get_subagents", success: true, data: { subagents: [{ id: "SubagentA", index: 0, agent: "task", agentSource: "bundled", status: "running", lastUpdate: 1 }] } });
return;
}
if (frame.type === "get_subagent_messages") {
write({ id: frame.id, type: "response", command: "get_subagent_messages", success: true, data: { sessionFile: frame.sessionFile || "/tmp/subagent.jsonl", fromByte: frame.fromByte || 0, nextByte: 0, reset: false, entries: [], messages: [] } });
return;
}
if (frame.type === "prompt") {
write({ id: frame.id, type: "response", command: "prompt", success: true });
write({ type: "notice", level: "info", message: "subagent test" });
write({ type: "subagent_lifecycle", payload: { id: "SubagentA", index: 0, agent: "task", agentSource: "bundled", status: "started", sessionFile: "/tmp/subagent.jsonl" } });
write({ type: "subagent_progress", payload: { index: 0, agent: "task", agentSource: "bundled", task: "Do work", assignment: "Implement work", sessionFile: "/tmp/subagent.jsonl", progress } });
write({ type: "subagent_event", payload: { id: "SubagentA", index: 0, agent: "task", agentSource: "bundled", task: "Do work", assignment: "Implement work", sessionFile: "/tmp/subagent.jsonl", event: { type: "agent_start" } } });
write({ type: "agent_end", messages: [] });
}
}
`,
);
using client = new RpcClient({ cliPath: scriptPath });
const lifecycleIds: string[] = [];
const progressTasks: string[] = [];
const rawEventTypes: string[] = [];
const sessionEventTypes: string[] = [];
client.onSubagentLifecycle(payload => lifecycleIds.push(payload.id));
client.onSubagentProgress(payload => progressTasks.push(payload.task));
client.onSubagentEvent(payload => rawEventTypes.push(payload.event.type));
client.onSessionEvent(event => sessionEventTypes.push(event.type));
await client.start();
await expect(client.setSubagentSubscription("events")).resolves.toBe("events");
await client.promptAndWait("Trigger subagent frames");
expect(await client.getSubagents()).toHaveLength(1);
expect(await client.getSubagentMessages({ sessionFile: "/tmp/subagent.jsonl" })).toMatchObject({
sessionFile: "/tmp/subagent.jsonl",
});
expect(lifecycleIds).toEqual(["SubagentA"]);
expect(progressTasks).toEqual(["Do work"]);
expect(rawEventTypes).toEqual(["agent_start"]);
expect(sessionEventTypes).toContain("notice");
});
});