Implement approved specification #32 and tickets #33-#37. Keep host authentication server-verified and pin session interaction language. Compile scoped base selectors for browser compatibility and retain full gutters during CSS pruning.
647 lines
26 KiB
TypeScript
647 lines
26 KiB
TypeScript
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 { parse, parseAllDocuments } from "yaml";
|
|
import { clearPrincipalEnvironment, principalEnvironment, type PrincipalContext } from "../auth/principal.js";
|
|
import { secretValue, type SecretBundleConfig } from "../config/secret-bundle.js";
|
|
import { renderWorkspaceRuntimeFromSnapshotPath } from "../workspaces/runtime-config-lease.js";
|
|
import {
|
|
DEFAULT_SEMANTIC_RUNTIME,
|
|
type RuntimeInstallationOverlay,
|
|
type RuntimePaths,
|
|
type SemanticRuntimeConfig,
|
|
} from "../workspaces/runtime-renderer.js";
|
|
import {
|
|
parseWorkspaceYaml,
|
|
validateOperationalWorkspace,
|
|
type WorkspaceDescriptor,
|
|
} from "../workspaces/schema.js";
|
|
import { reconcileCollection, type CollectionMode } from "../workspaces/qdrant-collection.js";
|
|
import type { WorkspaceSecretStore } from "../workspaces/secret-store.js";
|
|
import type { CatalogRepository } from "../catalog/types.js";
|
|
import { preprocessingInputFingerprint } from "../workspaces/effective-config.js";
|
|
import { workspaceVectorCollections } from "../workspaces/vector-collections.js";
|
|
|
|
export interface ThtConfig extends SecretBundleConfig {
|
|
thtBin: string;
|
|
harnessDir: string;
|
|
configPath: string;
|
|
dataRoot?: string;
|
|
runtimeSnapshotRoot?: string;
|
|
secretRoots?: readonly string[];
|
|
semanticRuntime: SemanticRuntimeConfig;
|
|
qdrantRequest?: typeof fetch;
|
|
/** "self_heal" for session admission (create missing collections/indexes), default "require_existing". */
|
|
qdrantCollectionMode?: "self_heal" | "require_existing";
|
|
workspaceSecretStore?: WorkspaceSecretStore;
|
|
catalogRepository?: CatalogRepository;
|
|
}
|
|
|
|
export interface RuntimeConfigLease {
|
|
path: string;
|
|
workspaceId: string;
|
|
workspaceRevision: string;
|
|
inputFingerprint: 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;
|
|
interaction_language?: 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;
|
|
}
|
|
|
|
export type SemanticReadinessCode = "workspace_not_activatable" | "semantic_index_incompatible";
|
|
|
|
export interface QdrantEnsureResult {
|
|
ok: boolean;
|
|
code?: SemanticReadinessCode;
|
|
state?: "ready" | "created" | "repaired" | "upgraded";
|
|
}
|
|
|
|
const REQUIRED_QDRANT_PAYLOAD_INDEXES = [
|
|
"content_hash",
|
|
"document_id",
|
|
"kind",
|
|
"record_key",
|
|
"record_kind",
|
|
"vector_generation",
|
|
"workspace_id",
|
|
"workspace_revision",
|
|
] as const;
|
|
|
|
interface RuntimeSnapshot {
|
|
path: string;
|
|
dev: number;
|
|
ino: number;
|
|
size: number;
|
|
mode: number;
|
|
digest: string;
|
|
}
|
|
|
|
export class ThtRunner {
|
|
private readonly runtimeSnapshots = new Map<string, RuntimeSnapshot>();
|
|
|
|
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. 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<typeof parseWorkspaceYaml>;
|
|
workspaceId: string;
|
|
workspaceRevision: string;
|
|
revisionContentRoot: 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 = validateOperationalWorkspace(parseWorkspaceYaml(source));
|
|
if (workspace.workspace.id !== identity.workspaceId) {
|
|
throw new Error("workspace snapshot identity does not match its path");
|
|
}
|
|
return { workspace, ...identity, revisionContentRoot: dirname(path) };
|
|
} 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 <dataRoot>/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"),
|
|
memory: join(root, "memory"),
|
|
};
|
|
}
|
|
|
|
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<string, unknown>;
|
|
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. */
|
|
async acquireWorkspaceRuntime(workspaceConfigPath: string): Promise<RuntimeConfigLease> {
|
|
const identity = this.assertWorkspaceSnapshot(workspaceConfigPath);
|
|
const catalogDatabase = this.cfg.catalogRepository === undefined
|
|
? undefined
|
|
: await this.cfg.catalogRepository.getByWorkspace(identity.workspaceId);
|
|
if (this.cfg.catalogRepository !== undefined && !catalogDatabase) {
|
|
throw new Error("workspace database is not configured in the Catalog");
|
|
}
|
|
const rendered = renderWorkspaceRuntimeFromSnapshotPath({
|
|
snapshotPath: workspaceConfigPath,
|
|
harnessDir: this.cfg.harnessDir,
|
|
configPath: this.cfg.configPath,
|
|
dataRoot: this.cfg.dataRoot ?? (() => {
|
|
throw new Error("registry workspace runtime requires an absolute data root");
|
|
})(),
|
|
secretRoots: this.cfg.secretRoots ?? [],
|
|
semanticRuntime: this.cfg.semanticRuntime ?? DEFAULT_SEMANTIC_RUNTIME,
|
|
workspaceSecretStore: this.cfg.workspaceSecretStore,
|
|
catalogDatabase,
|
|
});
|
|
let path: string;
|
|
try {
|
|
path = this.createRuntimeSnapshot(rendered.renderedConfig);
|
|
} catch (error) {
|
|
rendered.releaseSecrets();
|
|
throw error;
|
|
}
|
|
let released = false;
|
|
return {
|
|
path,
|
|
workspaceId: rendered.workspaceId,
|
|
workspaceRevision: rendered.workspaceRevision,
|
|
inputFingerprint: preprocessingInputFingerprint(
|
|
rendered.workspaceId,
|
|
rendered.workspaceRevision,
|
|
parse(rendered.renderedConfig),
|
|
),
|
|
release: () => {
|
|
if (released) return;
|
|
released = true;
|
|
this.cleanupRuntimeSnapshot(path);
|
|
rendered.releaseSecrets();
|
|
},
|
|
};
|
|
}
|
|
|
|
async workspaceInputFingerprint(workspaceConfigPath: string): Promise<string> {
|
|
const lease = await this.acquireWorkspaceRuntime(workspaceConfigPath);
|
|
try {
|
|
return lease.inputFingerprint;
|
|
} finally {
|
|
lease.release();
|
|
}
|
|
}
|
|
|
|
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)
|
|
) {
|
|
return this.acquireWorkspaceRuntime(workspaceConfigPath).then((runtime) => (
|
|
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<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, timeoutMs?: number): Promise<T> {
|
|
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<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;
|
|
interactionLanguage?: 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],
|
|
["--interaction-language", o.interactionLanguage],
|
|
["--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<void> {
|
|
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<SessionRow[]>(["session", "list", "--json"], workspace);
|
|
}
|
|
|
|
sessionShow(id: string, workspace?: string) {
|
|
return this.json<unknown>(["session", "show", id, "--json"], workspace);
|
|
}
|
|
|
|
ensureInteractionLanguage(id: string, workspace?: string) {
|
|
return this.json<{ interaction_language: string }>(
|
|
["session", "ensure-interaction-language", 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<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> {
|
|
// 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<OllamaEnsureResult> | 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}`),
|
|
};
|
|
}
|
|
|
|
async qdrantEnsure(
|
|
workspace: WorkspaceDescriptor,
|
|
timeoutSec: number,
|
|
mode: CollectionMode = "require_existing",
|
|
): Promise<QdrantEnsureResult> {
|
|
let descriptor;
|
|
try {
|
|
descriptor = validateOperationalWorkspace(workspace);
|
|
} catch {
|
|
return { ok: false, code: "workspace_not_activatable" };
|
|
}
|
|
const collections = workspaceVectorCollections(descriptor.workspace.id);
|
|
const controller = new AbortController();
|
|
const timer = setTimeout(() => controller.abort(), Math.max(1, timeoutSec) * 1000);
|
|
try {
|
|
const checked = await Promise.all(Object.entries(collections).map(async ([purpose, collection]) =>
|
|
await reconcileCollection({
|
|
baseUrl: this.cfg.semanticRuntime.internalQdrantUrl,
|
|
collection,
|
|
dimensions: this.cfg.semanticRuntime.internalEmbeddingDimensions,
|
|
distance: "cosine",
|
|
mode: mode === "evidence_maintenance" && purpose === "memory" ? "self_heal" : mode,
|
|
request: this.cfg.qdrantRequest ?? fetch,
|
|
signal: controller.signal,
|
|
})));
|
|
if (!checked.every((result) => result.ok)) {
|
|
return { ok: false, code: "semantic_index_incompatible" };
|
|
}
|
|
const state = (["upgraded", "repaired", "created", "ready"] as const)
|
|
.find((candidate) => checked.some((result) => result.state === candidate));
|
|
return { ok: true, ...(state ? { state } : {}) };
|
|
} catch {
|
|
return { ok: false, code: "workspace_not_activatable" };
|
|
} finally {
|
|
clearTimeout(timer);
|
|
controller.abort();
|
|
}
|
|
}
|
|
}
|