import { spawn, type ChildProcess } from "node:child_process"; import { createInterface, type Interface } from "node:readline"; import type { AcpInitializeResult, AcpPromptResult, AcpSessionResult, AcpSessionUpdate, JsonRpcId, JsonRpcMessage, ReasoningEffort, } from "./types.js"; export type ReverseRequestHandler = (method: string, params: unknown, signal: AbortSignal) => Promise; export type NotificationHandler = (message: JsonRpcMessage) => void; interface PendingRequest { resolve(value: unknown): void; reject(error: Error): void; } export class AcpJsonRpcClient { private nextId = 1; private readonly pending = new Map(); private readonly notificationHandlers = new Set(); private readonly reverseControllers = new Set(); private readonly lines: Interface; private reverseRequestHandler: ReverseRequestHandler | undefined; private closed = false; private stderr = ""; constructor(private readonly process: ChildProcess) { if (!process.stdout || !process.stdin) throw new Error("ACP process must have piped stdin and stdout"); this.lines = createInterface({ input: process.stdout, crlfDelay: Infinity, terminal: false }); this.lines.on("line", (line) => this.handleLine(line)); process.stderr?.setEncoding("utf8"); process.stderr?.on("data", (chunk: string) => { this.stderr = `${this.stderr}${chunk}`.slice(-32_000); }); process.on("error", (error) => this.rejectAll(error)); process.on("close", (code) => { this.closed = true; this.cancelReverseRequests(); if (this.pending.size === 0) return; const detail = this.stderr.trim(); this.rejectAll( new Error( detail ? `Grok ACP process exited (code ${code}): ${detail}` : `Grok ACP process exited before responding (code ${code})`, ), ); }); } setReverseRequestHandler(handler: ReverseRequestHandler): void { this.reverseRequestHandler = handler; } onNotification(handler: NotificationHandler): () => void { this.notificationHandlers.add(handler); return () => this.notificationHandlers.delete(handler); } request(method: string, params: Record, signal?: AbortSignal): Promise { if (this.closed) return Promise.reject(new Error("Grok ACP client is closed")); if (signal?.aborted) return Promise.reject(new Error("aborted")); const id = this.nextId++; const promise = new Promise((resolve, reject) => { const onAbort = () => { this.pending.delete(id); reject(new Error("aborted")); }; signal?.addEventListener("abort", onAbort, { once: true }); this.pending.set(id, { resolve: (value) => { signal?.removeEventListener("abort", onAbort); resolve(value); }, reject: (error) => { signal?.removeEventListener("abort", onAbort); reject(error); }, }); }); try { this.write({ jsonrpc: "2.0", id, method, params }); } catch (error) { this.pending.delete(id); return Promise.reject(error); } return promise; } notify(method: string, params: Record): void { if (this.closed) return; this.write({ jsonrpc: "2.0", method, params }); } cancelReverseRequests(): void { for (const controller of this.reverseControllers) controller.abort(); } getStderr(): string { return this.stderr; } dispose(): void { if (this.closed) return; this.closed = true; this.cancelReverseRequests(); this.rejectAll(new Error("Grok ACP client disposed")); this.lines.close(); if (this.process.exitCode === null && this.process.signalCode === null) { this.process.kill("SIGTERM"); setTimeout(() => { if (this.process.exitCode === null && this.process.signalCode === null) this.process.kill("SIGKILL"); }, 1_500).unref(); } } private write(message: JsonRpcMessage): void { if (!this.process.stdin?.writable) throw new Error("Grok ACP stdin is not writable"); this.process.stdin.write(`${JSON.stringify(message)}\n`); } private handleLine(line: string): void { let message: JsonRpcMessage; try { message = JSON.parse(line) as JsonRpcMessage; } catch { return; } if (message.method && message.id !== undefined && message.id !== null) { void this.handleReverseRequest(message.id, message.method, message.params); return; } if (message.id !== undefined && message.id !== null) { const pending = this.pending.get(message.id); if (!pending) return; this.pending.delete(message.id); if (message.error) { pending.reject(new Error(message.error.message ?? `ACP request failed: ${message.error.code ?? "unknown"}`)); } else { pending.resolve(message.result); } return; } for (const handler of this.notificationHandlers) { try { handler(message); } catch { // Notification observers cannot break the transport. } } } private async handleReverseRequest(id: JsonRpcId, method: string, params: unknown): Promise { if (!this.reverseRequestHandler) { this.write({ jsonrpc: "2.0", id, error: { code: -32601, message: `Unsupported ACP client method: ${method}` } }); return; } const controller = new AbortController(); this.reverseControllers.add(controller); try { const result = await this.reverseRequestHandler(method, params, controller.signal); if (!this.closed) this.write({ jsonrpc: "2.0", id, result }); } catch (error) { if (!this.closed) { this.write({ jsonrpc: "2.0", id, error: { code: -32603, message: error instanceof Error ? error.message : "ACP client request failed", }, }); } } finally { this.reverseControllers.delete(controller); } } private rejectAll(error: Error): void { for (const pending of this.pending.values()) pending.reject(error); this.pending.clear(); } } export function buildGrokAcpArgs(modelId: string, reasoningEffort?: ReasoningEffort): string[] { const args = ["--permission-mode", "default", "agent", "--no-leader", "--model", modelId]; if (reasoningEffort && reasoningEffort !== "none") args.push("--reasoning-effort", reasoningEffort); args.push("stdio"); return args; } function chooseAuthMethod(result: AcpInitializeResult): string { const methods = new Set(result.authMethods?.map((method) => method.id).filter((id): id is string => Boolean(id)) ?? []); if (methods.has("cached_token")) return "cached_token"; if (process.env.XAI_API_KEY && methods.has("xai.api_key")) return "xai.api_key"; throw new Error(`Grok Build is not logged in. Run \`grok login\`. Offered methods: ${[...methods].join(", ") || "none"}`); } export interface LiveAcpSession { client: AcpJsonRpcClient; process: ChildProcess; sessionId: string; dispose(): void; } export async function createLiveAcpSession(options: { binary: string; modelId: string; reasoningEffort?: ReasoningEffort; cwd: string; reverseRequestHandler: ReverseRequestHandler; signal?: AbortSignal; }): Promise { const process = spawn(options.binary, buildGrokAcpArgs(options.modelId, options.reasoningEffort), { cwd: options.cwd, env: globalThis.process.env, stdio: ["pipe", "pipe", "pipe"], }); const client = new AcpJsonRpcClient(process); client.setReverseRequestHandler(options.reverseRequestHandler); try { const initialized = (await client.request( "initialize", { protocolVersion: 1, clientCapabilities: { fs: { readTextFile: false, writeTextFile: false }, terminal: false, }, clientInfo: { name: "pi-grok-build-acp", version: "0.1.0" }, }, options.signal, )) as AcpInitializeResult; await client.request("authenticate", { methodId: chooseAuthMethod(initialized) }, options.signal); const session = (await client.request( "session/new", { cwd: options.cwd, mcpServers: [] }, options.signal, )) as AcpSessionResult; if (!session.sessionId) throw new Error("Grok ACP did not return a session id"); return { client, process, sessionId: session.sessionId, dispose: () => client.dispose() }; } catch (error) { client.dispose(); throw error; } } export async function promptAcpSession( session: LiveAcpSession, prompt: string, options: { signal?: AbortSignal; onUpdate(update: AcpSessionUpdate): void }, ): Promise { const unsubscribe = session.client.onNotification((message) => { if (message.method !== "session/update" || typeof message.params !== "object" || message.params === null) return; const update = (message.params as { update?: AcpSessionUpdate }).update; if (update) options.onUpdate(update); }); const onAbort = () => { session.client.notify("session/cancel", { sessionId: session.sessionId }); session.client.cancelReverseRequests(); }; options.signal?.addEventListener("abort", onAbort, { once: true }); try { return (await session.client.request( "session/prompt", { sessionId: session.sessionId, prompt: [{ type: "text", text: prompt }], }, options.signal, )) as AcpPromptResult; } finally { options.signal?.removeEventListener("abort", onAbort); unsubscribe(); } }