Files
ThothII/backend/src/workspace-maintenance.ts
T

337 lines
14 KiB
TypeScript

import { spawn } from "node:child_process";
import { closeSync, constants as fsConstants, openSync } from "node:fs";
import { readdir, readFile } from "node:fs/promises";
import { join } from "node:path";
import { parse } from "yaml";
import { loadConfig } from "./config.js";
import { ThtRunner } from "./tht/tht-runner.js";
import { WorkspaceRegistry } from "./workspaces/registry.js";
import { publishDeterministicRuntimeConfigLease, renderActiveWorkspaceRuntime } from "./workspaces/runtime-config-lease.js";
import { WorkspaceSecretStore } from "./workspaces/secret-store.js";
import { WorkspacePreprocessingService, type WorkspaceOperationResult } from "./workspaces/preprocessing-service.js";
import type { SessionInventoryRow } from "./workspaces/preprocessing-state.js";
export interface WorkspaceMaintenanceIo {
stdin: string;
stdout: string[];
stderr: string[];
writeStdout(value: string): void;
writeStderr(value: string): void;
}
type Command = "inspect" | "preprocess-dwh" | "schema-suggest-fks" | "schema-check" | "schema-accept" | "index-schema" | "preprocess-evidence" | "preprocess-run" | "vector-inspect" | "vector-rebuild";
function failureResult(
operation: string,
workspaceId = "",
code: WorkspaceOperationResult["code"] = "workspace_not_activatable",
): WorkspaceOperationResult {
return {
schemaVersion: 1,
status: "failed",
code,
workspaceId,
workspaceRevision: "",
descriptorBlob: "",
operation,
completedStages: [],
};
}
const STATE_ERROR_CODES: Record<string, WorkspaceOperationResult["code"]> = {
preprocessing_resume_mismatch: "preprocessing_resume_mismatch",
preprocessing_conflict: "preprocessing_conflict",
effective_config_mismatch: "effective_config_mismatch",
};
function boundedJson(result: WorkspaceOperationResult): string {
const encoded = JSON.stringify(result);
if (Buffer.byteLength(encoded, "utf8") > 1024 * 1024) {
return JSON.stringify(failureResult(result.operation || "unknown", result.workspaceId));
}
return encoded;
}
function sanitizeStderr(_error: unknown): string {
// Never return raw exception text: it may embed endpoints, tokens, or SQL.
return "workspace maintenance failed\n";
}
function parseRequest(command: string, stdin: string): Record<string, unknown> {
const parsed = JSON.parse(stdin) as Record<string, unknown>;
if (!parsed || typeof parsed !== "object" || Array.isArray(parsed) || parsed.schemaVersion !== 1) {
throw new Error("invalid request");
}
const allowedByCommand: Record<string, readonly string[]> = {
inspect: ["schemaVersion", "workspaceId"],
"preprocess-dwh": ["schemaVersion", "workspaceId", "resumeRunId"],
"schema-suggest-fks": ["schemaVersion", "workspaceId", "fromSql", "assume", "resumeRunId"],
"schema-check": ["schemaVersion", "workspaceId", "annotationsYaml", "reviewedCandidatesDigest"],
"schema-accept": ["schemaVersion", "workspaceId", "runId", "yes"],
"index-schema": ["schemaVersion", "workspaceId", "resumeRunId"],
"preprocess-evidence": ["schemaVersion", "workspaceId", "dryRun", "resumeRunId"],
"preprocess-run": ["schemaVersion", "workspaceId", "resumeRunId"],
"vector-inspect": ["schemaVersion", "workspaceId"],
"vector-rebuild": ["schemaVersion", "workspaceId", "collection", "confirm", "destroy"],
};
const allowed = allowedByCommand[command];
if (!allowed) throw new Error("unknown command");
if (typeof parsed.workspaceId !== "string") throw new Error("invalid workspace id");
for (const key of Object.keys(parsed)) if (!allowed.includes(key)) throw new Error("unexpected request field");
return parsed;
}
function exitCodeFor(result: WorkspaceOperationResult): number {
if (["succeeded", "unchanged", "dry_run"].includes(result.status)) return 0;
if (result.status === "blocked") return 3;
return 1;
}
async function dispatch(command: Command, service: WorkspacePreprocessingService, request: Record<string, unknown>): Promise<WorkspaceOperationResult> {
switch (command) {
case "inspect":
return await service.inspect({ workspaceId: request.workspaceId as string });
case "preprocess-dwh":
return await service.preprocessDwh({
workspaceId: request.workspaceId as string,
resumeRunId: request.resumeRunId as string | undefined,
});
case "schema-suggest-fks":
return await service.suggestFks({
workspaceId: request.workspaceId as string,
fromSql: request.fromSql as any,
assume: request.assume as any,
resumeRunId: request.resumeRunId as string | undefined,
});
case "schema-check":
return await service.checkSchema({
workspaceId: request.workspaceId as string,
annotationsYaml: request.annotationsYaml as string | undefined,
reviewedCandidatesDigest: request.reviewedCandidatesDigest as string | undefined,
});
case "schema-accept":
return await service.acceptSchema({
workspaceId: request.workspaceId as string,
runId: request.runId as string,
yes: request.yes === true,
});
case "index-schema":
return await service.indexSchema({
workspaceId: request.workspaceId as string,
resumeRunId: request.resumeRunId as string | undefined,
});
case "preprocess-evidence":
return await service.preprocessEvidence({
workspaceId: request.workspaceId as string,
dryRun: request.dryRun as boolean | undefined,
resumeRunId: request.resumeRunId as string | undefined,
});
case "preprocess-run":
return await service.run({
workspaceId: request.workspaceId as string,
resumeRunId: request.resumeRunId as string | undefined,
});
case "vector-inspect":
return await service.vectorInspect({ workspaceId: request.workspaceId as string });
case "vector-rebuild":
return await service.vectorRebuild({
workspaceId: request.workspaceId as string,
collection: request.collection as string | undefined,
confirm: request.confirm as string | undefined,
destroy: request.destroy === true,
});
}
}
export async function runWorkspaceMaintenanceCli(
argv: readonly string[],
service: WorkspacePreprocessingService,
io: WorkspaceMaintenanceIo,
): Promise<number> {
const command = argv[2];
if (!command) {
const result = failureResult("unknown");
io.writeStdout(boundedJson(result));
return 2;
}
try {
const request = parseRequest(command, io.stdin);
const result = await dispatch(command as Command, service, request);
io.writeStdout(boundedJson(result));
return exitCodeFor(result);
} catch (error) {
const failureCode = error instanceof Error
&& "code" in error
&& typeof (error as { code?: unknown }).code === "string"
&& (error as { code: string }).code in STATE_ERROR_CODES
? STATE_ERROR_CODES[(error as { code: string }).code]
: "workspace_not_activatable";
const result = failureResult(command, (() => {
try { return JSON.parse(io.stdin).workspaceId ?? ""; } catch { return ""; }
})(), failureCode);
io.writeStdout(boundedJson(result));
io.writeStderr(sanitizeStderr(error));
const message = String((error as Error).message ?? "");
const requestError = error instanceof SyntaxError
|| message === "invalid request"
|| message === "unknown command"
|| message === "unexpected request field"
|| message === "invalid workspace id";
return command in {
inspect: true, "preprocess-dwh": true, "schema-suggest-fks": true, "schema-check": true,
"schema-accept": true,
"index-schema": true, "preprocess-evidence": true, "preprocess-run": true,
} ? (requestError ? 2 : 1) : 2;
}
}
async function readSessionInventory(dataRoot: string, workspaceId: string): Promise<readonly SessionInventoryRow[]> {
const directory = join(dataRoot, "sessions", workspaceId, "sessions");
try {
const entries = await readdir(directory, { withFileTypes: true });
const rows: SessionInventoryRow[] = [];
for (const entry of entries) {
if (!entry.isDirectory() || entry.isSymbolicLink()) continue;
try {
const source = await readFile(join(directory, entry.name, "session_manifest.yaml"), "utf8");
const manifest = parse(source) as Record<string, unknown>;
rows.push({
id: entry.name,
status: typeof manifest.status === "string" ? manifest.status : "open",
archived: manifest.archived === true,
workspaceRevision: typeof manifest.workspace_revision === "string" ? manifest.workspace_revision : null,
});
} catch {
// fail closed at mutation time by ignoring unreadable manifests from the resumable scan
}
}
return rows;
} catch {
return [];
}
}
function createProductionService(): WorkspacePreprocessingService {
const config = loadConfig(process.env);
const registry = new WorkspaceRegistry(config.workspaceRegistry);
const workspaceSecretStore = new WorkspaceSecretStore({
root: config.workspaceSecretStoreRoot,
runtimeRoot: config.workspaceSecretRuntimeRoot,
installationId: config.workspaceRegistry.installationId,
});
const runner = new ThtRunner({
thtBin: config.thtBin,
harnessDir: config.harnessDir,
configPath: process.env.THT_CONFIG ?? "config/tht.yaml",
dataRoot: config.dataRoot,
runtimeSnapshotRoot: join(config.workspaceRegistry.root, "snapshots", "runtime"),
secretRoots: config.workspaceRegistry.secretRoots,
secretsFile: config.secretsFile,
secretFiles: config.secretFiles,
workspaceSecretStore,
semanticRuntime: {
internalQdrantUrl: config.internalQdrantUrl,
internalEmbeddingUrl: config.internalEmbeddingUrl,
internalEmbeddingModel: config.internalEmbeddingModel,
internalEmbeddingDimensions: config.internalEmbeddingDimensions,
},
});
return new WorkspacePreprocessingService({
dataRoot: config.dataRoot ?? "/data",
httpPrivateHostAllowlist: (process.env.THT_EVIDENCE_PRIVATE_HOST_ALLOWLIST ?? "")
.split(",").map((value) => value.trim()).filter((value) => value.length > 0),
acquireActiveRuntime: async (workspaceId) => {
const active = await renderActiveWorkspaceRuntime({
workspaceId,
registry,
registryConfig: config.workspaceRegistry,
harnessDir: config.harnessDir,
configPath: process.env.THT_CONFIG ?? "config/tht.yaml",
dataRoot: config.dataRoot ?? "/data",
secretRoots: config.workspaceRegistry.secretRoots,
workspaceSecretStore,
semanticRuntime: {
internalQdrantUrl: config.internalQdrantUrl,
internalEmbeddingUrl: config.internalEmbeddingUrl,
internalEmbeddingModel: config.internalEmbeddingModel,
internalEmbeddingDimensions: config.internalEmbeddingDimensions,
},
});
const configLease = await publishDeterministicRuntimeConfigLease({
workspaceId,
registry,
registryConfig: config.workspaceRegistry,
harnessDir: config.harnessDir,
configPath: process.env.THT_CONFIG ?? "config/tht.yaml",
dataRoot: config.dataRoot ?? "/data",
secretRoots: config.workspaceRegistry.secretRoots,
semanticRuntime: {
internalQdrantUrl: config.internalQdrantUrl,
internalEmbeddingUrl: config.internalEmbeddingUrl,
internalEmbeddingModel: config.internalEmbeddingModel,
internalEmbeddingDimensions: config.internalEmbeddingDimensions,
},
workspaceSecretStore,
});
return {
workspace: active.workspace,
workspaceId: active.workspaceId,
workspaceRevision: active.workspaceRevision,
descriptorBlob: active.descriptorBlob,
catalogBlob: active.catalogBlob,
configLease,
};
},
runChild: async ({ argv, configPath }) => {
const configFd = openSync(configPath, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW);
try {
return await new Promise((resolve) => {
const child = spawn(config.thtBin, argv, {
cwd: config.harnessDir,
env: { ...process.env, ...(config.dataRoot ? { THT_DATA_ROOT: config.dataRoot } : {}) },
stdio: ["ignore", "pipe", "pipe", configFd],
});
let stdout = "";
let stderr = "";
child.stdout?.on("data", (chunk: Buffer) => { stdout += chunk.toString("utf8"); });
child.stderr?.on("data", (chunk: Buffer) => { stderr += chunk.toString("utf8"); });
child.on("close", (code) => resolve({ exitCode: code ?? 0, stdout, stderr: stderr.slice(0, 64 * 1024) }));
child.on("error", (error) => resolve({ exitCode: 1, stdout, stderr: String(error.message).slice(0, 4096) }));
});
} finally {
closeSync(configFd);
}
},
listSessions: async (workspaceId) => await readSessionInventory(config.dataRoot ?? "/data", workspaceId),
semanticPreflight: async (workspace) => {
const result = await runner.qdrantEnsure(workspace, 30);
return result.ok ? { ok: true as const } : { ok: false as const, code: result.code ?? "workspace_not_activatable" };
},
evidencePreflight: async (workspace) => {
const result = await runner.qdrantEnsure(workspace, 30, "evidence_maintenance");
return result.ok ? { ok: true as const } : { ok: false as const, code: result.code ?? "workspace_not_activatable" };
},
});
}
if (process.argv[1] && import.meta.url === new URL(`file://${process.argv[1]}`).href) {
const stdout: string[] = [];
const stderr: string[] = [];
const io: WorkspaceMaintenanceIo = {
stdin: await new Promise<string>((resolve) => {
let input = "";
process.stdin.setEncoding("utf8");
process.stdin.on("data", (chunk) => { input += chunk; });
process.stdin.on("end", () => resolve(input));
}),
stdout,
stderr,
writeStdout: (value) => { stdout.push(value); },
writeStderr: (value) => { stderr.push(value); },
};
const exitCode = await runWorkspaceMaintenanceCli(process.argv, createProductionService(), io);
process.stdout.write(stdout.join(""));
if (stderr.length > 0) process.stderr.write(stderr.join("").slice(0, 64 * 1024));
process.exit(exitCode);
}