import { spawn } from "node:child_process"; import { createHash, randomUUID } from "node:crypto"; import { closeSync, constants as fsConstants, existsSync, fchmodSync, fstatSync, fsyncSync, lstatSync, mkdirSync, openSync, readFileSync, readSync, realpathSync, statSync, unlinkSync, writeFileSync, } from "node:fs"; import { dirname, isAbsolute, join, relative, resolve } from "node:path"; import { parseAllDocuments } from "yaml"; import { clearPrincipalEnvironment, principalEnvironment, type PrincipalContext } from "../auth/principal.js"; import { secretValue, type SecretBundleConfig } from "../config/secret-bundle.js"; import { resolveRuntimeBindings } from "../workspaces/bindings.js"; import { renderRuntimeConfig, type RuntimeInstallationOverlay, type RuntimePaths } from "../workspaces/runtime-renderer.js"; import { parseWorkspaceYaml } from "../workspaces/schema.js"; export interface ThtConfig extends SecretBundleConfig { thtBin: string; harnessDir: string; configPath: string; dataRoot?: string; runtimeSnapshotRoot?: string; secretRoots?: readonly string[]; } export interface RuntimeConfigLease { path: string; workspaceId: string; workspaceRevision: string; release(): void; } export interface SessionRow { id: string; status: string; question: string; summary: string | null; created_at: string; updated_at: string | null; author: string | null; workspace_id?: string | null; workspace_revision?: string | null; archived?: boolean; } export interface SessionDocument { phase: string; key: string; title: string; format: string; content: string; } export interface OllamaEnsureResult { ok: boolean; stage?: string; error?: string; server?: string; model?: string; model_name?: string; } interface RuntimeSnapshot { path: string; dev: number; ino: number; size: number; mode: number; digest: string; } export class ThtRunner { private readonly runtimeSnapshots = new Map(); constructor(private cfg: ThtConfig, private principal?: PrincipalContext) {} /** Bind one trusted request principal to every child spawned by this runner. */ withPrincipal(principal: PrincipalContext): ThtRunner { return new ThtRunner(this.cfg, principal); } /** * Resolve the `-c ` args. A named workspace MUST exist: silently falling * back to the default config would point every operation at the wrong workspace * (wrong DB, wrong sessions dir) — fail loud instead. */ private configArg(workspaceConfigPath?: string): string[] { if (workspaceConfigPath) { if (isAbsolute(workspaceConfigPath)) { if (this.runtimeSnapshots.has(workspaceConfigPath)) this.assertTrustedRuntimeSnapshot(workspaceConfigPath); else this.assertWorkspaceSnapshot(workspaceConfigPath); return ["-c", workspaceConfigPath]; } if (workspaceConfigPath.includes("/")) { throw new Error("workspace snapshot config path must be absolute"); } const workspace = workspaceConfigPath; if (!existsSync(join(this.cfg.harnessDir, "workspaces", `${workspace}.yaml`))) { throw new Error(`workspace non trovato: workspaces/${workspace}.yaml (harness: ${this.cfg.harnessDir})`); } return ["-c", `workspaces/${workspace}.yaml`]; } return ["-c", this.cfg.configPath]; } private assertWorkspaceSnapshot(path: string): { workspaceId: string; workspaceRevision: string } { if (!this.cfg.runtimeSnapshotRoot) throw new Error("workspace snapshot root is not configured"); const snapshotsRoot = dirname(this.cfg.runtimeSnapshotRoot); const pathRelative = relative(snapshotsRoot, path); const match = /^([0-9a-f]{40})\/([a-z][a-z0-9-]{2,62})\.yaml$/.exec(pathRelative); if ( pathRelative.startsWith("..") || isAbsolute(pathRelative) || !match ) throw new Error("config path is not a trusted runtime snapshot"); const entry = lstatSync(path); if (!entry.isFile() || entry.isSymbolicLink()) { throw new Error("config path is not a trusted runtime snapshot"); } return { workspaceRevision: match[1], workspaceId: match[2] }; } private readCanonicalWorkspaceSnapshot(path: string): { workspace: ReturnType; workspaceId: string; workspaceRevision: string; } { const identity = this.assertWorkspaceSnapshot(path); const fd = openSync(path, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW); try { const before = fstatSync(fd); if (!before.isFile()) throw new Error("workspace snapshot is not a file"); const source = readFileSync(fd, "utf8"); const after = fstatSync(fd); if (before.dev !== after.dev || before.ino !== after.ino || before.size !== after.size) { throw new Error("workspace snapshot changed while reading"); } const workspace = parseWorkspaceYaml(source); if (workspace.workspace.id !== identity.workspaceId) { throw new Error("workspace snapshot identity does not match its path"); } return { workspace, ...identity }; } finally { closeSync(fd); } } private runtimePaths(workspaceId: string): RuntimePaths { if (!this.cfg.dataRoot || !isAbsolute(this.cfg.dataRoot)) { throw new Error("registry workspace runtime requires an absolute data root"); } // The portable stack persists one `sessions` store at /sessions. Keep every // workspace's mutable harness roots below that mounted boundary. const root = join(this.cfg.dataRoot, "sessions", workspaceId); return { sessions: join(root, "sessions"), artifacts: join(root, "artifacts"), indexes: join(root, "indexes"), }; } private installationOverlay(): RuntimeInstallationOverlay { const path = isAbsolute(this.cfg.configPath) ? this.cfg.configPath : resolve(this.cfg.harnessDir, this.cfg.configPath); if (!existsSync(path)) return {}; const documents = parseAllDocuments(readFileSync(path, "utf8"), { uniqueKeys: true }); if (documents.length !== 1) throw new Error("installation config must contain one YAML document"); const document = documents[0]; if (document.errors.length > 0 || document.warnings.length > 0) { throw new Error("installation config contains invalid YAML"); } const parsed = document.toJSON(); if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) { throw new Error("installation config must be a YAML mapping"); } const source = parsed as Record; return { ...(source.session_storage === undefined ? {} : { session_storage: source.session_storage }), ...(source.profile === undefined ? {} : { profile: source.profile }), }; } /** Render one immutable canonical registry revision into a backend-owned harness config. */ acquireWorkspaceRuntime(workspaceConfigPath: string): RuntimeConfigLease { const canonical = this.readCanonicalWorkspaceSnapshot(workspaceConfigPath); const bindings = resolveRuntimeBindings( canonical.workspace, process.env, this.cfg.secretRoots ?? [], ); const config = renderRuntimeConfig( canonical.workspace, bindings, this.runtimePaths(canonical.workspaceId), canonical, this.installationOverlay(), ); const path = this.createRuntimeSnapshot(config); let released = false; return { path, workspaceId: canonical.workspaceId, workspaceRevision: canonical.workspaceRevision, release: () => { if (released) return; released = true; this.cleanupRuntimeSnapshot(path); }, }; } private runtimeSnapshotDirectory(): string { if (!this.cfg.runtimeSnapshotRoot) throw new Error("runtime snapshot root is not configured"); if (!isAbsolute(this.cfg.runtimeSnapshotRoot)) throw new Error("runtime snapshot root must be absolute"); mkdirSync(this.cfg.runtimeSnapshotRoot, { recursive: true, mode: 0o700 }); const directory = lstatSync(this.cfg.runtimeSnapshotRoot); if (!directory.isDirectory() || directory.isSymbolicLink() || (directory.mode & 0o077) !== 0) { throw new Error("runtime snapshot root is not trusted"); } return realpathSync(this.cfg.runtimeSnapshotRoot); } private static isRestrictiveMode(mode: number): boolean { const permissions = mode & 0o777; return (permissions & 0o400) !== 0 && (permissions & ~0o600) === 0; } private assertTrustedRuntimeSnapshot(path: string): RuntimeSnapshot { const snapshot = this.runtimeSnapshots.get(path); if (!snapshot) throw new Error("config path is not a trusted runtime snapshot"); try { const entry = lstatSync(path); const stat = statSync(path); if ( !entry.isFile() || entry.isSymbolicLink() || stat.dev !== snapshot.dev || stat.ino !== snapshot.ino || stat.size !== snapshot.size || (stat.mode & 0o777) !== snapshot.mode || !ThtRunner.isRestrictiveMode(stat.mode) || createHash("sha256").update(readFileSync(path)).digest("hex") !== snapshot.digest ) throw new Error("changed runtime snapshot"); return snapshot; } catch { throw new Error("config path is not a trusted runtime snapshot"); } } /** Create an opaque, backend-owned temporary config that is safe to hand to `tht`. */ createRuntimeSnapshot(config: string): string { const directory = this.runtimeSnapshotDirectory(); const path = join(directory, `runtime-${randomUUID()}.yaml`); const fd = openSync( path, fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_EXCL | fsConstants.O_NOFOLLOW, 0o600, ); try { writeFileSync(fd, config, "utf8"); fsyncSync(fd); fchmodSync(fd, 0o400); const stat = fstatSync(fd); this.runtimeSnapshots.set(path, { path, dev: stat.dev, ino: stat.ino, size: stat.size, mode: stat.mode & 0o777, digest: createHash("sha256").update(config, "utf8").digest("hex"), }); return path; } catch (error) { try { unlinkSync(path); } catch { /* creation did not produce a removable file */ } throw error; } finally { closeSync(fd); } } cleanupRuntimeSnapshot(path: string): void { const snapshot = this.runtimeSnapshots.get(path); if (!snapshot) return; this.runtimeSnapshots.delete(path); try { unlinkSync(snapshot.path); } catch { /* a changed path is never removed recursively */ } } private openTrustedRuntimeSnapshot(path: string): number { const snapshot = this.assertTrustedRuntimeSnapshot(path); const fd = openSync(path, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW); try { const stat = fstatSync(fd); if ( !stat.isFile() || stat.dev !== snapshot.dev || stat.ino !== snapshot.ino || stat.size !== snapshot.size || (stat.mode & 0o777) !== snapshot.mode || !ThtRunner.isRestrictiveMode(stat.mode) ) throw new Error("changed runtime snapshot"); const contents = Buffer.alloc(snapshot.size); let offset = 0; while (offset < contents.length) { const bytes = readSync(fd, contents, offset, contents.length - offset, offset); if (bytes === 0) throw new Error("truncated runtime snapshot"); offset += bytes; } if (createHash("sha256").update(contents).digest("hex") !== snapshot.digest) { throw new Error("changed runtime snapshot"); } return fd; } catch { closeSync(fd); throw new Error("config path is not a trusted runtime snapshot"); } } async runWithRuntimeSnapshot( args: string[], config: string, timeoutMs: number = ThtRunner.DEFAULT_TIMEOUT_MS, ): Promise<{ code: number; stdout: string; stderr: string }> { const snapshot = this.createRuntimeSnapshot(config); try { return await this.run(args, snapshot, timeoutMs); } finally { this.cleanupRuntimeSnapshot(snapshot); } } /** * Build the full argv for a `tht` invocation. `--config`/`-c` is a PER-COMMAND * option in the `tht` CLI (there is NO global `-c`), so it MUST be appended * AFTER the subcommand + its flags, never prepended. */ buildArgv(args: string[], workspace?: string): string[] { return [...args, ...this.configArg(workspace)]; } // Every route awaits these children; without a ceiling, one hung DWH/vector call // (dropped VPN mid-connect) wedges its HTTP request forever. Session/file commands // get the default; DWH-touching commands pass a wider explicit budget. static readonly DEFAULT_TIMEOUT_MS = 60_000; static readonly DWH_TIMEOUT_MS = 120_000; run( args: string[], workspaceConfigPath?: string, timeoutMs: number = ThtRunner.DEFAULT_TIMEOUT_MS, ): Promise<{ code: number; stdout: string; stderr: string }> { if ( workspaceConfigPath && isAbsolute(workspaceConfigPath) && !this.runtimeSnapshots.has(workspaceConfigPath) ) { let runtime: RuntimeConfigLease; try { runtime = this.acquireWorkspaceRuntime(workspaceConfigPath); } catch (error) { return Promise.reject(error); } return this.run(args, runtime.path, timeoutMs).finally(runtime.release); } return new Promise((resolve) => { const env: NodeJS.ProcessEnv = { ...process.env }; delete env.THT_DATA_ROOT; clearPrincipalEnvironment(env); if (this.cfg.dataRoot !== undefined) env.THT_DATA_ROOT = this.cfg.dataRoot; if (this.principal) Object.assign(env, principalEnvironment(this.principal)); 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); if (value !== undefined) env[name] = value; } const ca = secretValue(this.cfg, "THT_SSL_CA") ?? secretValue(this.cfg, "THT_CA"); if (ca !== undefined) { env.THT_CA = ca; env.THT_SSL_CA = ca; } let snapshotFd: number | undefined; let ch; try { snapshotFd = workspaceConfigPath && this.runtimeSnapshots.has(workspaceConfigPath) ? this.openTrustedRuntimeSnapshot(workspaceConfigPath) : undefined; ch = spawn( this.cfg.thtBin, snapshotFd === undefined ? this.buildArgv(args, workspaceConfigPath) : [...args, "-c", "/dev/fd/3"], { cwd: this.cfg.harnessDir, env, ...(snapshotFd === undefined ? {} : { stdio: ["ignore", "pipe", "pipe", snapshotFd] }), }, ); } finally { if (snapshotFd !== undefined) closeSync(snapshotFd); } let stdout = ""; let stderr = ""; let settled = false; let timer: ReturnType | undefined; const finish = (result: { code: number; stdout: string; stderr: string }) => { if (settled) return; settled = true; if (timer) clearTimeout(timer); resolve(result); }; if (timeoutMs !== undefined) { timer = setTimeout(() => { try { ch.kill("SIGKILL"); } catch { /* already gone */ } finish({ code: 124, stdout, stderr: stderr || `timed out after ${timeoutMs}ms` }); }, timeoutMs); } ch.stdout?.on("data", (d: Buffer) => (stdout += d)); ch.stderr?.on("data", (d: Buffer) => (stderr += d)); ch.on("error", (error) => finish({ code: 1, stdout, stderr: stderr || error.message })); ch.on("close", (code) => finish({ code: code ?? 0, stdout, stderr })); }); } private async json(args: string[], workspace?: string, timeoutMs?: number): Promise { const { code, stdout, stderr } = await this.run(args, workspace, timeoutMs); if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`); return JSON.parse(stdout) as T; } private async ok(args: string[], workspace?: string): Promise { const { code, stderr } = await this.run(args, workspace); if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`); } async sessionNew(o: { question: string; provider?: string; model?: string; thinking?: string; name?: string; /** Legacy named-workspace compatibility; pinned sessions use workspaceConfigPath. */ workspace?: string; workspaceConfigPath?: string; workspaceId?: string; workspaceRevision?: string; }) { const a = ["session", "new", o.question]; for (const [f, v] of [ ["--provider", o.provider], ["--model", o.model], ["--thinking", o.thinking], ["--name", o.name], ["--workspace-id", o.workspaceId], ["--workspace-revision", o.workspaceRevision], ] as const) { if (v) a.push(f, v); } a.push("--json"); return this.json<{ id: string }>(a, o.workspaceConfigPath ?? o.workspace); } /** Build and persist the deterministic F1 retrieval pack for a new session. */ async searchPack(question: string, sessionId: string, workspace?: string): Promise { const args = ["search", "pack", question, "--session", sessionId]; const { code, stderr } = await this.run(args, workspace, ThtRunner.DWH_TIMEOUT_MS); if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`); } /** * Probe DWH reachability via `tht db ping` (a REST health check). Never throws — * returns ok=false with the failure detail so the caller can refuse a new session * cleanly. Timed out to keep POST /sessions responsive when the host hangs. */ async dbPing(workspace?: string): Promise<{ ok: boolean; detail: string }> { const { code, stdout, stderr } = await this.run(["db", "ping"], workspace, 5_000); if (code === 0) return { ok: true, detail: stdout.trim() }; return { ok: false, detail: (stderr || stdout).trim() }; } sessionList(workspace?: string) { return this.json(["session", "list", "--json"], workspace); } sessionShow(id: string, workspace?: string) { return this.json(["session", "show", id, "--json"], workspace); } sqlPreview(id: string, p: { limit?: number; offset?: number }, workspace?: string) { // No positional FILE: the harness resolves sql_final.sql from the session // via _session_sql_file(cfg, session_id), which respects the workspace path. const a = ["sql", "preview", "--session", id, "--json"]; if (p.limit != null) a.push("--limit", String(p.limit)); if (p.offset != null) a.push("--offset", String(p.offset)); return this.json<{ columns: string[]; rows: unknown[][]; execution_ms: number; truncated: boolean; }>(a, workspace, ThtRunner.DWH_TIMEOUT_MS); } async sqlExport(id: string, workspace?: string) { const { code, stdout, stderr } = await this.run( ["sql", "export", "--session", id], workspace, ThtRunner.DWH_TIMEOUT_MS, ); if (code !== 0) throw new Error(`tht sql export exit ${code}: ${stderr.trim()}`); return { path: stdout.trim() }; } closeSession(id: string, workspace?: string) { return this.ok(["session", "close", id], workspace); } failSession(id: string, workspace?: string) { return this.ok(["session", "fail", id], workspace); } reopenSession(id: string, workspace?: string) { return this.ok(["session", "reopen", id], workspace); } setName(id: string, name: string, workspace?: string) { return this.ok(["session", "set-name", id, "--name", name], workspace); } setGroup(id: string, group: string, workspace?: string) { return this.ok(["session", "set-group", id, "--group", group], workspace); } archive(id: string, workspace?: string) { return this.ok(["session", "archive", id], workspace); } unarchive(id: string, workspace?: string) { return this.ok(["session", "unarchive", id], workspace); } async deleteSession(id: string, workspace?: string) { const { code, stderr } = await this.run(["session", "delete", id], workspace); if (code !== 0) throw new Error(`tht session delete exit ${code}: ${stderr.trim()}`); } documents(id: string, workspace?: string) { return this.json(["session", "documents", id, "--json"], workspace); } preferencesGet(workspace?: string) { return this.json>(["session", "preferences", "get", "--json"], workspace); } async preferencesSet(preferences: Record, workspace?: string): Promise { await this.ok(["session", "preferences", "set", "--json", JSON.stringify(preferences)], workspace); } async ollamaEnsure(workspace: string, timeoutSec: number): Promise { // Process budget wider than the CLI's own --timeout so the CLI reports its // failure itself; SIGKILL is only the backstop for a wedged child. const { code, stdout, stderr } = await this.run( ["ollama", "ensure", "--json", "--timeout", String(timeoutSec)], workspace, timeoutSec * 1000 + 30_000, ); let parsed: Partial | null = null; try { parsed = JSON.parse(stdout.trim()); } catch { /* not JSON */ } // Exit 0 with unparseable output is NOT a verified readiness: --json promises // pristine JSON, so treat the violation as a failed check, never as ok. if (code === 0 && parsed !== null) return { ok: true, ...parsed }; if (code === 0) { return { ok: false, error: `tht ollama ensure: output non-JSON: ${stdout.trim().slice(0, 200)}` }; } return { ok: false, stage: parsed?.stage, error: parsed?.error ?? (stderr.trim() || `tht ollama ensure exit ${code}`), }; } }