fix: bind pi runtime config at spawn

This commit is contained in:
2026-08-05 05:23:58 +02:00
parent 2703864572
commit 4422568f61
4 changed files with 375 additions and 46 deletions
+82 -8
View File
@@ -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 */ }
},
};
}
+51 -37
View File
@@ -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<string>;
private loadAuthProviders: (agentDir: string) => ReadonlySet<string>;
constructor(
private cfg: AppConfig,
opts?: { spawnFn?: SpawnFn; authProviders?: () => ReadonlySet<string> },
opts?: { spawnFn?: SpawnFn; authProviders?: (agentDir: string) => ReadonlySet<string> },
) {
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;
}
}
+134 -1
View File
@@ -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();
+108
View File
@@ -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({