{"version":3,"file":"rpc.d.ts","sourceRoot":"","sources":["../../../src/extension/rpc.ts"],"names":[],"mappings":"AACA,OAAO,KAAK,EAAE,eAAe,EAAE,MAAM,yBAAyB,CAAC;AAC/D,OAAO,KAAK,EAAE,gBAAgB,EAAE,MAAM,2BAA2B,CAAC;AAKlE,OAAO,KAAK,EAAE,kBAAkB,EAAE,MAAM,yCAAyC,CAAC;AAElF,OAAO,EAEN,KAAK,OAAO,EAKZ,KAAK,aAAa,EAClB,MAAM,oBAAoB,CAAC;AAM5B,eAAO,MAAM,6BAA6B,IAAI,CAAC;AAC/C,eAAO,MAAM,0BAA0B,6BAA6B,CAAC;AACrE,eAAO,MAAM,wBAAwB,2BAA2B,CAAC;AACjE,eAAO,MAAM,+BAA+B,4BAA4B,CAAC;AAEzE,eAAO,MAAM,oBAAoB,8EAA+E,CAAC;AACjH,MAAM,MAAM,iBAAiB,GAAG,CAAC,OAAO,oBAAoB,CAAC,CAAC,MAAM,CAAC,CAAC;AAEtE,MAAM,WAAW,0BAA0B;IAC1C,OAAO,EAAE,OAAO,6BAA6B,CAAC;IAC9C,SAAS,EAAE,MAAM,CAAC;IAClB,MAAM,EAAE,iBAAiB,CAAC;IAC1B,MAAM,CAAC,EAAE,OAAO,CAAC;IACjB,MAAM,CAAC,EAAE;QACR,SAAS,CAAC,EAAE,MAAM,CAAC;QACnB,CAAC,GAAG,EAAE,MAAM,GAAG,OAAO,CAAC;KACvB,CAAC;CACF;AAED,MAAM,MAAM,wBAAwB,CAAC,CAAC,GAAG,OAAO,IAC7C;IACA,OAAO,EAAE,OAAO,6BAA6B,CAAC;IAC9C,SAAS,EAAE,MAAM,CAAC;IAClB,MAAM,CAAC,EAAE,iBAAiB,CAAC;IAC3B,OAAO,EAAE,IAAI,CAAC;IACd,IAAI,EAAE,CAAC,CAAC;CACP,GACD;IACA,OAAO,EAAE,OAAO,6BAA6B,CAAC;IAC9C,SAAS,EAAE,MAAM,CAAC;IAClB,MAAM,CAAC,EAAE,iBAAiB,CAAC;IAC3B,OAAO,EAAE,KAAK,CAAC;IACf,KAAK,EAAE;QACN,IAAI,EAAE,oBAAoB,CAAC;QAC3B,OAAO,EAAE,MAAM,CAAC;KAChB,CAAC;CACD,CAAC;AAEL,KAAK,oBAAoB,GACtB,iBAAiB,GACjB,gBAAgB,GAChB,qBAAqB,GACrB,oBAAoB,GACpB,mBAAmB,GACnB,kBAAkB,GAClB,WAAW,GACX,eAAe,CAAC;AAEnB,UAAU,QAAQ;IACjB,EAAE,CAAC,KAAK,EAAE,MAAM,EAAE,OAAO,EAAE,CAAC,IAAI,EAAE,OAAO,KAAK,IAAI,GAAG,CAAC,MAAM,IAAI,CAAC,GAAG,SAAS,CAAC;IAC9E,IAAI,CAAC,KAAK,EAAE,MAAM,EAAE,IAAI,EAAE,OAAO,GAAG,IAAI,CAAC;CACzC;AAED,MAAM,WAAW,qBAAqB;IACrC,kFAAkF;IAClF,GAAG,EAAE,MAAM,CAAC;IACZ,sCAAsC;IACtC,KAAK,EAAE,MAAM,CAAC;IACd,IAAI,CAAC,EAAE,MAAM,CAAC;IACd,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,SAAS,EAAE,MAAM,CAAC;IAClB,MAAM,EAAE;QAAE,KAAK,EAAE,MAAM,CAAC;QAAC,MAAM,EAAE,MAAM,CAAC;QAAC,KAAK,EAAE,MAAM,CAAA;KAAE,CAAC;IACzD,IAAI,CAAC,EAAE,MAAM,CAAC;CACd;AAED,MAAM,WAAW,sBAAsB;IACtC,OAAO,EAAE,CAAC,CAAC;IACX,OAAO,EAAE,qBAAqB,EAAE,CAAC;IACjC,+DAA+D;IAC/D,WAAW,EAAE,MAAM,CAAC;IACpB,OAAO,EAAE,MAAM,CAAC;CAChB;AAwMD,UAAU,gCAAgC;IACzC,MAAM,EAAE,QAAQ,CAAC;IACjB,UAAU,EAAE,MAAM,gBAAgB,GAAG,IAAI,CAAC;IAC1C,OAAO,EAAE,CACR,EAAE,EAAE,MAAM,EACV,MAAM,EAAE,kBAAkB,EAC1B,MAAM,EAAE,WAAW,EACnB,QAAQ,EAAE,CAAC,CAAC,MAAM,EAAE,eAAe,CAAC,OAAO,CAAC,KAAK,IAAI,CAAC,GAAG,SAAS,EAClE,GAAG,EAAE,gBAAgB,KACjB,OAAO,CAAC,eAAe,CAAC,OAAO,CAAC,CAAC,CAAC;IACvC,YAAY,CAAC,EAAE,MAAM,CAAC;IACtB,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,IAAI,CAAC,EAAE,CAAC,GAAG,EAAE,MAAM,EAAE,MAAM,CAAC,EAAE,MAAM,CAAC,OAAO,GAAG,CAAC,KAAK,OAAO,CAAC;IAC7D,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;IACnB,qFAAqF;IACrF,KAAK,CAAC,EAAE,aAAa,CAAC;CACtB;AAcD,wBAAgB,qBAAqB,CAAC,SAAS,EAAE,MAAM,GAAG,MAAM,CAE/D;AA2UD,wBAAgB,yBAAyB,CAAC,OAAO,EAAE,gCAAgC,GAAG;IACrF,SAAS,EAAE,CAAC,GAAG,CAAC,EAAE,gBAAgB,GAAG,IAAI,KAAK,IAAI,CAAC;IACnD,OAAO,EAAE,MAAM,IAAI,CAAC;CACpB,CA4BA","sourcesContent":["import * as path from \"node:path\";\nimport type { AgentToolResult } from \"@lpb-work/pi-agent-core\";\nimport type { ExtensionContext } from \"@lpb-work/pi-coding-agent\";\nimport { Compile } from \"typebox/compile\";\nimport { resolveAsyncRunLocation } from \"../runs/background/async-resume.ts\";\nimport { deliverStopRequest } from \"../runs/background/control-channel.ts\";\nimport { reconcileAsyncRun } from \"../runs/background/stale-run-reconciler.ts\";\nimport type { SubagentParamsLike } from \"../runs/foreground/subagent-executor.ts\";\nimport { resolveCurrentSessionId } from \"../shared/session-identity.ts\";\nimport {\n\ttype AsyncJobStep,\n\ttype Details,\n\tDIRS,\n\tSUBAGENT_ASYNC_COMPLETE_EVENT,\n\tSUBAGENT_LIFECYCLE_ARTIFACT_VERSION,\n\tSUBAGENT_PROCESS_TERMINAL_EVENT,\n\ttype SubagentState,\n} from \"../shared/types.ts\";\nimport { readStatus } from \"../shared/utils.ts\";\nimport { formatWorkflowJsonPreview } from \"../workflows/scripted-workflow.ts\";\nimport { normalizePublicSubagentExecution } from \"./public-execution.ts\";\nimport { SubagentParams } from \"./schemas.ts\";\n\nexport const SUBAGENT_RPC_PROTOCOL_VERSION = 1;\nexport const SUBAGENT_RPC_REQUEST_EVENT = \"subagents:rpc:v1:request\";\nexport const SUBAGENT_RPC_READY_EVENT = \"subagents:rpc:v1:ready\";\nexport const SUBAGENT_RPC_REPLY_EVENT_PREFIX = \"subagents:rpc:v1:reply:\";\n\nexport const SUBAGENT_RPC_METHODS = [\"ping\", \"status\", \"spawn\", \"steer\", \"interrupt\", \"stop\", \"resume\"] as const;\nexport type SubagentRpcMethod = (typeof SUBAGENT_RPC_METHODS)[number];\n\nexport interface SubagentRpcRequestEnvelope {\n\tversion: typeof SUBAGENT_RPC_PROTOCOL_VERSION;\n\trequestId: string;\n\tmethod: SubagentRpcMethod;\n\tparams?: unknown;\n\tsource?: {\n\t\textension?: string;\n\t\t[key: string]: unknown;\n\t};\n}\n\nexport type SubagentRpcReplyEnvelope<T = unknown> =\n\t| {\n\t\t\tversion: typeof SUBAGENT_RPC_PROTOCOL_VERSION;\n\t\t\trequestId: string;\n\t\t\tmethod?: SubagentRpcMethod;\n\t\t\tsuccess: true;\n\t\t\tdata: T;\n\t  }\n\t| {\n\t\t\tversion: typeof SUBAGENT_RPC_PROTOCOL_VERSION;\n\t\t\trequestId: string;\n\t\t\tmethod?: SubagentRpcMethod;\n\t\t\tsuccess: false;\n\t\t\terror: {\n\t\t\t\tcode: SubagentRpcErrorCode;\n\t\t\t\tmessage: string;\n\t\t\t};\n\t  };\n\ntype SubagentRpcErrorCode =\n\t| \"invalid_request\"\n\t| \"invalid_params\"\n\t| \"unsupported_version\"\n\t| \"unsupported_method\"\n\t| \"no_active_session\"\n\t| \"execution_failed\"\n\t| \"not_found\"\n\t| \"invalid_state\";\n\ninterface EventBus {\n\ton(event: string, handler: (data: unknown) => void): (() => void) | undefined;\n\temit(event: string, data: unknown): void;\n}\n\nexport interface SubagentRpcFleetEntry {\n\t/** Opaque key for client-side reconciliation; never a run or async identifier. */\n\tkey: string;\n\t/** Resolved child agent/role name. */\n\tagent: string;\n\trole?: string;\n\tmodel?: string;\n\teffort?: string;\n\tstartedAt: number;\n\ttokens: { input: number; output: number; total: number };\n\tgoal?: string;\n}\n\nexport interface SubagentRpcFleetStatus {\n\tversion: 1;\n\tentries: SubagentRpcFleetEntry[];\n\t/** Total active children before the bounded entries window. */\n\ttotalActive: number;\n\tomitted: number;\n}\n\nconst MAX_FLEET_ENTRIES = 16;\nconst MAX_FLEET_CANDIDATES = 256;\nconst MAX_AGENT_LENGTH = 96;\nconst MAX_GOAL_LENGTH = 512;\nconst MAX_METADATA_LENGTH = 128;\n\nfunction displayText(value: unknown, maxLength: number): string | undefined {\n\tif (typeof value !== \"string\") return undefined;\n\t// Strip complete CSI/OSC/DCS/APC/PM strings and C1 controls before collapsing\n\t// whitespace; never leave CSI parameters behind after removing ESC.\n\tconst normalized = value\n\t\t.slice(0, 4_096)\n\t\t.replace(\n\t\t\t/\\x1b\\[[0-?]*[ -/]*[@-~]|\\x9b[0-?]*[ -/]*[@-~]|\\x1b][\\s\\S]*?(?:\\x07|\\x1b\\\\)|\\x1b[PX^_][\\s\\S]*?\\x1b\\\\|[\\u0000-\\u001f\\u007f-\\u009f]/g,\n\t\t\t\" \",\n\t\t)\n\t\t.replace(/\\s+/g, \" \")\n\t\t.trim();\n\treturn normalized ? normalized.slice(0, maxLength) : undefined;\n}\n\nfunction publicTokens(value: unknown): { input: number; output: number; total: number } {\n\tconst record = isRecord(value) ? value : {};\n\tconst count = (field: \"input\" | \"output\" | \"total\") => {\n\t\tconst raw = record[field];\n\t\treturn typeof raw === \"number\" && Number.isFinite(raw) && raw >= 0\n\t\t\t? Math.min(Number.MAX_SAFE_INTEGER, Math.floor(raw))\n\t\t\t: 0;\n\t};\n\tconst input = count(\"input\");\n\tconst output = count(\"output\");\n\tconst sum = Math.min(Number.MAX_SAFE_INTEGER, input + output);\n\treturn { input, output, total: Math.max(sum, count(\"total\")) };\n}\n\nfunction activeState(value: unknown): boolean {\n\treturn value === \"running\" || value === \"queued\" || value === \"pending\";\n}\n\ninterface FleetKeyState {\n\tsessionId: string | null;\n\tnext: number;\n\tkeys: Map<string, string>;\n}\n\ninterface FleetCandidate {\n\tinternalKey: string;\n\tagent: unknown;\n\trole?: unknown;\n\tmodel?: unknown;\n\teffort?: unknown;\n\tstartedAt: unknown;\n\ttokens?: unknown;\n\tgoal?: unknown;\n}\n\nfunction buildFleetStatus(\n\tstate: SubagentState | undefined,\n\tkeyState: FleetKeyState,\n\tsessionId: string | null | undefined,\n): SubagentRpcFleetStatus {\n\tconst authoritativeSessionId = sessionId ?? null;\n\tif (keyState.sessionId !== authoritativeSessionId) {\n\t\tkeyState.sessionId = authoritativeSessionId;\n\t\tkeyState.next = 0;\n\t\tkeyState.keys.clear();\n\t}\n\tif (!state || !authoritativeSessionId || state.currentSessionId !== authoritativeSessionId) {\n\t\tkeyState.keys.clear();\n\t\treturn { version: 1, entries: [], totalActive: 0, omitted: 0 };\n\t}\n\n\tlet totalActive = 0;\n\tconst candidates: FleetCandidate[] = [];\n\tconst addCandidate = (candidate: FleetCandidate) => {\n\t\ttotalActive += 1;\n\t\tif (candidates.length < MAX_FLEET_CANDIDATES) candidates.push(candidate);\n\t};\n\tfor (const control of state.foregroundControls.values()) {\n\t\tif (control.sessionId !== authoritativeSessionId) continue;\n\t\tif (control.activeChildren?.size) {\n\t\t\tfor (const child of control.activeChildren.values())\n\t\t\t\taddCandidate({\n\t\t\t\t\tinternalKey: `foreground:${control.runId}:${child.index}`,\n\t\t\t\t\tagent: child.agent,\n\t\t\t\t\tmodel: child.model,\n\t\t\t\t\teffort: child.thinking,\n\t\t\t\t\tstartedAt: child.startedAt,\n\t\t\t\t\ttokens: { input: child.inputTokens ?? 0, output: child.outputTokens ?? 0, total: child.tokens ?? 0 },\n\t\t\t\t\tgoal: child.description ?? control.description,\n\t\t\t\t});\n\t\t} else {\n\t\t\taddCandidate({\n\t\t\t\tinternalKey: `foreground:${control.runId}:${control.currentIndex ?? 0}`,\n\t\t\t\tagent: control.currentAgent ?? control.mode,\n\t\t\t\tmodel: control.model,\n\t\t\t\teffort: control.thinking,\n\t\t\t\tstartedAt: control.startedAt,\n\t\t\t\ttokens: { input: control.inputTokens ?? 0, output: control.outputTokens ?? 0, total: control.tokens ?? 0 },\n\t\t\t\tgoal: control.description,\n\t\t\t});\n\t\t}\n\t}\n\tfor (const job of state.asyncJobs.values()) {\n\t\tif (job.sessionId !== authoritativeSessionId || !activeState(job.status)) continue;\n\t\tconst startedAt = job.startedAt ?? job.updatedAt;\n\t\tif (job.mode === \"workflow\") {\n\t\t\tconst latestEmit = job.workflow?.emits?.length\n\t\t\t\t? formatWorkflowJsonPreview(job.workflow.emits.at(-1), 120)\n\t\t\t\t: undefined;\n\t\t\taddCandidate({\n\t\t\t\tinternalKey: `async:${job.asyncId}`,\n\t\t\t\tagent: \"workflow\",\n\t\t\t\tstartedAt,\n\t\t\t\ttokens: job.totalTokens,\n\t\t\t\tgoal: latestEmit !== undefined ? `latest emit: ${latestEmit}` : job.description,\n\t\t\t});\n\t\t\tcontinue;\n\t\t}\n\t\tconst steps: AsyncJobStep[] | undefined = job.steps?.length\n\t\t\t? job.steps\n\t\t\t: job.agents?.map((agent, index) => ({\n\t\t\t\t\tagent,\n\t\t\t\t\tindex,\n\t\t\t\t\tstatus: job.status === \"queued\" ? \"pending\" : \"running\",\n\t\t\t\t}));\n\t\tif (!steps?.length) {\n\t\t\taddCandidate({\n\t\t\t\tinternalKey: `async:${job.asyncId}`,\n\t\t\t\tagent: job.mode ?? \"subagent\",\n\t\t\t\tstartedAt,\n\t\t\t\ttokens: job.totalTokens,\n\t\t\t\tgoal: job.description,\n\t\t\t});\n\t\t\tcontinue;\n\t\t}\n\t\tfor (const [offset, step] of steps.entries()) {\n\t\t\tif (!activeState(step.status)) continue;\n\t\t\tconst index = step.index ?? offset;\n\t\t\tif (\n\t\t\t\tstep.status === \"pending\" &&\n\t\t\t\tjob.mode === \"chain\" &&\n\t\t\t\t!job.activeParallelGroup &&\n\t\t\t\tindex !== (job.currentStep ?? 0)\n\t\t\t)\n\t\t\t\tcontinue;\n\t\t\taddCandidate({\n\t\t\t\tinternalKey: `async:${job.asyncId}:${index}`,\n\t\t\t\tagent: step.agent,\n\t\t\t\trole: step.label,\n\t\t\t\tmodel: step.model,\n\t\t\t\teffort: step.thinking,\n\t\t\t\tstartedAt: step.startedAt ?? startedAt,\n\t\t\t\ttokens: step.tokens ?? (steps.length === 1 ? job.totalTokens : undefined),\n\t\t\t\tgoal: job.description,\n\t\t\t});\n\t\t}\n\t}\n\n\tcandidates.sort((left, right) => {\n\t\tconst leftStarted = typeof left.startedAt === \"number\" ? left.startedAt : Number.MAX_SAFE_INTEGER;\n\t\tconst rightStarted = typeof right.startedAt === \"number\" ? right.startedAt : Number.MAX_SAFE_INTEGER;\n\t\treturn leftStarted - rightStarted || left.internalKey.localeCompare(right.internalKey);\n\t});\n\tconst activeKeys = new Set(candidates.map((candidate) => candidate.internalKey));\n\tconst entries: SubagentRpcFleetEntry[] = [];\n\tfor (const candidate of candidates) {\n\t\tif (entries.length >= MAX_FLEET_ENTRIES) break;\n\t\tconst agent = displayText(candidate.agent, MAX_AGENT_LENGTH);\n\t\tconst startedAt = candidate.startedAt;\n\t\tif (!agent || typeof startedAt !== \"number\" || !Number.isSafeInteger(startedAt) || startedAt < 0) continue;\n\t\tlet key = keyState.keys.get(candidate.internalKey);\n\t\tif (!key) {\n\t\t\tkey = `fleet-${++keyState.next}`;\n\t\t\tkeyState.keys.set(candidate.internalKey, key);\n\t\t}\n\t\tconst role = displayText(candidate.role, MAX_AGENT_LENGTH);\n\t\tconst model = displayText(candidate.model, MAX_METADATA_LENGTH);\n\t\tconst effort = displayText(candidate.effort, MAX_METADATA_LENGTH);\n\t\tconst goal = displayText(candidate.goal, MAX_GOAL_LENGTH);\n\t\tentries.push({\n\t\t\tkey,\n\t\t\tagent,\n\t\t\t...(role ? { role } : {}),\n\t\t\t...(model ? { model } : {}),\n\t\t\t...(effort ? { effort } : {}),\n\t\t\tstartedAt,\n\t\t\ttokens: publicTokens(candidate.tokens),\n\t\t\t...(goal ? { goal } : {}),\n\t\t});\n\t}\n\tfor (const internalKey of keyState.keys.keys()) {\n\t\tif (!activeKeys.has(internalKey)) keyState.keys.delete(internalKey);\n\t}\n\tconst omitted = Math.max(0, totalActive - entries.length);\n\treturn { version: 1, entries, totalActive, omitted };\n}\n\ninterface RegisterSubagentRpcBridgeOptions {\n\tevents: EventBus;\n\tgetContext: () => ExtensionContext | null;\n\texecute: (\n\t\tid: string,\n\t\tparams: SubagentParamsLike,\n\t\tsignal: AbortSignal,\n\t\tonUpdate: ((result: AgentToolResult<Details>) => void) | undefined,\n\t\tctx: ExtensionContext,\n\t) => Promise<AgentToolResult<Details>>;\n\tasyncDirRoot?: string;\n\tresultsDir?: string;\n\tkill?: (pid: number, signal?: NodeJS.Signals | 0) => boolean;\n\tnow?: () => number;\n\t/** Native live state, projected into the optional public fleet-status capability. */\n\tstate?: SubagentState;\n}\n\nclass SubagentRpcError extends Error {\n\treadonly code: SubagentRpcErrorCode;\n\n\tconstructor(code: SubagentRpcErrorCode, message: string) {\n\t\tsuper(message);\n\t\tthis.name = \"SubagentRpcError\";\n\t\tthis.code = code;\n\t}\n}\n\nconst subagentParamsValidator = Compile(SubagentParams);\n\nexport function subagentRpcReplyEvent(requestId: string): string {\n\treturn `${SUBAGENT_RPC_REPLY_EVENT_PREFIX}${requestId}`;\n}\n\nfunction isRecord(value: unknown): value is Record<string, unknown> {\n\treturn Boolean(value) && typeof value === \"object\" && !Array.isArray(value);\n}\n\nfunction assertRequestId(value: unknown): string {\n\tif (typeof value !== \"string\" || value.trim().length === 0 || /[\\r\\n]/.test(value)) {\n\t\tthrow new SubagentRpcError(\"invalid_request\", \"RPC requestId must be a non-empty string without newlines.\");\n\t}\n\treturn value;\n}\n\nfunction assertRecordParams(params: unknown, method: SubagentRpcMethod): Record<string, unknown> {\n\tif (params === undefined) return {};\n\tif (!isRecord(params)) throw new SubagentRpcError(\"invalid_params\", `RPC ${method} params must be an object.`);\n\treturn params;\n}\n\nfunction assertSubagentParams(params: SubagentParamsLike, label: string): void {\n\tif (subagentParamsValidator.Check(params)) return;\n\tconst messages = [...subagentParamsValidator.Errors(params)].slice(0, 4).map((error) => error.message);\n\tthrow new SubagentRpcError(\"invalid_params\", `${label}: ${messages.join(\"; \") || \"invalid subagent parameters\"}`);\n}\n\nfunction textFromToolResult(result: AgentToolResult<Details>): string {\n\treturn result.content\n\t\t.filter((part): part is { type: \"text\"; text: string } => part.type === \"text\")\n\t\t.map((part) => part.text)\n\t\t.join(\"\\n\");\n}\n\ntype ToolResultWithError = AgentToolResult<Details> & { isError?: boolean };\n\nfunction dataFromToolResult(result: ToolResultWithError): { text: string; details?: Details; isError?: boolean } {\n\treturn {\n\t\ttext: textFromToolResult(result),\n\t\t...(result.details ? { details: result.details } : {}),\n\t\t...(result.isError ? { isError: true } : {}),\n\t};\n}\n\nfunction failIfToolError(result: ToolResultWithError): void {\n\tif (!result.isError) return;\n\tthrow new SubagentRpcError(\"execution_failed\", textFromToolResult(result) || \"Subagent RPC execution failed.\");\n}\n\nfunction normalizeTargetParams(\n\tparams: unknown,\n\tmethod: SubagentRpcMethod,\n): Pick<SubagentParamsLike, \"id\" | \"runId\" | \"dir\" | \"index\"> {\n\tconst input = assertRecordParams(params, method);\n\tconst output: Pick<SubagentParamsLike, \"id\" | \"runId\" | \"dir\" | \"index\"> = {};\n\tif (input.id !== undefined) output.id = input.id as string;\n\tif (input.runId !== undefined) output.runId = input.runId as string;\n\tif (input.dir !== undefined) output.dir = input.dir as string;\n\tif (input.index !== undefined) output.index = input.index as number;\n\treturn output;\n}\n\nfunction sessionData(ctx: ExtensionContext | null): { cwd?: string; sessionId?: string; sessionFile?: string | null } {\n\tif (!ctx) return {};\n\treturn {\n\t\tcwd: ctx.cwd,\n\t\tsessionId: ctx.sessionManager.getSessionId() ?? undefined,\n\t\tsessionFile: ctx.sessionManager.getSessionFile() ?? null,\n\t};\n}\n\nfunction pingData(ctx: ExtensionContext | null) {\n\treturn {\n\t\tversion: SUBAGENT_RPC_PROTOCOL_VERSION,\n\t\tmethods: [...SUBAGENT_RPC_METHODS],\n\t\tcapabilities: {\n\t\t\tstatus: true,\n\t\t\tfleetStatus: { version: 1 },\n\t\t\tasyncSpawn: true,\n\t\t\tsteer: true,\n\t\t\tnonRecoveringSteer: true,\n\t\t\tinterrupt: true,\n\t\t\tstop: true,\n\t\t\tresume: true,\n\t\t\tlaunchResolvedExtensions: { version: 1, source: \"launch-resolved\" },\n\t\t\truntimeAcknowledgedExtensions: {\n\t\t\t\tversion: 1,\n\t\t\t\tsource: \"child-runtime\",\n\t\t\t\tevent: \"subagent:acknowledge-extension\",\n\t\t\t},\n\t\t\tprocessTerminalProof: { version: 1, lifecycleArtifactVersion: SUBAGENT_LIFECYCLE_ARTIFACT_VERSION },\n\t\t},\n\t\tevents: {\n\t\t\tready: SUBAGENT_RPC_READY_EVENT,\n\t\t\trequest: SUBAGENT_RPC_REQUEST_EVENT,\n\t\t\treplyPrefix: SUBAGENT_RPC_REPLY_EVENT_PREFIX,\n\t\t\tasyncComplete: SUBAGENT_ASYNC_COMPLETE_EVENT,\n\t\t\tprocessTerminal: SUBAGENT_PROCESS_TERMINAL_EVENT,\n\t\t},\n\t\tsession: sessionData(ctx),\n\t};\n}\n\nasync function executeChecked(\n\toptions: RegisterSubagentRpcBridgeOptions,\n\tctx: ExtensionContext,\n\trequestId: string,\n\tmethod: SubagentRpcMethod,\n\tparams: SubagentParamsLike,\n): Promise<{ text: string; details?: Details; isError?: boolean }> {\n\tassertSubagentParams(params, `RPC ${method} params`);\n\tconst controller = new AbortController();\n\tconst result = await options.execute(`rpc-${method}-${requestId}`, params, controller.signal, undefined, ctx);\n\tfailIfToolError(result);\n\treturn dataFromToolResult(result);\n}\n\nfunction spawnParams(params: unknown): SubagentParamsLike {\n\tconst input = assertRecordParams(params, \"spawn\");\n\tconst normalized = normalizePublicSubagentExecution(input);\n\tif (!normalized.ok) throw new SubagentRpcError(\"invalid_params\", normalized.error);\n\tif (normalized.params.action !== undefined) {\n\t\tthrow new SubagentRpcError(\n\t\t\t\"invalid_params\",\n\t\t\t\"RPC spawn does not accept management/control actions. Use status or interrupt RPC methods instead.\",\n\t\t);\n\t}\n\tif (input.async === false) {\n\t\tthrow new SubagentRpcError(\n\t\t\t\"invalid_params\",\n\t\t\t\"RPC spawn only supports detached async launches; omit async or set async: true.\",\n\t\t);\n\t}\n\treturn { ...(normalized.params as SubagentParamsLike), async: true };\n}\n\nfunction steerParams(params: unknown): SubagentParamsLike {\n\tconst input = assertRecordParams(params, \"steer\");\n\tif (typeof input.message !== \"string\" || !input.message.trim())\n\t\tthrow new SubagentRpcError(\"invalid_params\", \"RPC steer requires a non-empty message.\");\n\tconst target = normalizeTargetParams(input, \"steer\");\n\tif (!target.id && !target.runId && !target.dir)\n\t\tthrow new SubagentRpcError(\"invalid_params\", \"RPC steer requires id, runId, or dir.\");\n\tif (input.mode !== undefined && input.mode !== \"steer\" && input.mode !== \"follow_up\" && input.mode !== \"auto\")\n\t\tthrow new SubagentRpcError(\"invalid_params\", \"RPC steer mode must be steer, follow_up, or auto.\");\n\treturn {\n\t\taction: \"steer\",\n\t\t...target,\n\t\tmessage: input.message.trim(),\n\t\t...(typeof input.mode === \"string\" ? { mode: input.mode as \"steer\" | \"follow_up\" | \"auto\" } : {}),\n\t\tsteeringRecovery: false,\n\t};\n}\n\nfunction resumeParams(params: unknown): SubagentParamsLike {\n\tconst input = assertRecordParams(params, \"resume\");\n\tif (typeof input.message !== \"string\" || !input.message.trim())\n\t\tthrow new SubagentRpcError(\"invalid_params\", \"RPC resume requires a non-empty message.\");\n\tconst target = normalizeTargetParams(input, \"resume\");\n\tif (!target.id && !target.runId && !target.dir)\n\t\tthrow new SubagentRpcError(\"invalid_params\", \"RPC resume requires id, runId, or dir.\");\n\tif (input.output !== undefined && (typeof input.output !== \"string\" || !input.output.trim()))\n\t\tthrow new SubagentRpcError(\"invalid_params\", \"RPC resume output must be a non-empty path.\");\n\tif (input.outputMode !== undefined && input.outputMode !== \"file-only\")\n\t\tthrow new SubagentRpcError(\"invalid_params\", \"RPC resume supports only file-only output mode.\");\n\treturn {\n\t\taction: \"resume\",\n\t\t...target,\n\t\tmessage: input.message.trim(),\n\t\t...(typeof input.output === \"string\" ? { output: input.output.trim(), outputMode: \"file-only\" } : {}),\n\t};\n}\n\nfunction stopAsyncRun(\n\tparams: unknown,\n\toptions: RegisterSubagentRpcBridgeOptions,\n\tctx: ExtensionContext,\n): { runId: string; asyncDir: string; previousState: string; state: \"stopping\"; message: string } {\n\tconst target = normalizeTargetParams(params, \"stop\");\n\tassertSubagentParams({ action: \"status\", ...target }, \"RPC stop target params\");\n\tconst asyncDirRoot = options.asyncDirRoot ?? DIRS.async;\n\tconst resultsDir = options.resultsDir ?? DIRS.results;\n\tlet location;\n\ttry {\n\t\tlocation = resolveAsyncRunLocation(target, asyncDirRoot, resultsDir);\n\t} catch (error) {\n\t\tthrow new SubagentRpcError(\"invalid_params\", error instanceof Error ? error.message : String(error));\n\t}\n\tif (!location.asyncDir) {\n\t\tthrow new SubagentRpcError(\n\t\t\t\"not_found\",\n\t\t\t\"Async run not found or already completed; stop requires a live async run directory.\",\n\t\t);\n\t}\n\n\tconst currentSessionId = resolveCurrentSessionId(ctx.sessionManager);\n\tconst initialStatus = readStatus(location.asyncDir);\n\tconst initialRunId = initialStatus?.runId ?? location.resolvedId ?? path.basename(location.asyncDir);\n\tif (!initialStatus)\n\t\tthrow new SubagentRpcError(\"not_found\", `Status file not found for async run '${initialRunId}'.`);\n\tif (!currentSessionId || initialStatus.sessionId !== currentSessionId) {\n\t\tthrow new SubagentRpcError(\"not_found\", `Async run '${initialRunId}' was not found in the active session.`);\n\t}\n\n\tlet status;\n\ttry {\n\t\tstatus = reconcileAsyncRun(location.asyncDir, { resultsDir, kill: options.kill, now: options.now }).status;\n\t} catch (error) {\n\t\tthrow new SubagentRpcError(\"execution_failed\", error instanceof Error ? error.message : String(error));\n\t}\n\tconst runId = status?.runId ?? initialRunId;\n\tif (!status) throw new SubagentRpcError(\"not_found\", `Status file not found for async run '${runId}'.`);\n\tif (status.sessionId !== currentSessionId) {\n\t\tthrow new SubagentRpcError(\"not_found\", `Async run '${runId}' was not found in the active session.`);\n\t}\n\tif (status.state !== \"running\") {\n\t\tthrow new SubagentRpcError(\n\t\t\t\"invalid_state\",\n\t\t\t`Async run ${runId} is ${status.state}; stop only supports running async runs.`,\n\t\t);\n\t}\n\n\ttry {\n\t\tdeliverStopRequest({\n\t\t\tasyncDir: location.asyncDir,\n\t\t\tpid: status.pid,\n\t\t\tkill: options.kill,\n\t\t\tnow: options.now,\n\t\t\tsource: \"rpc-stop\",\n\t\t});\n\t} catch (error) {\n\t\tthrow new SubagentRpcError(\"execution_failed\", error instanceof Error ? error.message : String(error));\n\t}\n\n\treturn {\n\t\trunId,\n\t\tasyncDir: location.asyncDir,\n\t\tpreviousState: status.state,\n\t\tstate: \"stopping\",\n\t\tmessage: `Stop requested for async run ${runId}.`,\n\t};\n}\n\nasync function handleRequest(\n\trequest: SubagentRpcRequestEnvelope,\n\toptions: RegisterSubagentRpcBridgeOptions,\n\tfleetKeys: FleetKeyState,\n): Promise<unknown> {\n\tconst ctx = options.getContext();\n\tif (request.method === \"ping\") return pingData(ctx);\n\tif (!ctx) throw new SubagentRpcError(\"no_active_session\", \"No active extension context for subagent RPC.\");\n\n\tif (request.method === \"spawn\") {\n\t\treturn executeChecked(options, ctx, request.requestId, request.method, spawnParams(request.params));\n\t}\n\tif (request.method === \"status\") {\n\t\tconst status = await executeChecked(options, ctx, request.requestId, request.method, {\n\t\t\taction: \"status\",\n\t\t\t...normalizeTargetParams(request.params, \"status\"),\n\t\t});\n\t\treturn {\n\t\t\t...status,\n\t\t\tfleet: buildFleetStatus(options.state, fleetKeys, resolveCurrentSessionId(ctx.sessionManager)),\n\t\t};\n\t}\n\tif (request.method === \"steer\") {\n\t\treturn executeChecked(options, ctx, request.requestId, request.method, steerParams(request.params));\n\t}\n\tif (request.method === \"interrupt\") {\n\t\treturn executeChecked(options, ctx, request.requestId, request.method, {\n\t\t\taction: \"interrupt\",\n\t\t\t...normalizeTargetParams(request.params, \"interrupt\"),\n\t\t});\n\t}\n\tif (request.method === \"stop\") {\n\t\treturn stopAsyncRun(request.params, options, ctx);\n\t}\n\tif (request.method === \"resume\") {\n\t\treturn executeChecked(options, ctx, request.requestId, request.method, resumeParams(request.params));\n\t}\n\tthrow new SubagentRpcError(\"unsupported_method\", `Unsupported subagent RPC method: ${String(request.method)}`);\n}\n\nfunction parseRequest(raw: unknown): SubagentRpcRequestEnvelope {\n\tif (!isRecord(raw)) throw new SubagentRpcError(\"invalid_request\", \"Subagent RPC request must be an object.\");\n\tconst requestId = assertRequestId(raw.requestId);\n\tif (raw.version !== SUBAGENT_RPC_PROTOCOL_VERSION) {\n\t\tthrow new SubagentRpcError(\"unsupported_version\", `Unsupported subagent RPC version: ${String(raw.version)}.`);\n\t}\n\tif (typeof raw.method !== \"string\" || !(SUBAGENT_RPC_METHODS as readonly string[]).includes(raw.method)) {\n\t\tthrow new SubagentRpcError(\"unsupported_method\", `Unsupported subagent RPC method: ${String(raw.method)}.`);\n\t}\n\treturn {\n\t\tversion: SUBAGENT_RPC_PROTOCOL_VERSION,\n\t\trequestId,\n\t\tmethod: raw.method as SubagentRpcMethod,\n\t\t...(raw.params !== undefined ? { params: raw.params } : {}),\n\t\t...(isRecord(raw.source) ? { source: raw.source as SubagentRpcRequestEnvelope[\"source\"] } : {}),\n\t};\n}\n\nfunction safeReplyRequestId(raw: unknown): string {\n\tif (!isRecord(raw)) return \"unknown\";\n\tconst requestId = raw.requestId;\n\treturn typeof requestId === \"string\" && requestId.trim().length > 0 && !/[\\r\\n]/.test(requestId)\n\t\t? requestId\n\t\t: \"unknown\";\n}\n\nfunction errorReply(raw: unknown, error: unknown): SubagentRpcReplyEnvelope {\n\tconst requestId = safeReplyRequestId(raw);\n\tconst method =\n\t\tisRecord(raw) &&\n\t\ttypeof raw.method === \"string\" &&\n\t\t(SUBAGENT_RPC_METHODS as readonly string[]).includes(raw.method)\n\t\t\t? (raw.method as SubagentRpcMethod)\n\t\t\t: undefined;\n\tconst rpcError =\n\t\terror instanceof SubagentRpcError\n\t\t\t? error\n\t\t\t: new SubagentRpcError(\"execution_failed\", error instanceof Error ? error.message : String(error));\n\treturn {\n\t\tversion: SUBAGENT_RPC_PROTOCOL_VERSION,\n\t\trequestId,\n\t\t...(method ? { method } : {}),\n\t\tsuccess: false,\n\t\terror: {\n\t\t\tcode: rpcError.code,\n\t\t\tmessage: rpcError.message,\n\t\t},\n\t};\n}\n\nexport function registerSubagentRpcBridge(options: RegisterSubagentRpcBridgeOptions): {\n\temitReady: (ctx?: ExtensionContext | null) => void;\n\tdispose: () => void;\n} {\n\tconst fleetKeys: FleetKeyState = { sessionId: null, next: 0, keys: new Map() };\n\tconst unsubscribe = options.events.on(SUBAGENT_RPC_REQUEST_EVENT, async (raw) => {\n\t\tlet request: SubagentRpcRequestEnvelope | undefined;\n\t\ttry {\n\t\t\trequest = parseRequest(raw);\n\t\t\tconst data = await handleRequest(request, options, fleetKeys);\n\t\t\toptions.events.emit(subagentRpcReplyEvent(request.requestId), {\n\t\t\t\tversion: SUBAGENT_RPC_PROTOCOL_VERSION,\n\t\t\t\trequestId: request.requestId,\n\t\t\t\tmethod: request.method,\n\t\t\t\tsuccess: true,\n\t\t\t\tdata,\n\t\t\t} satisfies SubagentRpcReplyEnvelope);\n\t\t} catch (error) {\n\t\t\tconst reply = errorReply(request ?? raw, error);\n\t\t\toptions.events.emit(subagentRpcReplyEvent(reply.requestId), reply);\n\t\t}\n\t});\n\n\treturn {\n\t\temitReady: (ctx) => {\n\t\t\toptions.events.emit(SUBAGENT_RPC_READY_EVENT, pingData(ctx ?? options.getContext()));\n\t\t},\n\t\tdispose: () => {\n\t\t\tif (typeof unsubscribe === \"function\") unsubscribe();\n\t\t},\n\t};\n}\n"]}