New session now refuses to spawn a Pi runtime that would only die in bootstrap
retrieval when the DWH/vector host is unreachable (e.g. a dropped VPN). Before
`session new`, POST /sessions probes the DWH via `tht db ping`; if it is down it
returns 503 {code:"dwh_unreachable"} with a clear message and creates nothing.
- Gated behind the THT_DWH_PRECHECK flag (default off), enabled only by the local
dev launcher (run-stack.sh) — containers/CI never pay the probe, and existing
tests that don't set it are unaffected.
- ThtRunner.dbPing() runs `tht db ping` with a 10s timeout (run() gains an optional
timeout that SIGKILLs a hung child).
- Frontend: apiFetch throws a typed ApiError (status + parsed payload); the new-
session composer shows the specific alert on `dwh_unreachable` instead of the
generic retry hint, keeping the question for retry.
Verified live on an isolated backend (precheck on + broken DWH host → 503
dwh_unreachable, no session created) and via unit tests (backend 228, frontend 308).
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
214 lines
8.4 KiB
TypeScript
214 lines
8.4 KiB
TypeScript
import { spawn } from "node:child_process";
|
|
import { existsSync } from "node:fs";
|
|
import { join } from "node:path";
|
|
import { clearPrincipalEnvironment, principalEnvironment, type PrincipalContext } from "../auth/principal.js";
|
|
|
|
export interface ThtConfig {
|
|
thtBin: string;
|
|
harnessDir: string;
|
|
configPath: string;
|
|
dataRoot?: string;
|
|
}
|
|
|
|
export interface SessionRow {
|
|
id: string;
|
|
status: string;
|
|
question: string;
|
|
summary: string | null;
|
|
created_at: string;
|
|
updated_at: string | null;
|
|
author: string | null;
|
|
}
|
|
|
|
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;
|
|
}
|
|
|
|
export class ThtRunner {
|
|
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 <config>` args. If `workspace` is given AND a matching
|
|
* `workspaces/<workspace>.yaml` exists under harnessDir, select it; otherwise
|
|
* fall back to the default configPath.
|
|
*/
|
|
private configArg(workspace?: string): string[] {
|
|
if (workspace && existsSync(join(this.cfg.harnessDir, "workspaces", `${workspace}.yaml`))) {
|
|
return ["-c", `workspaces/${workspace}.yaml`];
|
|
}
|
|
return ["-c", this.cfg.configPath];
|
|
}
|
|
|
|
/**
|
|
* 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)];
|
|
}
|
|
|
|
run(
|
|
args: string[], workspace?: string, timeoutMs?: number,
|
|
): Promise<{ code: number; stdout: string; stderr: string }> {
|
|
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));
|
|
const ch = spawn(this.cfg.thtBin, this.buildArgv(args, workspace), {
|
|
cwd: this.cfg.harnessDir,
|
|
env,
|
|
});
|
|
let stdout = "";
|
|
let stderr = "";
|
|
let settled = false;
|
|
let timer: ReturnType<typeof setTimeout> | 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<T>(args: string[], workspace?: string): Promise<T> {
|
|
const { code, stdout, stderr } = await this.run(args, workspace);
|
|
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<void> {
|
|
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;
|
|
workspace?: string;
|
|
}) {
|
|
const a = ["session", "new", o.question];
|
|
for (const [f, v] of [
|
|
["--provider", o.provider],
|
|
["--model", o.model],
|
|
["--thinking", o.thinking],
|
|
["--name", o.name],
|
|
] as const) {
|
|
if (v) a.push(f, v);
|
|
}
|
|
a.push("--json");
|
|
return this.json<{ id: string }>(a, o.workspace);
|
|
}
|
|
|
|
/** Build and persist the deterministic F1 retrieval pack for a new session. */
|
|
async searchPack(question: string, sessionId: string, workspace?: string): Promise<void> {
|
|
const args = ["search", "pack", question, "--session", sessionId];
|
|
const { code, stderr } = await this.run(args, workspace);
|
|
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, 10_000);
|
|
if (code === 0) return { ok: true, detail: stdout.trim() };
|
|
return { ok: false, detail: (stderr || stdout).trim() };
|
|
}
|
|
|
|
sessionList(workspace?: string) {
|
|
return this.json<SessionRow[]>(["session", "list", "--json"], workspace);
|
|
}
|
|
|
|
sessionShow(id: string, workspace?: string) {
|
|
return this.json<unknown>(["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);
|
|
}
|
|
|
|
async sqlExport(id: string, workspace?: string) {
|
|
const { code, stdout, stderr } = await this.run(["sql", "export", "--session", id], workspace);
|
|
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<SessionDocument[]>(["session", "documents", id, "--json"], workspace); }
|
|
|
|
preferencesGet(workspace?: string) { return this.json<Record<string, unknown>>(["session", "preferences", "get", "--json"], workspace); }
|
|
async preferencesSet(preferences: Record<string, unknown>, workspace?: string): Promise<void> {
|
|
await this.ok(["session", "preferences", "set", "--json", JSON.stringify(preferences)], workspace);
|
|
}
|
|
|
|
async ollamaEnsure(workspace: string, timeoutSec: number): Promise<OllamaEnsureResult> {
|
|
const { code, stdout, stderr } = await this.run(
|
|
["ollama", "ensure", "--json", "--timeout", String(timeoutSec)],
|
|
workspace,
|
|
);
|
|
let parsed: Partial<OllamaEnsureResult> = {};
|
|
try { parsed = JSON.parse(stdout.trim() || "{}"); } catch { /* leave {} */ }
|
|
if (code === 0) return { ok: true, ...parsed };
|
|
return {
|
|
ok: false,
|
|
stage: parsed.stage,
|
|
error: parsed.error ?? (stderr.trim() || `tht ollama ensure exit ${code}`),
|
|
};
|
|
}
|
|
}
|