From 23cce69f4b904bc97b772a9364018dd74fce931d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Korm=C3=A1kur?= Date: Tue, 9 Jun 2026 23:42:26 +0000 Subject: [PATCH] Add RPC subagent subscriptions --- packages/coding-agent/CHANGELOG.md | 1 + packages/coding-agent/src/main.ts | 3 +- packages/coding-agent/src/modes/index.ts | 12 + .../coding-agent/src/modes/rpc/rpc-client.ts | 154 +++++++++++- .../coding-agent/src/modes/rpc/rpc-mode.ts | 45 ++++ .../src/modes/rpc/rpc-subagents.ts | 170 +++++++++++++ .../coding-agent/src/modes/rpc/rpc-types.ts | 82 ++++++- packages/coding-agent/src/task/executor.ts | 32 ++- packages/coding-agent/src/task/index.ts | 3 + packages/coding-agent/src/task/types.ts | 16 ++ .../coding-agent/test/rpc-subagents.test.ts | 228 ++++++++++++++++++ 11 files changed, 730 insertions(+), 16 deletions(-) create mode 100644 packages/coding-agent/src/modes/rpc/rpc-subagents.ts create mode 100644 packages/coding-agent/test/rpc-subagents.test.ts diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index ee70e90fe..ea7fd9783 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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 ``. 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 diff --git a/packages/coding-agent/src/main.ts b/packages/coding-agent/src/main.ts index 3fe16d8b0..b8473ba75 100644 --- a/packages/coding-agent/src/main.ts +++ b/packages/coding-agent/src/main.ts @@ -85,6 +85,7 @@ type RunPrintMode = (session: AgentSession, options: PrintModeOptions) => Promis type RunRpcMode = ( session: AgentSession, setToolUIContext?: (uiContext: ExtensionUIContext, hasUI: boolean) => void, + eventBus?: EventBus, ) => Promise; 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); diff --git a/packages/coding-agent/src/modes/index.ts b/packages/coding-agent/src/modes/index.ts index ac9f12896..8de01c219 100644 --- a/packages/coding-agent/src/modes/index.ts +++ b/packages/coding-agent/src/modes/index.ts @@ -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", () => { diff --git a/packages/coding-agent/src/modes/rpc/rpc-client.ts b/packages/coding-agent/src/modes/rpc/rpc-client.ts index 29df4b683..ff7b88740 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-client.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-client.ts @@ -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; 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 { toolCallId: string; @@ -92,6 +102,23 @@ const agentEventTypes = new Set([ "tool_execution_end", ]); +const sessionEventTypes = new Set([ + ...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(result: RpcClientToolResult): A export class RpcClient { #process: ptree.ChildProcess | null = null; #eventListeners: RpcEventListener[] = []; + #sessionEventListeners: RpcSessionEventListener[] = []; + #subagentLifecycleListeners = new Set(); + #subagentProgressListeners = new Set(); + #subagentEventListeners = new Set(); #pendingRequests: Map 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 { + 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 { + 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 { + const response = await this.#send({ + type: "get_subagent_messages", + subagentId: selector.subagentId, + sessionFile: selector.sessionFile, + fromByte: selector.fromByte, + }); + return this.#getData(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); } diff --git a/packages/coding-agent/src/modes/rpc/rpc-mode.ts b/packages/coding-agent/src/modes/rpc/rpc-mode.ts index 778cfa649..a74880cd4 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-mode.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-mode.ts @@ -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, output: RpcOutput, @@ -169,6 +176,7 @@ export function requestRpcEditor( export async function runRpcMode( session: AgentSession, setToolUIContext?: (uiContext: ExtensionUIContext, hasUI: boolean) => void, + eventBus?: EventBus, ): Promise { // 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(); 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); } diff --git a/packages/coding-agent/src/modes/rpc/rpc-subagents.ts b/packages/coding-agent/src/modes/rpc/rpc-subagents.ts new file mode 100644 index 000000000..614c696bb --- /dev/null +++ b/packages/coding-agent/src/modes/rpc/rpc-subagents.ts @@ -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 { + 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(); + #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"); + } +} diff --git a/packages/coding-agent/src/modes/rpc/rpc-types.ts b/packages/coding-agent/src/modes/rpc/rpc-types.ts index 5ff67a664..8ee5f421f 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-types.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-types.ts @@ -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) // ============================================================================ diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index 5fcc075ce..2f1fae54d 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -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 { + 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 { + emitSubagentEvent(event); if (event.type === "auto_retry_start") { progress.retryState = { attempt: event.attempt, @@ -1704,6 +1711,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise { + for (const tempPath of tempPaths.splice(0)) { + fs.rmSync(tempPath, { recursive: true, force: true }); + } +}); + +function createProgress(overrides: Partial = {}): 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"); + }); +});