import { resolve } from "node:path"; import { fileURLToPath } from "node:url"; import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-agent"; import { createPiAgentsClient, PROTOCOL_VERSION, RPC_REPLY_PREFIX, RPC_REQUEST_CHANNEL, RUN_EVENT_CHANNEL, type StartRpcParams } from "pi-agents/api"; import { buildModelCatalog, resolveModelReference } from "pi-agents/src/catalog/models.js"; import { createSubprocessSpawnEngine } from "pi-agents/src/engine/subprocess.js"; import type { SpawnEngine } from "pi-agents/src/engine/types.js"; import type { AgentNode, FlowNode } from "pi-agents/src/model/ast.js"; import { validateFlow } from "pi-agents/src/model/validate.js"; import { RunManager } from "pi-agents/src/run/runs.js"; import type { RunEvent } from "pi-agents/src/run/events.js"; import type { TaskEnvelope, Topology, UltraConfig } from "../types.js"; import { fitEnvelope } from "../context/envelope.js"; import { agentResultSchema } from "../context/result-contracts.js"; import { boardBatchSchema, type WarRoomWorkerBinding, WAR_ROOM_MEMBERS } from "../workflows/war-room.js"; import { SCOPED_READ_TOOLS, SCOPED_WRITE_TOOLS, stripWorkerScopeMarker, workerScopeFromTask, workerScopeMarker, workerWarRoomFromTask, workerWarRoomMarker, WORKER_ROOT_ENV, WORKER_SCOPE_ENV, WORKER_WAR_ROOM_ENV } from "../security/worker-scope.js"; import { writerModelForLens } from "../models/model-policy.js"; type StartWithBudgets = StartRpcParams & { budgets?: UltraConfig["piAgentsBudgets"] }; type RpcRequest = { protocol?: unknown; id?: unknown; op?: unknown; params?: unknown; caller?: unknown }; function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } function taskFor(envelope: TaskEnvelope): string { return `Complete this bounded task. Do not delegate. Use only guarded_bash for repository commands; no other shell command is permitted. Its output is supporting evidence; the controller remains authoritative for verification. Submit the required structured result, including attributionManifest with used/rejected input fact IDs, accepted decisions, and changed paths.\n\n${workerScopeMarker(envelope.scope)}\n\n${JSON.stringify(envelope)}`; } function scoutTaskFor(envelope: TaskEnvelope, lens: string): string { const { acceptance: _acceptance, ...readOnlyEnvelope } = envelope; return `You are a read-only scout focused on ${lens}. Inspect the bounded task and submit concise structured facts, evidence, and unknowns for the writer. Do not edit files, invoke shell commands, or attempt the acceptance command. Use only the available read-only tools.\n\n${workerScopeMarker(envelope.scope)}\n\n${JSON.stringify({ ...readOnlyEnvelope, evidenceTarget: "read-only evidence" })}`; } function warRoomTaskFor(envelope: TaskEnvelope, author: string, round: number, feed?: string, binding?: WarRoomWorkerBinding): string { const { acceptance: _acceptance, ...readOnlyEnvelope } = envelope; return `You are the read-only War Room ${author} in round ${round}. You may call ultra_signal for evidence, cursor-linked challenges, concrete specialist requests, or blockers; target one member when steering is useful and never broadcast. Communicate only by returning exactly one typed blackboard event (hypothesis, evidence, challenge, confirmation, request, or blocker) authored exactly as ${author}. Challenge and confirmation events must include targetCursor for an existing hypothesis; confirmation must include evidence. Workers cannot emit decisions. Do not return chat, prose outside JSON, AgentResult, edits, or shell commands.${feed ? ` Read only these relevant events after your cursor: ${feed}.` : " Establish the smallest useful initial hypothesis or evidence."}\n\n${workerScopeMarker(envelope.scope)}${binding ? `\n${workerWarRoomMarker(binding)}` : ""}\n\n${JSON.stringify(readOnlyEnvelope)}`; } export type WarRoomFeeds = Partial>; export type WarRoomBindings = Partial>; export function flowFor(topology: Topology, envelope: TaskEnvelope, config: UltraConfig, cwd: string, writerModel?: string): FlowNode { const readOnly = [...SCOPED_READ_TOOLS]; const writer = [...readOnly, ...SCOPED_WRITE_TOOLS, "guarded_bash"]; const assignedEnvelope = (nodeId: string): TaskEnvelope => fitEnvelope({ ...envelope, trace: { ...envelope.trace, nodeId } }, config.context.envelopeHardCapTokens); const agent = (model: string, thinking: "low" | "medium" | "high" | "xhigh", tools: string[], nodeId: string): AgentNode => { const assigned = assignedEnvelope(nodeId); return { kind: "agent", task: taskFor(assigned), json: agentResultSchema(assigned, assigned.trace.nodeId, false, "writer"), model, thinking, tools, cwd }; }; const scout = (model: string, thinking: "low" | "medium" | "high" | "xhigh", lens: string, nodeId: string): AgentNode => { const assigned = assignedEnvelope(nodeId); return { kind: "agent", task: scoutTaskFor(assigned, lens), json: agentResultSchema(assigned, assigned.trace.nodeId, false, "scout", true), model, thinking, tools: readOnly, cwd }; }; if (topology === "deep" && envelope.lens === "first-proven-result") { const writerAgent = agent(writerModel ?? writerModelForLens(config, envelope.lens, config.deep.model), "high", writer, "$.steps[1]"); return { kind: "sequence", steps: [ { kind: "parallel", as: "candidate", mode: "any", onError: "collect", concurrency: 2, branches: { reproduce: scout(config.scout.model, "high", "independently reproduce and evidence the root cause", "$.steps[0].branches.reproduce"), workingPath: scout(config.scout.model, "high", "independently find and evidence a working path", "$.steps[0].branches.workingPath"), } }, { ...writerAgent, task: `${writerAgent.task}\n\nUse the first proven candidate below; the other speculative branch was cancelled.\n{candidate}` }, ], }; } if (topology === "deep") return agent(writerModel ?? writerModelForLens(config, envelope.lens, config.deep.model), "high", writer, "$"); if (topology === "warroom") return warRoomRoundFlow(envelope, config, cwd, 1); if (topology === "swarm") return agent(writerModel ?? writerModelForLens(config, envelope.lens), "high", writer, "$"); return agent(writerModelForLens(config, envelope.lens), "high", writer, "$"); } export function warRoomRoundFlow(envelope: TaskEnvelope, config: UltraConfig, cwd: string, round: 1 | 2 | 3, feeds: WarRoomFeeds = {}, bindings: WarRoomBindings = {}): FlowNode { const readOnly = [...SCOPED_READ_TOOLS]; const member = (model: string, author: "lead" | "flow" | "invariants"): AgentNode => { const assigned = fitEnvelope({ ...envelope, trace: { ...envelope.trace, nodeId: `$.branches.${author}` } }, config.context.envelopeHardCapTokens); return { kind: "agent", task: warRoomTaskFor(assigned, author, round, feeds[author], bindings[author]), json: boardBatchSchema(author), model, thinking: "high", tools: [...readOnly, "ultra_signal"], cwd }; }; return { kind: "parallel", mode: "all", onError: "collect", concurrency: Math.min(WAR_ROOM_MEMBERS.length, config.piAgentsBudgets.maxParallelism), branches: { lead: member(config.root.model, "lead"), flow: member(config.scout.model, "flow"), invariants: member(config.scout.model, "invariants"), }, }; } export function warRoomSpecialistFlow(envelope: TaskEnvelope, config: UltraConfig, cwd: string, author: string, gap: string, events: readonly unknown[], binding?: WarRoomWorkerBinding): FlowNode { const assigned = fitEnvelope({ ...envelope, trace: { ...envelope.trace, nodeId: "$.specialist" } }, config.context.envelopeHardCapTokens); const task = `${warRoomTaskFor(assigned, author, 1, JSON.stringify(events), binding)}\nResolve this exact gap: ${gap}`; return { kind: "agent", task, json: boardBatchSchema(author), model: config.scout.model, thinking: "high", tools: [...SCOPED_READ_TOOLS, "ultra_signal"], cwd }; } export interface PlannedAgent { nodeId: string; task: string; model?: string; thinking?: string; tools?: string[]; } export function plannedAgents(flow: FlowNode): PlannedAgent[] { const agents: PlannedAgent[] = []; const visit = (node: FlowNode, nodeId: string): void => { if (node.kind === "agent") { agents.push({ nodeId, task: node.task, ...(node.model ? { model: node.model } : {}), ...(node.thinking ? { thinking: node.thinking } : {}), ...(node.tools ? { tools: node.tools } : {}) }); return; } if (node.kind === "sequence") node.steps.forEach((step, index) => visit(step, `${nodeId}.steps[${index}]`)); if (node.kind === "parallel") for (const [name, branch] of Object.entries(node.branches)) visit(branch, `${nodeId}.branches.${name}`); }; visit(flow, "$"); return agents; } export class PiAgentsHost { private context: ExtensionContext | undefined; private readonly manager: RunManager; private readonly eventDeliveries = new Map>(); constructor(private readonly pi: ExtensionAPI, private readonly onEvent: (event: RunEvent) => void | Promise = () => {}) { this.manager = new RunManager({ engine: workerSpawnEngine(), publish: (event) => { this.enqueueEvent(event); } }); } private enqueueEvent(event: RunEvent): void { const runId = event.type === "run_created" ? event.run.id : event.runId; const prior = this.eventDeliveries.get(runId) ?? Promise.resolve(); const delivery = prior.then(async () => { try { await this.onEvent(event); } finally { this.pi.events.emit(RUN_EVENT_CHANNEL, { protocol: PROTOCOL_VERSION, event }); } }).catch(() => {}); this.eventDeliveries.set(runId, delivery); void delivery.then(() => { if (this.eventDeliveries.get(runId) === delivery) this.eventDeliveries.delete(runId); }); } install(): void { this.pi.events.on(RPC_REQUEST_CHANNEL, (raw: unknown) => { void this.handle(raw); }); this.pi.on("session_start", (_event, context) => { this.context = context; }); this.pi.on("session_shutdown", () => { this.manager.stopAll(); }); } private reply(id: string, success: boolean, data?: unknown, error?: string): void { this.pi.events.emit(`${RPC_REPLY_PREFIX}${id}`, success ? { protocol: PROTOCOL_VERSION, id, success, data } : { protocol: PROTOCOL_VERSION, id, success, error: error ?? "Unknown UltraPi RPC error" }); } private async handle(raw: unknown): Promise { if (!isRecord(raw)) return; const request = raw as RpcRequest; if (request.protocol !== PROTOCOL_VERSION || typeof request.id !== "string" || typeof request.op !== "string") return; try { if (request.op === "ping") this.reply(request.id, true, { protocol: PROTOCOL_VERSION, version: "ultrapi" }); else if (request.op === "list") this.reply(request.id, true, { runs: this.manager.state.order.map((id) => this.manager.state.runs.get(id)).filter(Boolean).map((run) => ({ runId: run!.header.id, label: run!.header.label, status: run!.status, source: run!.header.source })) }); else if (request.op === "stop") { const runId = isRecord(request.params) ? request.params.runId : undefined; if (typeof runId !== "string" || !this.manager.stop(runId)) throw new Error("Run is not live"); this.reply(request.id, true, { runId }); } else if (request.op === "steer") { const params = isRecord(request.params) ? request.params : {}; const runId = params.runId; const message = params.message; const instances = typeof runId === "string" ? this.manager.steerableInstances(runId) : []; const instance = typeof params.instance === "string" ? params.instance : instances.length === 1 ? instances[0] : undefined; if (typeof runId !== "string" || typeof message !== "string" || !instance) throw new Error("A live exact node instance is required for steering"); const result = await this.manager.steer(runId, instance, message, "rpc", typeof request.caller === "string" ? request.caller : undefined); if (result.status !== "queued") throw new Error(result.status === "rejected" ? result.error : result.reason); this.reply(request.id, true, { runId, instance }); } else if (request.op === "start") this.start(request.id, request.params, typeof request.caller === "string" ? request.caller : undefined); else throw new Error(`Unknown RPC operation ${request.op}`); } catch (error) { this.reply(request.id, false, undefined, error instanceof Error ? error.message : String(error)); } } private start(id: string, raw: unknown, caller?: string): void { if (!this.context || !isRecord(raw) || raw.flow === undefined || raw.workflow !== undefined) throw new Error("UltraPi requires one inline flow in an active session"); const cwd = typeof raw.cwd === "string" ? resolve(raw.cwd) : this.context.cwd; const flow = validateFlow(raw.flow); const catalog = buildModelCatalog(this.context.modelRegistry); const started = this.manager.start({ flow, cwd, budgets: isRecord(raw.budgets) ? raw.budgets as UltraConfig["piAgentsBudgets"] : undefined, source: { kind: "rpc", caller }, trusted: this.context.isProjectTrusted(), defaults: { model: this.context.model ? `${this.context.model.provider}/${this.context.model.id}` : undefined, thinking: this.context.thinkingLevel as never }, resolveModel: (ref) => resolveModelReference(ref, catalog) }); this.manager.markBackgrounded(started.runId); started.done.catch(() => {}); this.reply(id, true, { runId: started.runId }); } } const WORKER_SCOPE_EXTENSION = fileURLToPath(new URL("../security/worker-scope-extension.ts", import.meta.url)); export function withWorkerScope(engine: SpawnEngine): SpawnEngine { return { spawn(spec) { const scope = workerScopeFromTask(spec.task); if (!scope) throw new Error("Delegated worker task is missing its scope boundary"); const warRoom = workerWarRoomFromTask(spec.task); return engine.spawn({ ...spec, task: stripWorkerScopeMarker(spec.task), env: { ...spec.env, [WORKER_ROOT_ENV]: spec.cwd, [WORKER_SCOPE_ENV]: JSON.stringify(scope), ...(warRoom ? { [WORKER_WAR_ROOM_ENV]: JSON.stringify(warRoom) } : {}) } }); } }; } function workerSpawnEngine(): SpawnEngine { const engine = createSubprocessSpawnEngine({ extraExtensionPaths: [WORKER_SCOPE_EXTENSION] }); return withWorkerScope(engine); } export class PiAgentsBackend { private readonly client; constructor(pi: ExtensionAPI) { this.client = createPiAgentsClient(pi, { caller: "ultrapi" }); } async start(topology: Exclude, envelope: TaskEnvelope, config: UltraConfig, cwd: string, writerModel?: string, flow?: FlowNode): Promise { const resolvedCwd = resolve(cwd); const params: StartWithBudgets = { flow: flow ?? flowFor(topology, envelope, config, resolvedCwd, writerModel), cwd: resolvedCwd, label: `UltraPi ${topology}`, budgets: config.piAgentsBudgets }; const started = await this.client.start(params); return started.runId; } stop(runId: string) { return this.client.stop(runId); } steer(runId: string, message: string, instance?: string) { return this.client.steer({ runId, message, ...(instance ? { instance } : {}) }); } list() { return this.client.list(); } }