import fs from "node:fs/promises"; import type { Model } from "@earendil-works/pi-ai/compat"; import { AuthStorage, createAgentSession, DefaultResourceLoader, ModelRegistry, SessionManager, SettingsManager, } from "@earendil-works/pi-coding-agent"; import { installBrokerFetch } from "./broker-client.js"; import { decodeHarnessFrame, encodeHarnessFrame, HARNESS_ADAPTER_ID, HARNESS_PROFILE_ID, HARNESS_PROTOCOL_VERSION, HarnessRunPacketSchema, type HarnessRunPacket, type HarnessRunResult, MAX_HARNESS_PACKET_BYTES, } from "./protocol.js"; const SESSION_DIR = "/state/sessions"; const protocolWrite = process.stdout.write.bind(process.stdout); async function readStdin(maxBytes: number): Promise { const chunks: Buffer[] = []; let bytes = 0; for await (const raw of process.stdin) { const chunk = Buffer.from(raw as Buffer); bytes += chunk.length; if (bytes > maxBytes) throw new Error("Harness input exceeded its byte limit"); chunks.push(chunk); } return Buffer.concat(chunks); } function writeResult(result: HarnessRunResult): void { protocolWrite(encodeHarnessFrame(result)); } function failed(runId: string, errorCode: Extract["errorCode"]): HarnessRunResult { return { version: HARNESS_PROTOCOL_VERSION, runId, adapter: HARNESS_ADAPTER_ID, profile: HARNESS_PROFILE_ID, status: "failed", errorCode, }; } function buildModel(packet: HarnessRunPacket): Model<"openai-completions"> { return { id: packet.model.id, name: packet.model.id, api: "openai-completions", provider: packet.model.provider, baseUrl: `${packet.broker.virtualOrigin}/v1`, reasoning: packet.model.reasoning, ...(packet.model.reasoning ? { compat: { supportsDeveloperRole: false, supportsReasoningEffort: false, thinkingFormat: "qwen-chat-template" as const, }, } : {}), input: ["text"], cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, contextWindow: packet.model.contextWindow, maxTokens: packet.model.maxOutputTokens, }; } async function sessionManagerFor(packet: HarnessRunPacket): Promise { await fs.mkdir(SESSION_DIR, { recursive: true, mode: 0o700 }); if (packet.session.mode === "new") return SessionManager.create(packet.workspace.path, SESSION_DIR); const sessionId = packet.session.sessionId; const sessions = await SessionManager.list(packet.workspace.path, SESSION_DIR); const selected = sessions.find((session) => session.id === sessionId); if (!selected) throw new Error("session-not-found"); return SessionManager.open(selected.path, SESSION_DIR, packet.workspace.path); } async function main(): Promise { let runId = "invalid"; let restoreFetch: (() => number) | undefined; try { const input = await readStdin(MAX_HARNESS_PACKET_BYTES + 4); const packet = HarnessRunPacketSchema.parse(decodeHarnessFrame(input, MAX_HARNESS_PACKET_BYTES)); runId = packet.runId; restoreFetch = installBrokerFetch(packet); const authStorage = AuthStorage.inMemory({ [packet.model.provider]: { type: "api_key", key: "broker-placeholder" }, }); const modelRegistry = ModelRegistry.inMemory(authStorage); const settingsManager = SettingsManager.inMemory({ defaultProvider: packet.model.provider, defaultModel: packet.model.id, defaultThinkingLevel: "off", defaultProjectTrust: "never", compaction: { enabled: false }, retry: { enabled: false, maxRetries: 0, provider: { maxRetries: 0 } }, packages: [], extensions: [], skills: [], prompts: [], themes: [], enableSkillCommands: false, enableInstallTelemetry: false, enableAnalytics: false, }, { projectTrusted: false }); const resourceLoader = new DefaultResourceLoader({ cwd: packet.workspace.path, agentDir: "/state/agent", settingsManager, noExtensions: true, noSkills: true, noPromptTemplates: true, noThemes: true, noContextFiles: true, systemPrompt: packet.systemPrompt, appendSystemPrompt: [], }); await resourceLoader.reload(); const sessionManager = await sessionManagerFor(packet); const model = buildModel(packet); const { session } = await createAgentSession({ cwd: packet.workspace.path, agentDir: "/state/agent", authStorage, modelRegistry, model, thinkingLevel: "off", scopedModels: [{ model, thinkingLevel: "off" }], tools: packet.tools, resourceLoader, sessionManager, settingsManager, }); const before = session.getSessionStats().tokens; const toolReceipts: Array<{ sequence: number; tool: typeof packet.tools[number]; status: "completed" | "failed" }> = []; let receiptOverflow = false; const unsubscribe = session.subscribe((event) => { if (event.type !== "tool_execution_end") return; if (toolReceipts.length >= packet.limits.maxToolReceipts) { receiptOverflow = true; session.abort(); return; } if (!packet.tools.includes(event.toolName as typeof packet.tools[number])) { receiptOverflow = true; session.abort(); return; } toolReceipts.push({ sequence: toolReceipts.length + 1, tool: event.toolName as typeof packet.tools[number], status: event.isError ? "failed" : "completed", }); }); try { await session.prompt(packet.prompt, { expandPromptTemplates: false, source: "rpc" }); } finally { unsubscribe(); } if (receiptOverflow) throw new Error("tool-receipt-overflow"); const finalText = session.getLastAssistantText(); if (typeof finalText !== "string") throw new Error("missing-final-output"); if (Buffer.byteLength(finalText) > packet.limits.maxResultBytes) { writeResult(failed(runId, "output-oversize")); session.dispose(); return; } const after = session.getSessionStats().tokens; const usage = { inputTokens: Math.max(0, after.input - before.input), outputTokens: Math.max(0, after.output - before.output), cacheReadTokens: Math.max(0, after.cacheRead - before.cacheRead), cacheWriteTokens: Math.max(0, after.cacheWrite - before.cacheWrite), totalTokens: Math.max(0, after.total - before.total), }; const result: HarnessRunResult = { version: HARNESS_PROTOCOL_VERSION, runId, adapter: HARNESS_ADAPTER_ID, profile: HARNESS_PROFILE_ID, status: "completed", sessionId: session.sessionId, finalText, toolReceipts, usage, }; session.dispose(); if (encodeHarnessFrame(result).length - 4 > packet.limits.maxResultBytes) { writeResult(failed(runId, "output-oversize")); return; } writeResult(result); } catch (error) { const errorCode = error instanceof Error && error.message === "session-not-found" ? "session-not-found" : runId === "invalid" ? "invalid-packet" : error instanceof Error && /provider|model/i.test(error.message) ? "provider-failure" : "adapter-failure"; writeResult(failed(runId, errorCode)); } finally { restoreFetch?.(); } } void main();