Files
ThothII/backend/src/workspaces/preprocessing-service.ts
T
Codex 82e2c91f42
Publish documentation / publish (push) Successful in 1m27s
feat: implement memory and evidence administration with guided repairs
Add PostgreSQL-backed memory, editable evidence with source review and activation, and human-approved archive repairs across the harness, API, and UI. Include migrations, deployment support, regression coverage, and validation documentation.

Refresh permissions from validated session roles so existing administrator logins can access newly deployed archive management features.
2026-09-10 10:31:34 +02:00

421 lines
17 KiB
TypeScript

import { randomBytes } from "node:crypto";
import { renameSync, rmSync, writeFileSync, mkdirSync } from "node:fs";
import { join } from "node:path";
import {
continueEvidencePreprocessing,
evidencePolicy,
runEvidenceStage,
type EvidencePreprocessingDependencies,
type EvidencePreprocessingOutcome,
} from "./evidence/preprocessing.js";
import type { WorkspaceDescriptor } from "./schema.js";
import {
PreprocessingStateStore,
type PreprocessingJobState,
} from "./preprocessing-state.js";
import type { DeterministicRuntimeConfigLease } from "./runtime-config-lease.js";
import { buildCatalogMetadataSnapshot } from "../catalog/metadata-snapshot.js";
import type { CatalogRepository } from "../catalog/types.js";
export interface WorkspaceOperationResult {
schemaVersion: 1;
status: "succeeded" | "unchanged" | "dry_run" | "blocked" | "failed";
code:
| "ok" | "workspace_not_found" | "workspace_not_activatable"
| "binding_missing" | "preprocessing_conflict"
| "preprocessing_resume_mismatch" | "manual_review_required"
| "evidence_materialization_required" | "effective_config_mismatch"
| "semantic_index_incompatible" | "annotation_invalid"
| "egress_policy_refused" | "catalog_not_ready" | "preprocessing_clear_failed";
workspaceId: string;
workspaceRevision: string;
descriptorBlob: string;
operation: string;
runId?: string;
childRuns?: Record<string, string>;
completedStages: string[];
counts?: Record<string, number>;
artifactIdentities?: Array<{ kind: string; digest: string }>;
effectiveConfigIdentity?: string;
configFingerprint?: string;
inputFingerprint?: string;
/** Suggested FK annotations YAML for the operator to write to --output (schema suggest-fks). */
suggestedFksYaml?: string;
warnings?: string[];
}
export interface ChildProcessRequest {
argv: string[];
configPath: string;
}
interface ActiveRuntime {
workspace: WorkspaceDescriptor;
workspaceId: string;
workspaceRevision: string;
descriptorBlob: string;
catalogBlob: string;
configLease: DeterministicRuntimeConfigLease;
}
interface ChildProcessResult {
exitCode: number;
stdout: string;
stderr: string;
}
export interface WorkspacePreprocessingServiceDeps {
dataRoot: string;
acquireActiveRuntime(workspaceId: string): Promise<ActiveRuntime>;
runChild(request: ChildProcessRequest): Promise<ChildProcessResult>;
semanticPreflight(workspace: WorkspaceDescriptor): Promise<
{ ok: true } | { ok: false; code: "workspace_not_activatable" | "semantic_index_incompatible" }
>;
evidencePreflight(workspace: WorkspaceDescriptor): Promise<
{ ok: true } | { ok: false; code: "workspace_not_activatable" | "semantic_index_incompatible" }
>;
httpPrivateHostAllowlist?: readonly string[];
catalogRepository?: CatalogRepository;
}
interface RunScope {
runtime: ActiveRuntime;
state: PreprocessingStateStore;
job: PreprocessingJobState;
}
function baseResult(
runtime: ActiveRuntime,
operation: string,
status: WorkspaceOperationResult["status"],
code: WorkspaceOperationResult["code"],
extra: Omit<Partial<WorkspaceOperationResult>, "schemaVersion" | "status" | "code" | "workspaceId" | "workspaceRevision" | "descriptorBlob" | "operation"> = {},
): WorkspaceOperationResult {
return {
schemaVersion: 1,
status,
code,
workspaceId: runtime.workspaceId,
workspaceRevision: runtime.workspaceRevision,
descriptorBlob: runtime.descriptorBlob,
operation,
completedStages: [],
effectiveConfigIdentity: runtime.configLease.effectiveConfigIdentity,
configFingerprint: runtime.configLease.configFingerprint,
inputFingerprint: runtime.configLease.inputFingerprint,
...extra,
};
}
export class WorkspacePreprocessingService {
constructor(private readonly deps: WorkspacePreprocessingServiceDeps) {}
async evidenceSources(options: { workspaceId: string; action: "refresh" | "decide";
sourceId?: string; revision?: string; decision?: "keep" | "replace"; actor?: string }): Promise<WorkspaceOperationResult> {
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
try {
if (options.action === "refresh") {
const refused = evidencePolicy(runtime.workspace.evidence, this.deps.httpPrivateHostAllowlist);
if (refused) return baseResult(runtime, "evidence refresh", "failed", refused.code);
}
const argv = ["evidence", "sources", options.action, "--json", "-c", "/dev/fd/3"];
if (options.action === "decide") {
if (!/^[a-f0-9]{64}$/.test(options.sourceId ?? "") || !/^[a-f0-9]{64}$/.test(options.revision ?? "")
|| !["keep", "replace"].includes(options.decision ?? "")) throw new Error("Invalid source decision");
const preflight = await this.deps.evidencePreflight(runtime.workspace);
if (!preflight.ok) return baseResult(runtime, "evidence sources", "failed", preflight.code);
argv.push("--source-id", options.sourceId!, "--revision", options.revision!, "--decision", options.decision!);
}
argv.push("--actor", options.actor ?? "installation operator");
const result = await this.deps.runChild({ argv, configPath: runtime.configLease.path });
const payload = JSON.parse(result.stdout);
if (result.exitCode !== 0 || payload.status !== "succeeded") throw new Error(
typeof payload.error === "string" ? payload.error.slice(0, 1500) : "Evidence source operation failed");
return baseResult(runtime, `evidence ${options.action}`, "succeeded", "ok", {
counts: this.numberRecord(payload.counts),
warnings: [options.action === "refresh" ? "Source comparisons are ready in Evidence management. Active content is unchanged."
: "Source decision activated locally. Commit and push the Evidence tree manually."],
});
} catch (error) {
return baseResult(runtime, `evidence ${options.action}`, "failed", "evidence_materialization_required", {
warnings: [error instanceof Error ? error.message : "Evidence source operation failed"],
});
}
}
async consolidateEvidence(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
const preflight = await this.deps.evidencePreflight(runtime.workspace);
if (!preflight.ok) return baseResult(runtime, "evidence consolidate", "failed", preflight.code);
// Evidence has its own durable archive/corpus jobs. Do not alter Catalog readiness
// or the full preprocessing job when publishing this one component.
try {
const outcome = await runEvidenceStage({ evidence: runtime.workspace.evidence,
consolidate: true, job: { runId: randomBytes(16).toString("hex"), completedStages: [], childRuns: {} },
}, {
runStage: async argv => {
const result = await this.deps.runChild({ argv, configPath: runtime.configLease.path });
const payload = JSON.parse(result.stdout);
if (result.exitCode !== 0 || payload.status !== "succeeded") {
return Promise.reject(new Error(typeof payload.error === "string" ? payload.error.slice(0, 1500) : "Evidence consolidation failed. Retry the command."));
}
return payload;
},
persistJob: () => undefined,
evidencePreflight: async () => preflight,
requireRunId: value => this.requireRunId(value),
numberRecord: value => this.numberRecord(value),
});
return baseResult(runtime, "evidence consolidate", "succeeded", "ok", {
...outcome, warnings: ["Evidence is active locally. Catalog and Schema readiness are unchanged. Commit and push the Evidence files manually."],
});
} catch (error) {
return baseResult(runtime, "evidence consolidate", "failed", "evidence_materialization_required", {
warnings: [error instanceof Error ? error.message : "Evidence consolidation failed. Retry the command."],
});
}
}
async inspect(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
try {
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
return baseResult(runtime, "inspect", "succeeded", "ok", {
artifactIdentities: [
{ kind: "descriptor", digest: runtime.descriptorBlob },
{ kind: "catalog", digest: runtime.catalogBlob },
{ kind: "runtime_config", digest: runtime.configLease.configDigest },
],
});
} catch {
return {
schemaVersion: 1,
status: "failed",
code: "workspace_not_activatable",
workspaceId: options.workspaceId,
workspaceRevision: "",
descriptorBlob: "",
operation: "inspect",
completedStages: [],
};
}
}
async run(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
if (!this.deps.catalogRepository) {
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
return baseResult(runtime, "preprocess run", "failed", "catalog_not_ready");
}
return await this.runFromCatalog(options);
}
async clear(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
const repository = this.deps.catalogRepository;
if (!repository) return baseResult(runtime, "preprocess clear", "failed", "catalog_not_ready");
const cleared = await repository.clearPreprocessing(runtime.workspaceId);
if (cleared.kind !== "cleared") {
return baseResult(
runtime,
"preprocess clear",
"failed",
cleared.kind === "already_running" ? "preprocessing_conflict" : "catalog_not_ready",
);
}
const result = await this.deps.runChild({
argv: ["preprocess", "clear", "--json", "-c", "/dev/fd/3"],
configPath: runtime.configLease.path,
});
if (result.exitCode !== 0) {
return baseResult(runtime, "preprocess clear", "failed", "preprocessing_clear_failed");
}
const payload = JSON.parse(result.stdout) as Record<string, unknown>;
return baseResult(runtime, "preprocess clear", "succeeded", "ok", {
counts: this.numberRecord(payload.counts),
});
}
private async runFromCatalog(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
const scope = await this.startRun(options.workspaceId);
const repository = this.deps.catalogRepository!;
const fingerprint = scope.runtime.configLease.inputFingerprint;
const started = await repository.beginPreprocessing(scope.runtime.workspaceId, fingerprint);
if (started.kind !== "started") {
scope.job.status = "failed";
scope.state.writeJob(scope.job);
return baseResult(
scope.runtime,
"preprocess run",
"failed",
started.kind === "already_running" ? "preprocessing_conflict" : "catalog_not_ready",
{ warnings: [`catalog=${started.kind}`] },
);
}
const revision = started.database.metadataContentRevision;
let finished = false;
let failureCode = "catalog_snapshot_failed";
try {
const snapshot = await buildCatalogMetadataSnapshot(
repository,
scope.runtime.workspaceId,
revision,
);
const snapshotPath = this.publishCatalogSnapshot(
scope.runtime.workspaceId,
JSON.stringify(snapshot),
);
this.completeStages(scope, "catalog_snapshot");
const semantic = await this.deps.semanticPreflight(scope.runtime.workspace);
if (!semantic.ok) {
await repository.finishPreprocessing(
scope.runtime.workspaceId,
revision,
fingerprint,
{ status: "failed", errorCode: semantic.code },
);
scope.job.status = "failed";
scope.state.writeJob(scope.job);
finished = true;
return baseResult(scope.runtime, "preprocess run", "failed", semantic.code);
}
failureCode = "schema_index_failed";
const payload = await this.runJsonStage(scope.runtime, [
"preprocess",
"catalog",
"--catalog-metadata",
snapshotPath,
"--json",
"-c",
"/dev/fd/3",
]);
this.completeStages(scope, "catalog_metadata", "lsh", "schema_index");
failureCode = "evidence_preprocessing_failed";
const outcome = await continueEvidencePreprocessing(
{
evidence: scope.runtime.workspace.evidence,
job: scope.job,
httpPrivateHostAllowlist: this.deps.httpPrivateHostAllowlist,
priorCounts: this.numberRecord(payload.counts),
},
this.evidenceDependencies(scope),
);
if (!["succeeded", "unchanged"].includes(outcome.status)) {
await repository.finishPreprocessing(
scope.runtime.workspaceId,
revision,
fingerprint,
{ status: "failed", errorCode: outcome.code },
);
scope.job.status = "failed";
scope.state.writeJob(scope.job);
finished = true;
return this.evidenceResult(scope, outcome);
}
// Evidence has been published. The only remaining operation is the atomic Catalog commit.
this.completeStages(scope, "evidence");
const completed = await repository.finishPreprocessing(
scope.runtime.workspaceId,
revision,
fingerprint,
{ status: "succeeded" },
);
if (!completed) throw new Error("catalog preprocessing completion lost its lease");
scope.job.status = "succeeded";
scope.state.writeJob(scope.job);
finished = true;
return this.evidenceResult(scope, outcome);
} catch (error) {
if (!finished) {
await repository.finishPreprocessing(
scope.runtime.workspaceId,
revision,
fingerprint,
{ status: "failed", errorCode: failureCode },
);
}
scope.job.status = "failed";
scope.state.writeJob(scope.job);
throw error;
}
}
private async startRun(workspaceId: string): Promise<RunScope> {
const runtime = await this.deps.acquireActiveRuntime(workspaceId);
const state = this.state(runtime.workspaceId);
const job = await state.beginJob({
operation: "preprocess run",
workspaceRevision: runtime.workspaceRevision,
descriptorBlob: runtime.descriptorBlob,
catalogBlob: runtime.catalogBlob,
configDigest: runtime.configLease.configDigest,
bindingDigest: runtime.configLease.bindingDigest,
embeddingId: runtime.configLease.effectiveConfig.embedding.id,
embeddingDimensions: runtime.configLease.effectiveConfig.embedding.dimensions,
});
return { runtime, state, job };
}
private state(workspaceId: string): PreprocessingStateStore {
return new PreprocessingStateStore({ dataRoot: this.deps.dataRoot, workspaceId });
}
private completeStages(scope: RunScope, ...stages: string[]): void {
for (const stage of stages) {
if (!scope.job.completedStages.includes(stage)) scope.job.completedStages.push(stage);
}
scope.state.writeJob(scope.job);
}
private publishCatalogSnapshot(workspaceId: string, contents: string): string {
const root = join(this.deps.dataRoot, "sessions", workspaceId, "preprocessing");
mkdirSync(root, { recursive: true, mode: 0o700 });
const target = join(root, "catalog-metadata.json");
const staging = join(root, `.catalog-metadata-${randomBytes(6).toString("hex")}.json`);
try {
writeFileSync(staging, contents, { encoding: "utf8", flag: "wx", mode: 0o600 });
renameSync(staging, target);
return target;
} catch (error) {
rmSync(staging, { force: true });
throw error;
}
}
private evidenceDependencies(scope: RunScope): EvidencePreprocessingDependencies {
return {
runStage: async (argv) => await this.runJsonStage(scope.runtime, argv),
persistJob: () => this.state(scope.runtime.workspaceId).writeJob(scope.job),
evidencePreflight: async () => await this.deps.evidencePreflight(scope.runtime.workspace),
requireRunId: (value) => this.requireRunId(value),
numberRecord: (value) => this.numberRecord(value),
};
}
private evidenceResult(
scope: RunScope,
outcome: EvidencePreprocessingOutcome,
): WorkspaceOperationResult {
const { status, code, ...extra } = outcome;
return baseResult(scope.runtime, "preprocess run", status, code, extra);
}
private async runJsonStage(runtime: ActiveRuntime, argv: string[]): Promise<Record<string, unknown>> {
const result = await this.deps.runChild({ argv, configPath: runtime.configLease.path });
if (result.exitCode !== 0) throw new Error("workspace child failed");
return JSON.parse(result.stdout) as Record<string, unknown>;
}
private requireRunId(value: unknown): string {
if (typeof value !== "string" || !/^[0-9a-f]{32}$/.test(value)) throw new Error("child run id is invalid");
return value;
}
private numberRecord(value: unknown): Record<string, number> | undefined {
if (!value || typeof value !== "object" || Array.isArray(value)) return undefined;
return Object.fromEntries(Object.entries(value as Record<string, unknown>).map(([key, nested]) => [key, Number(nested)]));
}
}