From 4422568f61020bd762a7978cbc79cce771937674 Mon Sep 17 00:00:00 2001 From: mptyl Date: Wed, 5 Aug 2026 05:23:58 +0200 Subject: [PATCH] fix: bind pi runtime config at spawn --- backend/src/pi/managed-config.ts | 90 ++++++++++++++-- backend/src/pi/pi-process-manager.ts | 88 ++++++++------- backend/test/pi-process-manager.test.ts | 135 +++++++++++++++++++++++- backend/test/routes-sessions.test.ts | 108 +++++++++++++++++++ 4 files changed, 375 insertions(+), 46 deletions(-) diff --git a/backend/src/pi/managed-config.ts b/backend/src/pi/managed-config.ts index 2d1cbfde..42c12603 100644 --- a/backend/src/pi/managed-config.ts +++ b/backend/src/pi/managed-config.ts @@ -1,8 +1,9 @@ import { - closeSync, constants, fstatSync, lstatSync, openSync, readFileSync, + chmodSync, closeSync, constants, fstatSync, lstatSync, mkdtempSync, openSync, + readFileSync, readdirSync, rmSync, symlinkSync, writeFileSync, type Dirent, } from "node:fs"; -import { homedir } from "node:os"; -import { join } from "node:path"; +import { homedir, tmpdir } from "node:os"; +import { join, resolve } from "node:path"; const MAX_AGENT_CONFIG_BYTES = 1024 * 1024; @@ -55,13 +56,15 @@ export function validateDeclarativePiConfig(raw: string): void { assertDeclarativePiConfig(parsePiConfigJson(raw)); } -export function readConfiguredPiAgentFile(name: "auth.json"): string; -export function readConfiguredPiAgentFile(name: "models.json", optional: true): string | undefined; -export function readConfiguredPiAgentFile( +function configuredPiAgentDir(): string { + return resolve(process.env.PI_CODING_AGENT_DIR ?? join(homedir(), ".pi", "agent")); +} + +function readPiAgentFile( + configuredAgentDir: string, name: "auth.json" | "models.json", - optional = false, + optional: boolean, ): string | undefined { - const configuredAgentDir = process.env.PI_CODING_AGENT_DIR ?? join(homedir(), ".pi", "agent"); const path = join(configuredAgentDir, name); let fd: number | undefined; try { @@ -85,3 +88,74 @@ export function readConfiguredPiAgentFile( } } } + +export function readConfiguredPiAgentFile(name: "auth.json"): string; +export function readConfiguredPiAgentFile(name: "auth.json", optional: true): string | undefined; +export function readConfiguredPiAgentFile(name: "models.json", optional: true): string | undefined; +export function readConfiguredPiAgentFile( + name: "auth.json" | "models.json", + optional = false, +): string | undefined { + return readPiAgentFile(configuredPiAgentDir(), name, optional); +} + +export interface PiRuntimeAgentSnapshot { + agentDir: string; + sessionDir: string; + cleanup: () => void; +} + +/** + * Bind a session Pi process to the exact managed auth/model bytes validated at spawn time. + * Other agent resources remain live through symlinks, while session storage stays persistent. + */ +export function createPiRuntimeAgentSnapshot(): PiRuntimeAgentSnapshot { + const sourceAgentDir = configuredPiAgentDir(); + const auth = readPiAgentFile(sourceAgentDir, "auth.json", true); + const models = readPiAgentFile(sourceAgentDir, "models.json", true); + if (auth !== undefined) validateDeclarativePiConfig(auth); + if (models !== undefined) validateDeclarativePiConfig(models); + + let snapshotDir: string | undefined; + try { + snapshotDir = mkdtempSync(join(tmpdir(), "thoth-pi-runtime-agent-")); + chmodSync(snapshotDir, 0o700); + let entries: Dirent[]; + try { + entries = readdirSync(sourceAgentDir, { withFileTypes: true }); + } catch (error) { + if ((error as NodeJS.ErrnoException)?.code !== "ENOENT") throw error; + entries = []; + } + for (const entry of entries) { + if (entry.name === "auth.json" || entry.name === "models.json") continue; + symlinkSync( + join(sourceAgentDir, entry.name), + join(snapshotDir, entry.name), + entry.isDirectory() ? (process.platform === "win32" ? "junction" : "dir") : "file", + ); + } + if (auth !== undefined) { + writeFileSync(join(snapshotDir, "auth.json"), auth, { flag: "wx", mode: 0o600 }); + } + if (models !== undefined) { + writeFileSync(join(snapshotDir, "models.json"), models, { flag: "wx", mode: 0o600 }); + } + } catch { + if (snapshotDir !== undefined) { + try { rmSync(snapshotDir, { recursive: true, force: true }); } catch { /* sanitized */ } + } + throw new PiManagedConfigError(); + } + + let cleaned = false; + return { + agentDir: snapshotDir, + sessionDir: process.env.PI_CODING_AGENT_SESSION_DIR || join(sourceAgentDir, "sessions"), + cleanup: () => { + if (cleaned) return; + cleaned = true; + try { rmSync(snapshotDir, { recursive: true, force: true }); } catch { /* sanitized */ } + }, + }; +} diff --git a/backend/src/pi/pi-process-manager.ts b/backend/src/pi/pi-process-manager.ts index a12750e4..e4b2486b 100644 --- a/backend/src/pi/pi-process-manager.ts +++ b/backend/src/pi/pi-process-manager.ts @@ -7,6 +7,7 @@ import { buildPiChildEnv, canonicalPiProvider } from "./provider-credentials.js" import { loadPiAuthProviders } from "./auth-providers.js"; import { secretValue } from "../config/secret-bundle.js"; import { clearPrincipalEnvironment, principalEnvironment, type PrincipalContext } from "../auth/principal.js"; +import { createPiRuntimeAgentSnapshot } from "./managed-config.js"; export interface SessionRuntime { rpc: RpcClient; @@ -37,13 +38,14 @@ export class PiProcessManager { private spawnFn: ( sessionId: string, author: string, provider: string | undefined, principal?: PrincipalContext, ) => ChildProcessWithoutNullStreams; - private loadAuthProviders: () => ReadonlySet; + private loadAuthProviders: (agentDir: string) => ReadonlySet; constructor( private cfg: AppConfig, - opts?: { spawnFn?: SpawnFn; authProviders?: () => ReadonlySet }, + opts?: { spawnFn?: SpawnFn; authProviders?: (agentDir: string) => ReadonlySet }, ) { - this.loadAuthProviders = opts?.authProviders ?? (() => loadPiAuthProviders()); + this.loadAuthProviders = opts?.authProviders + ?? ((agentDir) => loadPiAuthProviders({ agentDir })); if (opts?.spawnFn) { this.spawnFn = (sessionId, author, provider, principal) => this.spawnPi(opts.spawnFn!, sessionId, author, provider, principal); @@ -56,45 +58,57 @@ export class PiProcessManager { private spawnPi( spawnFn: SpawnFn, sessionId: string, author: string, provider: string | undefined, principal?: PrincipalContext, ): ChildProcessWithoutNullStreams { - const env = buildPiChildEnv({ - provider, - authProviders: this.loadAuthProviders(), - credentialValue: secretValue(this.cfg, "THT_MODEL_API_KEY"), - credentialFile: this.cfg.modelApiKeyFile, - additions: { THT_SESSION: sessionId, THT_AUTHOR: author }, - }); - clearPrincipalEnvironment(env); - if (principal) Object.assign(env, principalEnvironment(principal)); - // The Thoth gate executes the deterministic `tht` CLI as a Pi tool. Give only - // this managed session process the adapter values already loaded by the core - // entrypoint; the generic provider helper continues to scrub them by default. - for (const name of [ - "THT_DWH_API_KEY", "THT_VEC_API_KEY", "THT_VEC_WRITE_API_KEY", - ] as const) { - const value = secretValue(this.cfg, name) ?? process.env[name]; - if (value !== undefined) env[name] = value; - } - const ca = secretValue(this.cfg, "THT_SSL_CA") - ?? secretValue(this.cfg, "THT_CA") - ?? process.env.THT_SSL_CA - ?? process.env.THT_CA; - if (ca !== undefined) { - env.THT_CA = ca; - env.THT_SSL_CA = ca; - } - delete env.THT_DATA_ROOT; - if (this.cfg.dataRoot !== undefined) env.THT_DATA_ROOT = this.cfg.dataRoot; - // pi 0.73 removed `--approve`: rpc mode is headless and its argv is intentionally minimal. - const child = spawnFn(this.cfg.piBin, ["--mode", "rpc"], { - cwd: this.cfg.harnessDir, - env, - }); + // This is the final shared boundary for createFor(), spawnFor(), and resume(). Validate + // before auth-provider inspection, then make Pi consume the exact copied bytes rather than + // reopening mutable mounted auth/models files after this check. + const agent = createPiRuntimeAgentSnapshot(); + let child: ChildProcessWithoutNullStreams | undefined; try { + const env = buildPiChildEnv({ + provider, + authProviders: this.loadAuthProviders(agent.agentDir), + credentialValue: secretValue(this.cfg, "THT_MODEL_API_KEY"), + credentialFile: this.cfg.modelApiKeyFile, + additions: { THT_SESSION: sessionId, THT_AUTHOR: author }, + }); + env.PI_CODING_AGENT_DIR = agent.agentDir; + env.PI_CODING_AGENT_SESSION_DIR = agent.sessionDir; + clearPrincipalEnvironment(env); + if (principal) Object.assign(env, principalEnvironment(principal)); + // The Thoth gate executes the deterministic `tht` CLI as a Pi tool. Give only + // this managed session process the adapter values already loaded by the core + // entrypoint; the generic provider helper continues to scrub them by default. + for (const name of [ + "THT_DWH_API_KEY", "THT_VEC_API_KEY", "THT_VEC_WRITE_API_KEY", + ] as const) { + const value = secretValue(this.cfg, name) ?? process.env[name]; + if (value !== undefined) env[name] = value; + } + const ca = secretValue(this.cfg, "THT_SSL_CA") + ?? secretValue(this.cfg, "THT_CA") + ?? process.env.THT_SSL_CA + ?? process.env.THT_CA; + if (ca !== undefined) { + env.THT_CA = ca; + env.THT_SSL_CA = ca; + } + delete env.THT_DATA_ROOT; + if (this.cfg.dataRoot !== undefined) env.THT_DATA_ROOT = this.cfg.dataRoot; + // pi 0.73 removed `--approve`: rpc mode is headless and its argv is intentionally minimal. + child = spawnFn(this.cfg.piBin, ["--mode", "rpc"], { + cwd: this.cfg.harnessDir, + env, + }); + child.once("exit", agent.cleanup); + child.once("close", agent.cleanup); // Log stderr for debugging (was silently drained) child.stderr.on("data", (d: Buffer) => console.error(`[pi:${sessionId}] stderr:`, d.toString().trim())); return child; } catch (error) { - try { child.kill(); } catch { /* preserve the initialization error */ } + if (child) { + try { child.kill(); } catch { /* preserve the initialization error */ } + } + agent.cleanup(); throw error; } } diff --git a/backend/test/pi-process-manager.test.ts b/backend/test/pi-process-manager.test.ts index 27c4b71a..51d1ce96 100644 --- a/backend/test/pi-process-manager.test.ts +++ b/backend/test/pi-process-manager.test.ts @@ -3,14 +3,36 @@ import { spawn } from "node:child_process"; import path from "node:path"; import { fileURLToPath } from "node:url"; import { EventEmitter } from "node:events"; -import { chmodSync, writeFileSync } from "node:fs"; +import { + chmodSync, existsSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; import { PiProcessManager } from "../src/pi/pi-process-manager.js"; import { loadConfig } from "../src/config.js"; +import { + PI_MANAGED_CONFIG_ERROR_MESSAGE, + validateDeclarativePiConfig, +} from "../src/pi/managed-config.js"; const __dirname = path.dirname(fileURLToPath(import.meta.url)); const FAKE = path.resolve(__dirname, "../../harness/tests/fake_pi/fake_pi_rpc.mjs"); const SCRIPT = path.resolve(__dirname, "../../harness/tests/fake_pi/scripts/f1_disambiguation.json"); +const SAFE_AUTH = '{"deepseek":{"type":"api_key","key":"safe-token"}}\n'; +const SAFE_MODELS = [ + "{", + ' "providers": {', + ' "local-qwen": {"baseUrl":"http://model.invalid/v1","models":[{"id":"qwen"}]}', + " }", + "}", + "", +].join("\n"); + +function writeSafeAgentConfig(agentDir: string): void { + writeFileSync(path.join(agentDir, "auth.json"), SAFE_AUTH, { mode: 0o600 }); + writeFileSync(path.join(agentDir, "models.json"), SAFE_MODELS, { mode: 0o600 }); +} + test("spawnFor avvia un runtime e il bridge emette il widget F1", async () => { const cfg = loadConfig({ THT_HARNESS_DIR: "../harness" }); const mgr = new PiProcessManager(cfg, { spawnFn: () => spawn("node", [FAKE, SCRIPT]) as any }); @@ -146,6 +168,117 @@ function recordingChild() { return ch; } +test.each([ + ["new", "auth.json", '{"deepseek":{"key":"!runtime-auth-command runtime-secret /private/runtime-auth"}}\n'], + ["new", "models.json", '{"providers":{"local-qwen":{"headers":["!runtime-model-command runtime-secret /private/runtime-model"]}}}\n'], + ["resume", "auth.json", '{"deepseek":{"key":"!resume-auth-command runtime-secret /private/resume-auth"}}\n'], + ["resume", "models.json", '{"providers":{"local-qwen":{"models":[{"apiKey":"!resume-model-command runtime-secret /private/resume-model"}]}}}\n'], +] as const)( + "%s runtime rejects post-admission executable %s before auth resolution or child spawn", + async (mode, changedFile, unsafeRaw) => { + const root = mkdtempSync(path.join(tmpdir(), "tht-runtime-managed-config-")); + const agentDir = path.join(root, "agent"); + mkdirSync(agentDir, { mode: 0o700 }); + writeSafeAgentConfig(agentDir); + vi.stubEnv("PI_CODING_AGENT_DIR", agentDir); + let authResolutions = 0; + let spawns = 0; + const mgr = new PiProcessManager(loadConfig({}), { + authProviders: () => { authResolutions += 1; return new Set(); }, + spawnFn: () => { spawns += 1; throw new Error("SPAWN_BOUNDARY_REACHED"); }, + }); + + try { + // Admission/model validation succeeded while the mounted files were still safe, and the + // session was then persisted. The operator-controlled mount changes before runtime start. + validateDeclarativePiConfig(readFileSync(path.join(agentDir, "auth.json"), "utf8")); + validateDeclarativePiConfig(readFileSync(path.join(agentDir, "models.json"), "utf8")); + writeFileSync(path.join(root, "session-created"), `${mode}\n`); + writeFileSync(path.join(agentDir, changedFile), unsafeRaw, { mode: 0o600 }); + + let failure: unknown; + try { + if (mode === "new") { + mgr.createFor("post-admission-new", { provider: "local-qwen" }); + } else { + await mgr.spawnFor("post-admission-resume", { + provider: "local-qwen", mode: "resume", + }); + } + } catch (error) { + failure = error; + } + const message = failure instanceof Error ? failure.message : String(failure); + + expect({ message, authResolutions, spawns, runtimes: mgr.count() }).toEqual({ + message: PI_MANAGED_CONFIG_ERROR_MESSAGE, + authResolutions: 0, + spawns: 0, + runtimes: 0, + }); + expect(message).not.toMatch(/runtime-secret|\/private\/|runtime-(?:auth|model)-command|resume-(?:auth|model)-command/); + } finally { + mgr.teardown("post-admission-new"); + mgr.teardown("post-admission-resume"); + vi.unstubAllEnvs(); + rmSync(root, { recursive: true, force: true }); + } + }, +); + +test("runtime Pi consumes exact validated auth/models snapshots and keeps persistent agent resources", async () => { + const root = mkdtempSync(path.join(tmpdir(), "tht-runtime-agent-snapshot-")); + const agentDir = path.join(root, "agent"); + const sessionsDir = path.join(agentDir, "sessions"); + const extensionDir = path.join(agentDir, "extensions"); + mkdirSync(sessionsDir, { recursive: true, mode: 0o700 }); + mkdirSync(extensionDir, { recursive: true, mode: 0o700 }); + writeSafeAgentConfig(agentDir); + const settings = '{"quietStartup":true}\n'; + writeFileSync(path.join(agentDir, "settings.json"), settings, { mode: 0o600 }); + writeFileSync(path.join(extensionDir, "runtime-extension.js"), "export default {};\n"); + vi.stubEnv("PI_CODING_AGENT_DIR", agentDir); + vi.stubEnv("PI_CODING_AGENT_SESSION_DIR", ""); + const child = recordingChild(); + let spawnEnv: NodeJS.ProcessEnv | undefined; + const mgr = new PiProcessManager(loadConfig({}), { + spawnFn: (_command, _args, options) => { + spawnEnv = options.env; + // This mutation happens after validation but before the child can open either source file. + writeFileSync(path.join(agentDir, "auth.json"), '{"deepseek":{"key":"!late-auth-command"}}\n'); + writeFileSync(path.join(agentDir, "models.json"), '{"providers":{"late":{"apiKey":"!late-model-command"}}}\n'); + return child as any; + }, + }); + let snapshotDir: string | undefined; + let childExited = false; + + try { + await mgr.spawnFor("snapshot-session", { provider: "local-qwen" }); + snapshotDir = spawnEnv?.PI_CODING_AGENT_DIR; + + expect(snapshotDir).toBeTruthy(); + expect(snapshotDir).not.toBe(agentDir); + expect(readFileSync(path.join(snapshotDir!, "auth.json"), "utf8")).toBe(SAFE_AUTH); + expect(readFileSync(path.join(snapshotDir!, "models.json"), "utf8")).toBe(SAFE_MODELS); + expect(readFileSync(path.join(snapshotDir!, "settings.json"), "utf8")).toBe(settings); + expect(readFileSync(path.join(snapshotDir!, "extensions", "runtime-extension.js"), "utf8")) + .toBe("export default {};\n"); + expect(spawnEnv?.PI_CODING_AGENT_SESSION_DIR).toBe(sessionsDir); + + mgr.teardown("snapshot-session"); + child.emit("exit", 0); + childExited = true; + expect(existsSync(snapshotDir!)).toBe(false); + } finally { + mgr.teardown("snapshot-session"); + if (!childExited) child.emit("exit", 0); + vi.unstubAllEnvs(); + if (snapshotDir && snapshotDir !== agentDir) rmSync(snapshotDir, { recursive: true, force: true }); + rmSync(root, { recursive: true, force: true }); + } +}); + test("createFor kills a spawned child when post-spawn initialization throws", () => { const child = recordingChild(); child.kill = vi.fn(); diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 8b8e2e8b..4d05cae0 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -8,6 +8,8 @@ import { buildApp as buildRealApp } from "../src/app.js"; import { loadConfig } from "../src/config.js"; import { SseHub } from "../src/sse/sse-hub.js"; import { MaintenanceBarrier } from "../src/runtime/maintenance-gate.js"; +import { PiProcessManager } from "../src/pi/pi-process-manager.js"; +import { validateDeclarativePiConfig } from "../src/pi/managed-config.js"; const FAKE = path.resolve("../harness/tests/fake_pi/fake_pi_rpc.mjs"); const SCRIPT = path.resolve("../harness/tests/fake_pi/scripts/f1_disambiguation.json"); @@ -2426,6 +2428,112 @@ test("POST /sessions marks a persisted session failed when runtime construction expect(failed).toBe(1); }); +test.each([ + [ + "new", + "auth.json", + '{"unrelated":{"credential":{"nested":["!post-session-auth-command route-secret /private/route-auth"]}}}\n', + "Session startup failed. Check configuration and connectivity, then Resume the session.", + ], + [ + "resume", + "models.json", + '{"providers":{"local-qwen":{"models":[{"id":"qwen3.6-35b-a3b","headers":{"x":"!post-session-model-command route-secret /private/route-model"}}]}}}\n', + "Session could not be resumed. Check configuration and connectivity, then try again.", + ], +] as const)( + "POST %s refuses executable %s changed after session persistence without spawning or leaking", + async (flow, changedFile, unsafeRaw, publicMessage) => { + const agentDir = mkdtempSync(path.join(tmpdir(), "tht-route-runtime-config-")); + const safeAuth = '{"deepseek":{"type":"api_key","key":"safe-token"}}\n'; + const safeModels = '{"providers":{"local-qwen":{"baseUrl":"http://model.invalid/v1","models":[{"id":"qwen3.6-35b-a3b"}]}}}\n'; + writeFileSync(path.join(agentDir, "auth.json"), safeAuth, { mode: 0o600 }); + writeFileSync(path.join(agentDir, "models.json"), safeModels, { mode: 0o600 }); + vi.stubEnv("PI_CODING_AGENT_DIR", agentDir); + const cfg = loadConfig({ THT_HARNESS_DIR: "../harness" }); + let authResolutions = 0; + let spawns = 0; + const mgr = new PiProcessManager(cfg, { + authProviders: () => { authResolutions += 1; return new Set(); }, + spawnFn: () => { spawns += 1; throw new Error("ROUTE_SPAWN_BOUNDARY_REACHED"); }, + }); + const hub = new SseHub(); + const events: Array<{ event: string; data: object }> = []; + const sessionId = `post-persistence-${flow}`; + hub.subscribe(sessionId, (event, data) => events.push({ event, data })); + const failSession = vi.fn(async () => {}); + const consoleError = vi.spyOn(console, "error").mockImplementation(() => undefined); + const validateEarlierState = () => { + validateDeclarativePiConfig(readFileSync(path.join(agentDir, "auth.json"), "utf8")); + validateDeclarativePiConfig(readFileSync(path.join(agentDir, "models.json"), "utf8")); + }; + + if (flow === "resume") { + // This session was admitted and persisted while both mounted files were safe. + validateEarlierState(); + writeFileSync(path.join(agentDir, changedFile), unsafeRaw, { mode: 0o600 }); + } + + const app = buildApp(cfg, { + mgr, + hub, + thtRunner: { + sessionNew: async () => { + // Model admission completed immediately above this persistence boundary. + writeFileSync(path.join(agentDir, changedFile), unsafeRaw, { mode: 0o600 }); + return { id: sessionId }; + }, + failSession, + searchPack: async () => {}, + sessionShow: async () => ({ + id: sessionId, + status: "open", + archived: false, + provider: "local-qwen", + model: "qwen3.6-35b-a3b", + thinking: "low", + }), + reopenSession: async () => {}, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ + workspace: "w", + provider: "local-qwen", + model: "qwen3.6-35b-a3b", + thinking: "low", + }) as any, + listModels: async () => { + validateEarlierState(); + return [{ + provider: "local-qwen", id: "qwen3.6-35b-a3b", name: "Qwen", reasoning: true, + }]; + }, + }); + + try { + const response = flow === "new" + ? await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } }) + : await app.inject({ method: "POST", url: `/sessions/${sessionId}/resume` }); + const logs = consoleError.mock.calls.flat().map(String).join(" "); + + expect(response.statusCode).toBe(503); + expect(response.json()).toEqual({ error: publicMessage }); + expect({ authResolutions, spawns, runtimes: mgr.count() }).toEqual({ + authResolutions: 0, spawns: 0, runtimes: 0, + }); + expect(failSession).toHaveBeenCalledTimes(flow === "new" ? 1 : 0); + expect(events).toEqual([]); + expect(`${response.body}\n${logs}`).not.toContain(unsafeRaw.trim()); + expect(`${response.body}\n${logs}`).not.toMatch(/route-secret|post-session-(?:auth|model)-command|\/private\/route-|tht-route-runtime-config/); + } finally { + await app.close(); + consoleError.mockRestore(); + vi.unstubAllEnvs(); + rmSync(agentDir, { recursive: true, force: true }); + } + }, +); + test("POST /sessions/:id/resume returns 409 for a read-only session without calling ollamaEnsure", async () => { let ensureCalled = false; const app = mutApp({