import { spawn } from "node:child_process"; import { constants } from "node:fs"; import { access } from "node:fs/promises"; import { join } from "node:path"; import type { BashOperations } from "@earendil-works/pi-coding-agent"; import type { Client } from "@modelcontextprotocol/sdk/client/index.js"; import type { StdioClientTransport } from "@modelcontextprotocol/sdk/client/stdio.js"; import { dbg, span } from "./src/debug.ts"; import { findExecutable } from "./src/pi-ext-executable.ts"; const SAFE_NU_CODES = new Set([ "nu::shell::error", "nu::shell::external_command", "nu::shell::division_by_zero", ]); const FRESH_ARGS = [ "--no-config-file", "--no-history", "--error-style", "fancy", "--table-mode", "markdown", "-c", ]; export function findNu(): string | undefined { return findExecutable(process.platform === "win32" ? ["nu.exe"] : ["nu"]); } function stopProcess(pid: number | undefined): void { if (!pid) return; if (process.platform === "win32") { const root = process.env.SystemRoot ?? "C:\\Windows"; spawn( join(root, "System32", "taskkill.exe"), ["/F", "/T", "/PID", String(pid)], { stdio: "ignore", windowsHide: true, }, ).unref(); } else { try { process.kill(-pid, "SIGKILL"); } catch { try { process.kill(pid, "SIGKILL"); } catch { /* already exited */ } } } } export function freshOperations(binary: string): BashOperations { return { async exec(command, cwd, { onData, signal, timeout, env }) { const finish = span?.( "fresh", timeout === undefined ? undefined : { timeoutMs: timeout * 1_000 }, ); let exitCode: number | null = null; let timedOut = false; let outputBytes = 0; try { if (signal?.aborted) throw new Error("aborted"); if ( timeout !== undefined && (!Number.isFinite(timeout) || timeout <= 0 || timeout * 1000 > 2_147_483_647) ) throw new Error( "Invalid timeout: must be a positive number of seconds within the timer range", ); await access(cwd, constants.F_OK); const child = spawn(binary, [...FRESH_ARGS, command], { cwd, env, detached: process.platform !== "win32", stdio: ["ignore", "pipe", "pipe"], windowsHide: true, }); const onAbort = () => stopProcess(child.pid); const timer = timeout === undefined ? undefined : setTimeout(() => { timedOut = true; onAbort(); }, timeout * 1000); signal?.addEventListener("abort", onAbort, { once: true }); if (signal?.aborted) onAbort(); if (finish) { const forward = (chunk: Buffer) => { outputBytes += chunk.length; onData(chunk); }; child.stdout.on("data", forward); child.stderr.on("data", forward); } else { child.stdout.on("data", onData); child.stderr.on("data", onData); } try { exitCode = await new Promise((resolve, reject) => { child.once("error", reject); child.once("exit", resolve); }); if (signal?.aborted) throw new Error("aborted"); if (timedOut) throw new Error(`timeout:${timeout}`); return { exitCode: exitCode ?? (child.signalCode ? 1 : null) }; } finally { if (timer) clearTimeout(timer); signal?.removeEventListener("abort", onAbort); } } finally { const status = signal?.aborted ? "aborted" : timedOut ? "timeout" : exitCode === 0 ? "success" : exitCode === null ? "failed" : "nonzero"; finish?.( status === "success" || status === "nonzero" ? "finish" : "error", { status, count: outputBytes, }, ); } }, }; } type Worker = { client: Client; transport: StdioClientTransport; types: typeof import("@modelcontextprotocol/sdk/types.js"); }; function errorKind( error: unknown, worker: Worker | undefined, signal: AbortSignal | undefined, ): "aborted" | "mcp" | "nu" | "failure" { if (signal?.aborted) return "aborted"; if (worker && error instanceof worker.types.McpError) return "mcp"; if ( error instanceof Error && SAFE_NU_CODES.has(/\bnu::[a-z_]+::[a-z_]+\b/.exec(error.message)?.[0] ?? "") ) return "nu"; return "failure"; } export class NuSession { private worker: Worker | undefined; private pending: Promise = Promise.resolve(); private generation = 0; private async start(binary: string, cwd: string): Promise { const finish = span?.("session.connect"); const { Client } = await import( "@modelcontextprotocol/sdk/client/index.js" ); const { StdioClientTransport } = await import( "@modelcontextprotocol/sdk/client/stdio.js" ); // Load MCP schemas with the first worker, not during inert Pi startup. const types = await import("@modelcontextprotocol/sdk/types.js"); const client = new Client({ name: "pi-ext-nushell", version: "1" }); const transport = new StdioClientTransport({ command: binary, args: [ "--no-config-file", "--no-history", "--mcp", "--mcp-transport", "stdio", ], cwd, env: Object.fromEntries( Object.entries(process.env).filter( (entry): entry is [string, string] => entry[1] !== undefined, ), ), stderr: "pipe", }); // Nu MCP logs requests and responses to stderr, including command text. // Drain it; never print it into Pi's TUI or JSON mode. transport.stderr?.on("data", () => {}); try { await client.connect(transport); finish?.(); } catch (error) { finish?.("error", { kind: errorKind(error, undefined, undefined) }); await client.close(); throw error; } return { client, transport, types }; } async evaluate( binary: string, cwd: string, command: string, signal?: AbortSignal, ): Promise { const work = async () => { const finish = span?.("session.request"); let worker: Worker | undefined; let startedCall = false; try { if (signal?.aborted) throw new Error("Nu session request aborted before execution"); const generation = this.generation; worker = this.worker ?? (await this.start(binary, cwd)); if (generation !== this.generation || signal?.aborted) { if (worker !== this.worker) await worker.client.close(); throw new Error("Nu session reset before execution"); } this.worker = worker; startedCall = true; const { CallToolResultSchema } = worker.types; const result = CallToolResultSchema.parse( await worker.client.callTool( { name: "evaluate", arguments: { input: command } }, CallToolResultSchema, { signal, timeout: 110_000 }, ), ); const text = result.content .filter((part) => part.type === "text") .map((part) => part.text) .join("\n"); if (result.isError) throw new Error(text || "Nushell evaluation failed"); finish?.("finish", { count: Buffer.byteLength(text) }); return text || "(no output)"; } catch (error) { finish?.("error", { kind: errorKind(error, worker, signal) }); if (startedCall && worker) { const { McpError, ErrorCode } = worker.types; if ( signal?.aborted || (error instanceof McpError && (error.code === ErrorCode.RequestTimeout || (error.code === ErrorCode.InternalError && error.message.includes( "Operation promoted to background job (id:", )))) ) { await this.reset("interrupted"); throw new Error( "Nu session interrupted; state reset. Detached external commands may still be running; inspect before retrying.", ); } if (worker.transport.pid === null) await this.reset("worker_exit"); } throw error; } }; const next = this.pending.then(work, work); this.pending = next.catch(() => {}); return next; } async reset( reason: | "requested" | "interrupted" | "worker_exit" | "explicit" | "toggle_off" | "session_start" | "session_tree" | "session_shutdown" = "requested", ): Promise { this.generation++; const worker = this.worker; this.worker = undefined; dbg?.("session.reset", { kind: reason }); if (worker) await worker.client.close(); } }