feat: complete catalog-driven preprocessing
Publish documentation / publish (push) Successful in 2m12s

This commit is contained in:
Codex
2026-09-06 17:49:35 +02:00
parent 8707ae1d46
commit cffa60772e
141 changed files with 5898 additions and 3015 deletions
@@ -14,7 +14,7 @@ const consolidationSchema = z.discriminatedUnion("target", [
target: z.enum(["tables", "columns"]),
targetIds: z.array(idSchema).min(1).max(10_000),
}).strict(),
z.object({ target: z.literal("database_columns") }).strict(),
z.object({ target: z.enum(["database", "database_columns"]) }).strict(),
]);
function manage(request: FastifyRequest, reply: FastifyReply) {
+63 -8
View File
@@ -11,8 +11,8 @@ import type { WorkspaceRegistry } from "../workspaces/registry.js";
import { validateOperationalWorkspace, type WorkspaceDescriptor } from "../workspaces/schema.js";
import type { MaintenanceBarrier } from "../runtime/maintenance-gate.js";
import { hasPermission, isPrincipalContext, requirePermission } from "../auth/authorization.js";
import type { EffectiveRelationshipSnapshotProvider } from "../catalog/effective-relationship-snapshot.js";
import { splitCanonicalModelId, type RuntimeModelCatalog } from "../models/runtime-model-catalog.js";
import type { CatalogRepository } from "../catalog/types.js";
const BOOTSTRAP_FAILURE_MESSAGE =
"Session startup failed. Check configuration and connectivity, then Resume the session.";
@@ -42,12 +42,40 @@ export function sessionRoutes(
/** Fail-closed installation/runtime transport capability check. */
workspaceRuntimeSupport: (workspace: WorkspaceDescriptor) => boolean;
maintenanceBarrier: MaintenanceBarrier;
/** Optional only for narrow route-test stubs and installations without a Catalog database. */
effectiveRelationships?: EffectiveRelationshipSnapshotProvider;
modelCatalog: RuntimeModelCatalog;
/** PostgreSQL authority for mandatory preprocessing admission. */
catalogRepository?: CatalogRepository;
},
) {
const lifecycleTails = new Map<string, Promise<void>>();
const preprocessingIsCurrent = async (
workspaceId: string,
inputFingerprint?: string,
): Promise<boolean> => {
if (!d.catalogRepository) return true;
const database = await d.catalogRepository.getByWorkspace(workspaceId);
return database !== undefined
&& database.preprocessingStatus === "succeeded"
&& database.preprocessedMetadataRevision === database.metadataContentRevision
&& (inputFingerprint === undefined
|| database.preprocessingInputFingerprint === inputFingerprint);
};
const catalogTransportSupportsSessionRuntime = async (workspaceId: string): Promise<boolean> => {
if (!d.catalogRepository) return true;
const database = await d.catalogRepository.getByWorkspace(workspaceId);
return database === undefined || database.binding.transport !== "ssh_tunnel";
};
const currentInputFingerprint = async (
runner: any,
workspaceConfigPath: string,
): Promise<string | undefined> => (
typeof runner.workspaceInputFingerprint === "function"
? await runner.workspaceInputFingerprint(workspaceConfigPath)
: undefined
);
const boundRuntimes = new Map<
string,
ReturnType<PiProcessManager["createFor"]>
@@ -92,16 +120,13 @@ export function sessionRoutes(
const optionsWithRuntimeConfig = async (
runner: any,
workspaceConfigPath: string | undefined,
workspaceId: string | undefined,
_workspaceId: string | undefined,
options: any,
) => {
if (!workspaceConfigPath || typeof runner.acquireWorkspaceRuntime !== "function") return options;
const effectiveRelationships = workspaceId && d.effectiveRelationships
? await d.effectiveRelationships.render(workspaceId)
: undefined;
return {
...options,
runtimeConfig: runner.acquireWorkspaceRuntime(workspaceConfigPath, effectiveRelationships),
runtimeConfig: await runner.acquireWorkspaceRuntime(workspaceConfigPath),
};
};
@@ -379,6 +404,19 @@ export function sessionRoutes(
workspaceId = resolved.revision.id;
workspaceRevision = resolved.revision.commit;
workspaceDescriptor = resolved.workspace;
if (!await catalogTransportSupportsSessionRuntime(workspaceId)) {
return reply.code(409).send({
error: "This workspace transport is not available to runtime sessions.",
code: "workspace_not_activatable",
});
}
const fingerprint = await currentInputFingerprint(runner, workspaceConfigPath);
if (!await preprocessingIsCurrent(workspaceId, fingerprint)) {
return reply.code(409).send({
error: "Run workspace preprocessing before starting the core.",
code: "preprocessing_required",
});
}
} catch {
return reply.code(409).send({
error: WORKSPACE_REVISION_UNAVAILABLE_MESSAGE,
@@ -615,6 +653,23 @@ export function sessionRoutes(
workspaceDescriptor = resolved.workspace;
}
catch { return unavailableWorkspaceReply(reply); }
if (d.catalogRepository && saved.workspace_id
&& !await catalogTransportSupportsSessionRuntime(saved.workspace_id)) {
return reply.code(409).send({
error: "This workspace transport is not available to runtime sessions.",
code: "workspace_not_activatable",
});
}
const fingerprint = d.catalogRepository
? await currentInputFingerprint(runner, workspaceConfigPath)
: undefined;
if (d.catalogRepository
&& (!saved.workspace_id || !await preprocessingIsCurrent(saved.workspace_id, fingerprint))) {
return reply.code(409).send({
error: "Run workspace preprocessing before starting the core.",
code: "preprocessing_required",
});
}
try { settings = await d.getSettings(principal); } catch { return storageFailure(reply); }
// This check belongs inside the per-session lock: a preceding cold Resume may have
// installed a running runtime while this request was waiting.
@@ -0,0 +1,402 @@
import type { FastifyInstance, FastifyReply, FastifyRequest } from "fastify";
import { z } from "zod";
import { isPrincipalContext, requirePermission } from "../auth/authorization.js";
import {
CatalogUnavailableError,
type CatalogRepository,
type WorkspaceDatabase,
} from "../catalog/types.js";
import type { ThtRunner } from "../tht/tht-runner.js";
import type { WorkspacePreprocessingService } from "../workspaces/preprocessing-service.js";
import type { PreprocessingJobState } from "../workspaces/preprocessing-state.js";
import type { WorkspaceRegistry } from "../workspaces/registry.js";
const workspaceIdSchema = z.string().regex(/^[a-z][a-z0-9-]{2,62}$/);
const ACTIVE_SYNC_STATES = new Set(["queued", "running", "awaiting_confirmation", "applying"]);
export type WorkspacePreprocessingUiState =
| "ready"
| "required"
| "running"
| "blocked"
| "failed";
export interface WorkspacePreprocessingStatus {
schemaVersion: 1;
workspaceId: string;
state: WorkspacePreprocessingUiState;
actionable: boolean;
clearable: boolean;
detail: string;
reason?: string;
nextStep?: string;
metadataRevision?: number;
preprocessedMetadataRevision?: number;
startedAt?: string;
finishedAt?: string;
progress?: {
stage: "catalog_snapshot" | "schema_index" | "evidence" | "finalizing";
step: number;
totalSteps: 4;
};
lastFailure?: {
stage: string;
errorCode: string;
finishedAt: string;
};
}
export interface WorkspacePreprocessingRouteDeps {
repository: CatalogRepository;
registry: WorkspaceRegistry;
service: Pick<WorkspacePreprocessingService, "run" | "clear">;
inputFingerprint?: Pick<ThtRunner, "workspaceInputFingerprint">;
readLatestJob?: (workspaceId: string) => PreprocessingJobState | undefined;
}
const PROGRESS_DETAILS = {
catalog_snapshot: "Preparing the PostgreSQL Catalog snapshot.",
schema_index: "Building schema vectors and LSH indexes.",
evidence: "Indexing Evidence.",
finalizing: "Publishing the completed preprocessing state.",
} as const;
function runningProgress(
workspaceId: string,
deps: WorkspacePreprocessingRouteDeps,
): NonNullable<WorkspacePreprocessingStatus["progress"]> {
let job: PreprocessingJobState | undefined;
try {
job = deps.readLatestJob?.(workspaceId);
} catch {
// Progress is supplemental. A damaged or temporarily unavailable checkpoint must not hide
// the authoritative running state held by PostgreSQL.
}
const completed = new Set(job?.status === "active" ? job.completedStages : []);
if (completed.has("evidence")) return { stage: "finalizing", step: 4, totalSteps: 4 };
if (completed.has("schema_index")) return { stage: "evidence", step: 3, totalSteps: 4 };
if (completed.has("catalog_snapshot")) return { stage: "schema_index", step: 2, totalSteps: 4 };
return { stage: "catalog_snapshot", step: 1, totalSteps: 4 };
}
function base(
workspaceId: string,
state: WorkspacePreprocessingUiState,
detail: string,
database?: WorkspaceDatabase,
): WorkspacePreprocessingStatus {
return {
schemaVersion: 1,
workspaceId,
state,
actionable: state === "ready" || state === "required" || state === "failed",
clearable: Boolean(
database
&& state !== "running"
&& database.preprocessingErrorCode !== "derived_data_cleared"
),
detail,
...(database ? {
metadataRevision: database.metadataContentRevision,
...(database.preprocessedMetadataRevision === undefined
? {}
: { preprocessedMetadataRevision: database.preprocessedMetadataRevision }),
...(database.preprocessingStartedAt ? { startedAt: database.preprocessingStartedAt } : {}),
...(database.preprocessingFinishedAt ? { finishedAt: database.preprocessingFinishedAt } : {}),
} : {}),
};
}
function blocked(
workspaceId: string,
detail: string,
reason: string,
nextStep: string,
database?: WorkspaceDatabase,
): WorkspacePreprocessingStatus {
return { ...base(workspaceId, "blocked", detail, database), reason, nextStep };
}
function failureDiagnostic(errorCode: string): {
stage: string;
detail: string;
reason: string;
nextStep: string;
} {
switch (errorCode) {
case "catalog_snapshot_failed":
return {
stage: "catalog_snapshot",
detail: "The Catalog snapshot could not be prepared.",
reason: "PostgreSQL Catalog metadata could not be read into a consistent preprocessing snapshot.",
nextStep: "Check Catalog availability, then retry preprocessing.",
};
case "schema_index_failed":
return {
stage: "schema_index",
detail: "The schema index could not be rebuilt.",
reason: "The schema indexing worker stopped before the Catalog snapshot was published to Qdrant.",
nextStep: "Open Last run details below, check the core service log for this error code, then retry.",
};
case "evidence_preprocessing_failed":
return {
stage: "evidence",
detail: "Evidence preprocessing did not complete.",
reason: "The Evidence indexing worker stopped before it finished publishing the current workspace data.",
nextStep: "Open Last run details below, check the core service log for this error code, then retry.",
};
case "semantic_index_incompatible":
return {
stage: "semantic_preflight",
detail: "The semantic index configuration is incompatible.",
reason: "Qdrant rejected the collection configuration for the active embedding model.",
nextStep: "Check embedding dimensions and Qdrant collection settings, then retry.",
};
case "egress_policy_refused":
return {
stage: "evidence",
detail: "The Evidence source was refused by policy.",
reason: "The configured Evidence endpoint is not allowed by the installation egress policy.",
nextStep: "Correct the Evidence source or its allowlist configuration, then retry.",
};
case "workspace_not_activatable":
return {
stage: "runtime_preflight",
detail: "The workspace runtime could not be prepared.",
reason: "The active workspace or one of its required runtime bindings is not usable.",
nextStep: "Check Workspace management and Database management, then retry.",
};
default:
return {
stage: "preprocessing",
detail: "Preprocessing did not complete.",
reason: "The preprocessing worker stopped before the current Catalog revision was published.",
nextStep: "Open Last run details below, check the core service log for this error code, then retry.",
};
}
}
export async function readWorkspacePreprocessingStatus(
workspaceId: string,
deps: WorkspacePreprocessingRouteDeps,
): Promise<WorkspacePreprocessingStatus> {
const database = await deps.repository.getByWorkspace(workspaceId);
if (!database) {
return blocked(
workspaceId,
"Database configuration is required.",
"This workspace has no database configuration in the PostgreSQL Catalog.",
"Open Database management and configure the workspace database.",
);
}
if (database.preprocessingStatus === "running") {
const progress = runningProgress(workspaceId, deps);
return {
...base(workspaceId, "running", PROGRESS_DETAILS[progress.stage], database),
progress,
};
}
if (database.binding.transport === "ssh_tunnel") {
return blocked(
workspaceId,
"The database transport is not supported by the core runtime.",
"The current database binding uses an SSH tunnel, which cannot be used by a ThothII session.",
"Open Database management and select a supported runtime transport.",
database,
);
}
if (database.schemaSyncedVersion !== database.version) {
return blocked(
workspaceId,
"Catalog synchronization is required.",
`Database configuration v${database.version} is newer than the latest Catalog synchronization${database.schemaSyncedVersion === undefined ? "." : ` v${database.schemaSyncedVersion}.`}`,
"Open Database management and run Synchronize schema.",
database,
);
}
const [syncRuns, activeDescriptionRun, sensitivityRuns] = await Promise.all([
deps.repository.listSyncRuns(database.id, 10),
deps.repository.getActiveDescriptionGenerationRun(),
deps.repository.listSensitivityAnalysisRuns(50),
]);
if (syncRuns.some((run) => ACTIVE_SYNC_STATES.has(run.state))) {
return blocked(
workspaceId,
"Catalog synchronization is in progress.",
"The Catalog is being synchronized and its metadata revision is not stable yet.",
"Wait for schema synchronization to finish, then run preprocessing.",
database,
);
}
if (activeDescriptionRun?.databaseId === database.id
&& ["queued", "running"].includes(activeDescriptionRun.status)) {
return blocked(
workspaceId,
"Description generation is in progress.",
"Catalog descriptions are still being generated for this database.",
"Wait for description generation to finish, then run preprocessing.",
database,
);
}
if (sensitivityRuns.some((run) => run.databaseId === database.id && run.status === "running")) {
return blocked(
workspaceId,
"Sensitivity analysis is in progress.",
"Catalog sensitivity metadata is still being analyzed for this database.",
"Wait for sensitivity analysis to finish, then run preprocessing.",
database,
);
}
let inputFingerprint: string | undefined;
try {
const record = await deps.registry.read(workspaceId);
if (deps.inputFingerprint) {
inputFingerprint = await deps.inputFingerprint.workspaceInputFingerprint(
record.revision.snapshotPath,
);
}
} catch {
return blocked(
workspaceId,
"The workspace runtime configuration is unavailable.",
"The active workspace revision or one of its required database secrets could not be resolved.",
"Check Workspace management and Database management before running preprocessing.",
database,
);
}
const current = database.preprocessingStatus === "succeeded"
&& database.preprocessedMetadataRevision === database.metadataContentRevision
&& (inputFingerprint === undefined
|| database.preprocessingInputFingerprint === inputFingerprint);
if (current) {
return base(
workspaceId,
"ready",
`Catalog revision ${database.metadataContentRevision} is indexed.`,
database,
);
}
if (database.preprocessingErrorCode === "derived_data_cleared") {
return base(
workspaceId,
"required",
"Reference vectors and LSH are empty. Memory is preserved.",
database,
);
}
if (database.preprocessingStatus === "failed"
&& database.preprocessingErrorCode
&& database.preprocessingErrorCode !== "catalog_changed"
&& database.preprocessingErrorCode !== "derived_data_cleared"
&& database.preprocessingFinishedAt) {
const diagnostic = failureDiagnostic(database.preprocessingErrorCode);
return {
...base(workspaceId, "failed", diagnostic.detail, database),
reason: diagnostic.reason,
nextStep: diagnostic.nextStep,
lastFailure: {
stage: diagnostic.stage,
errorCode: database.preprocessingErrorCode,
finishedAt: database.preprocessingFinishedAt,
},
};
}
return base(
workspaceId,
"required",
`Catalog revision ${database.metadataContentRevision} is not indexed.`,
database,
);
}
function safeError(reply: FastifyReply, error: unknown) {
if (error instanceof CatalogUnavailableError) {
return reply.code(503).send({
code: "catalog_unavailable",
message: "Database catalog is unavailable.",
});
}
if (error instanceof z.ZodError) {
return reply.code(400).send({
code: "preprocessing_request_invalid",
message: "Preprocessing request is invalid.",
});
}
return reply.code(500).send({
code: "preprocessing_run_failed",
message: "Preprocessing status could not be resolved.",
});
}
function workspaceIdFrom(request: FastifyRequest): string {
return workspaceIdSchema.parse((request.params as { workspaceId?: unknown }).workspaceId);
}
export function workspacePreprocessingRoutes(
app: FastifyInstance,
deps: WorkspacePreprocessingRouteDeps,
): void {
app.get("/workspaces/:workspaceId/preprocessing", async (request, reply) => {
if (!isPrincipalContext(requirePermission(request, reply, "session.use"))) return reply;
try {
return await readWorkspacePreprocessingStatus(workspaceIdFrom(request), deps);
} catch (error) {
return safeError(reply, error);
}
});
app.post("/workspaces/:workspaceId/preprocessing", async (request, reply) => {
if (!isPrincipalContext(requirePermission(request, reply, "database.manage"))) return reply;
try {
const workspaceId = workspaceIdFrom(request);
const before = await readWorkspacePreprocessingStatus(workspaceId, deps);
if (!before.actionable) {
return reply.code(409).send({ ...before, code: "preprocessing_blocked" });
}
const result = await deps.service.run({ workspaceId });
const after = await readWorkspacePreprocessingStatus(workspaceId, deps);
if (after.state === "ready") return after;
if (after.state === "failed") {
return reply.code(422).send({ ...after, code: "preprocessing_run_failed" });
}
if (after.state === "blocked" || after.state === "running") {
return reply.code(409).send({ ...after, code: "preprocessing_blocked" });
}
return reply.code(500).send({
code: "preprocessing_run_failed",
message: `Preprocessing ended with ${result.code}.`,
});
} catch (error) {
return safeError(reply, error);
}
});
app.delete("/workspaces/:workspaceId/preprocessing", async (request, reply) => {
if (!isPrincipalContext(requirePermission(request, reply, "database.manage"))) return reply;
try {
const workspaceId = workspaceIdFrom(request);
const before = await readWorkspacePreprocessingStatus(workspaceId, deps);
if (!before.clearable) {
return reply.code(409).send({ ...before, code: "preprocessing_clear_blocked" });
}
const result = await deps.service.clear({ workspaceId });
if (result.status !== "succeeded") {
return reply.code(result.code === "preprocessing_conflict" ? 409 : 500).send({
code: result.code,
message: "Preprocessing data could not be cleared.",
});
}
return await readWorkspacePreprocessingStatus(workspaceId, deps);
} catch (error) {
return safeError(reply, error);
}
});
}