import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { connect } from "node:net"; import { Type } from "typebox"; import { RESULT_TOOL_NAME } from "pi-agents/src/engine/result-tool.js"; import { registerGuardedBash } from "./guarded-bash.js"; import { assertWorkerPath, parseWorkerScope, scopedWriterToolDefinitions, WORKER_ROOT_ENV, WORKER_SCOPE_ENV, WORKER_WAR_ROOM_ENV } from "./worker-scope.js"; import type { WarRoomSignal, WarRoomWorkerBinding } from "../workflows/war-room.js"; import type { ParsedArgv } from "../verification/guarded-exec.js"; const SCOPED_READ_COMMANDS = new Set(["git", "rg", "grep", "ls", "pwd"]); async function assertWorkerCommandScope(root: string, scope: Parameters[2], { file, args }: ParsedArgv): Promise { if (SCOPED_READ_COMMANDS.has(file)) throw new Error("Use the scoped UltraPi read tools for repository inspection"); let paths: string[] = []; if (["npm", "pnpm", "yarn"].includes(file) || file === "bun" && args[0] === "run") { const expected = args[0] === "run" ? 2 : 1; if (args.length !== expected) throw new Error("Delegated package verification does not accept extra arguments"); } else if (file === "node") { const test = args.indexOf("--test"); paths = test < 0 ? [] : args.slice(test + 1); } else if (file === "pytest") { paths = args.filter((arg) => arg !== "-q"); } else if (file === "bun" && args[0] === "test") { paths = args.slice(1); } else if (file === "go" && args[0] === "test") { paths = args.slice(1); } if (paths.some((path) => path.startsWith("-"))) throw new Error("Delegated verification accepts only scoped path arguments"); for (const path of paths) await assertWorkerPath(root, path.endsWith("/...") ? path.slice(0, -4) || "." : path, scope); } export function installResultCorrectionLimit(pi: ExtensionAPI): void { let rejected = 0; pi.on("tool_result", (event, context) => { if (event.toolName === RESULT_TOOL_NAME && event.isError && ++rejected >= 2) context.abort(); }); } async function sendWarRoomSignal(signalPath: string, signal: WarRoomSignal): Promise { const socket = connect(signalPath); return await new Promise((resolve, reject) => { let buffer = ""; let settled = false; const fail = (error: unknown) => { if (settled) return; settled = true; socket.destroy(); reject(error instanceof Error ? error : new Error(String(error))); }; socket.setEncoding("utf8"); socket.setTimeout(5_000, () => fail(new Error("War Room signal controller timed out"))); socket.on("error", fail); socket.on("data", (chunk: string) => { buffer += chunk; if (buffer.length > 2_000) return fail(new Error("War Room signal acknowledgement is oversized")); const newline = buffer.indexOf("\n"); if (newline < 0) return; try { const reply = JSON.parse(buffer.slice(0, newline)) as { ok?: unknown; cursor?: unknown; error?: unknown }; const cursor = reply.cursor; if (reply.ok !== true || typeof cursor !== "number" || !Number.isSafeInteger(cursor) || cursor < 1) throw new Error(typeof reply.error === "string" ? reply.error : "War Room signal was rejected"); settled = true; socket.end(); resolve(cursor); } catch (error) { fail(error); } }); socket.write(`${JSON.stringify(signal)}\n`); }); } function registerWarRoomSignal(pi: ExtensionAPI, binding: WarRoomWorkerBinding): void { pi.registerTool({ name: "ultra_signal", label: "War Room signal", description: "Send one bounded, targeted War Room signal to the controller before submitting your result.", promptSnippet: "ultra_signal: send evidence, a cursor-linked challenge, a concrete request, or a blocker to one War Room member.", parameters: Type.Object({ type: Type.Union([Type.Literal("evidence"), Type.Literal("challenge"), Type.Literal("request"), Type.Literal("blocker")]), target: Type.Optional(Type.String({ minLength: 1, maxLength: 256 })), targetCursor: Type.Optional(Type.Integer({ minimum: 1 })), claim: Type.String({ minLength: 1, maxLength: 1000 }), evidence: Type.Optional(Type.Array(Type.String({ minLength: 1, maxLength: 256 }), { maxItems: 8 })), }), executionMode: "sequential", async execute(id, params) { const signal: WarRoomSignal = { credential: binding.credential, signalId: id, type: params.type, claim: params.claim, ...(params.target ? { target: params.target } : {}), ...(params.targetCursor !== undefined ? { targetCursor: params.targetCursor } : {}), ...(params.evidence ? { evidence: params.evidence } : {}) }; const cursor = await sendWarRoomSignal(binding.signalPath, signal); return { content: [{ type: "text", text: `War Room ${params.type} signal accepted at cursor ${cursor}${params.target ? ` for ${params.target}` : ""}.` }], details: { signalId: id, cursor } }; }, }); } export default function installWorkerScope(pi: ExtensionAPI): void { installResultCorrectionLimit(pi); const root = process.env[WORKER_ROOT_ENV]; const scope = parseWorkerScope(process.env[WORKER_SCOPE_ENV]); if (!root || !scope) return; const warRoom = process.env[WORKER_WAR_ROOM_ENV]; if (warRoom) { try { const parsed = JSON.parse(warRoom) as { credential?: unknown; signalPath?: unknown }; if (typeof parsed.credential === "string" && parsed.credential.length >= 16 && parsed.credential.length <= 128 && typeof parsed.signalPath === "string" && parsed.signalPath.length > 0) registerWarRoomSignal(pi, { credential: parsed.credential, signalPath: parsed.signalPath }); } catch {} } for (const tool of scopedWriterToolDefinitions(root, scope)) pi.registerTool(tool); registerGuardedBash(pi, (argv) => assertWorkerCommandScope(root, scope, argv)); }