Files
ThothII/backend/src/workspaces/diagnostics.ts
T

561 lines
20 KiB
TypeScript

import { readFile } from "node:fs/promises";
import { Client } from "pg";
import { MAX_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS } from "../config.js";
import { buildInstallationContract } from "./contracts.js";
import type { RuntimeBindings } from "./runtime-renderer.js";
import {
resolveDiagnosticUrl,
validateWorkspaceDescriptor,
type RestDiagnosticRequest,
type WorkspaceDescriptor,
} from "./schema.js";
import type { WorkspaceErrorCode } from "./types.js";
import type { SemanticRuntimeConfig } from "./runtime-renderer.js";
import type { AuthDiagnostics } from "../auth/diagnostics.js";
export interface Diagnostic {
level: "error" | "warning" | "info";
code: WorkspaceErrorCode | "binding_ok";
field?: string;
variable?: string;
message: string;
}
/** Connector-only result produced before the route aggregates authentication. */
export interface ConnectorDiagnostics {
activatable: boolean;
diagnostics: Diagnostic[];
}
/** The HTTP workspace diagnostic contract always includes the shared authentication report. */
export interface WorkspaceDiagnostics extends ConnectorDiagnostics {
authentication: AuthDiagnostics;
}
interface DiagnosticResource {
database?: string;
schema?: string;
}
type RestConnectorDiagnostic = RestDiagnosticRequest & {
response?: Record<string, string>;
};
export interface ConnectorDiagnosticRequest {
role: "dwh";
transport: "postgres_direct" | "rest_api";
host?: string;
port?: number;
baseUrl?: string;
user?: string;
credentialFile?: string;
tlsCaFile?: string;
tlsServername?: string;
resource: DiagnosticResource;
timeoutMs: number;
signal: AbortSignal;
diagnostic?: RestConnectorDiagnostic;
}
export interface ConnectorDiagnosticResult {
resolved: boolean;
tlsVerified: boolean;
authenticated: boolean;
resource: DiagnosticResource;
}
export interface QdrantDiagnosticRequest {
baseUrl: string;
collection: string;
timeoutMs: number;
signal: AbortSignal;
}
export interface QdrantDiagnosticResult {
collection: string;
dimensions?: number;
distance?: string;
}
export interface EmbeddingDiagnosticRequest {
baseUrl: string;
model: string;
timeoutMs: number;
signal: AbortSignal;
}
export interface EmbeddingDiagnosticResult {
available: boolean;
dimensions?: number;
}
export interface DirectProtocolFactory {
probe(request: ConnectorDiagnosticRequest): Promise<ConnectorDiagnosticResult>;
}
export interface DatabaseDiagnosticClient {
query(sql: string, values: readonly unknown[]): Promise<{ rows: Array<Record<string, unknown>> }>;
end(): Promise<void>;
}
export interface DatabaseDiagnosticClientFactory {
connect(request: {
host: string;
port: number;
database: string;
user: string;
credentialFile: string;
tlsCaFile?: string;
tlsServername?: string;
signal: AbortSignal;
}): Promise<DatabaseDiagnosticClient>;
}
export interface ConcreteDiagnosticAdapterDependencies {
directProtocol?: DirectProtocolFactory;
databaseClient?: DatabaseDiagnosticClientFactory;
}
/** Adapters retain only diagnostic metadata and never return credential contents or bodies. */
export interface DiagnosticAdapters {
probeConnector(request: ConnectorDiagnosticRequest): Promise<ConnectorDiagnosticResult>;
inspectQdrant(request: QdrantDiagnosticRequest): Promise<QdrantDiagnosticResult>;
probeEmbedding(request: EmbeddingDiagnosticRequest): Promise<EmbeddingDiagnosticResult>;
}
export const DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS = 5_000;
async function secretPresent(file: string): Promise<boolean> {
return (await readFile(file, "utf8")).trim().length > 0;
}
async function restHeaders(
diagnostic: RestDiagnosticRequest,
credentialFile: string | undefined,
): Promise<Record<string, string>> {
if (diagnostic.auth === "none") return {};
if (!credentialFile || !(await secretPresent(credentialFile))) throw new Error("REST probe failed");
const secret = (await readFile(credentialFile, "utf8")).trim();
return diagnostic.auth === "bearer" ? { authorization: `Bearer ${secret}` } : { "x-api-key": secret };
}
/**
* Concrete production adapters deliberately retain only probe metadata. Protocol failures and
* response bodies are discarded at this boundary; callers receive fixed diagnostics instead.
*/
export function createConcreteDiagnosticAdapters(
dependencies: ConcreteDiagnosticAdapterDependencies = {},
): DiagnosticAdapters {
const databaseClient = dependencies.databaseClient ?? {
async connect(request: {
host: string;
port: number;
database: string;
user: string;
credentialFile: string;
tlsCaFile?: string;
tlsServername?: string;
signal: AbortSignal;
}) {
const client = new Client({
host: request.host,
port: request.port,
database: request.database,
user: request.user,
password: (await readFile(request.credentialFile, "utf8")).trim(),
ssl: {
...(request.tlsCaFile ? { ca: await readFile(request.tlsCaFile, "utf8") } : {}),
...(request.tlsServername ? { servername: request.tlsServername } : {}),
rejectUnauthorized: true,
},
connectionTimeoutMillis: 5_000,
});
const abort = () => { void client.end(); };
request.signal.addEventListener("abort", abort, { once: true });
try {
await client.connect();
return {
query: async (sql: string, values: readonly unknown[]) => await client.query(sql, [...values]),
end: async () => {
request.signal.removeEventListener("abort", abort);
await client.end();
},
};
} catch (error) {
request.signal.removeEventListener("abort", abort);
await client.end().catch(() => undefined);
throw error;
}
},
};
const directProtocol = dependencies.directProtocol ?? {
async probe(request: ConnectorDiagnosticRequest): Promise<ConnectorDiagnosticResult> {
if (!request.host || !request.port || !request.user || !request.credentialFile
|| !(await secretPresent(request.credentialFile))) {
throw new Error("direct probe failed");
}
const database = request.resource.database;
const schema = request.resource.schema;
if (!database || !schema) throw new Error("direct probe failed");
const client = await databaseClient.connect({
host: request.host,
port: request.port,
database,
user: request.user,
credentialFile: request.credentialFile,
tlsCaFile: request.tlsCaFile,
tlsServername: request.tlsServername,
signal: request.signal,
});
try {
const result = await client.query(
"SELECT current_database() AS database, current_schema() AS schema",
[],
);
const row = result.rows[0];
if (row?.database !== database || row.schema !== schema) throw new Error("direct probe failed");
return {
resolved: true,
tlsVerified: true,
authenticated: true,
resource: request.resource,
};
} finally {
await client.end().catch(() => undefined);
}
},
};
return {
async probeConnector(request) {
if (request.transport === "rest_api") {
if (!request.baseUrl || !request.diagnostic || request.tlsCaFile) {
throw new Error("REST probe failed");
}
const endpoint = resolveDiagnosticUrl(request.baseUrl, request.diagnostic.path);
const response = await fetch(endpoint.toString(), {
method: request.diagnostic.method,
headers: await restHeaders(request.diagnostic, request.credentialFile),
signal: request.signal,
redirect: "error",
});
if (!response.ok) throw new Error("REST probe failed");
if ("response" in request.diagnostic) {
const payload = await response.json().catch(() => undefined) as Record<string, unknown> | undefined;
const declared = request.diagnostic.response;
// Validate a declared field only when the probe response actually carries it, so a
// health-style ping (2xx + JSON without database/schema identity) still proves a
// reachable, authenticated connector. Declared fields that are present must match.
const databaseMatches = declared?.database === undefined
|| payload?.[declared.database] === undefined
|| payload[declared.database] === request.resource.database;
const schemaMatches = declared?.schema === undefined
|| payload?.[declared.schema] === undefined
|| payload[declared.schema] === request.resource.schema;
if (!payload || !databaseMatches || !schemaMatches) {
throw new Error("REST probe failed");
}
}
return {
resolved: true,
tlsVerified: new URL(request.baseUrl).protocol === "https:",
authenticated: true,
resource: request.resource,
};
}
if (request.transport !== "postgres_direct") throw new Error("direct probe failed");
return await directProtocol.probe(request);
},
async inspectQdrant(request) {
const response = await fetch(
new URL(`/collections/${request.collection}`, `${request.baseUrl}/`).toString(),
{ method: "GET", signal: request.signal, redirect: "error" },
);
const payload = await response.json().catch(() => undefined) as {
result?: { config?: { params?: { vectors?: { size?: unknown; distance?: unknown } } } };
} | undefined;
const size = payload?.result?.config?.params?.vectors?.size;
const distance = payload?.result?.config?.params?.vectors?.distance;
if (!response.ok || !Number.isInteger(size) || typeof distance !== "string"
|| distance.length === 0) {
throw new Error("Qdrant metadata probe failed");
}
return {
collection: request.collection,
dimensions: size as number,
distance: distance.toLowerCase(),
};
},
async probeEmbedding(request) {
const response = await fetch(new URL("/api/embed", `${request.baseUrl}/`).toString(), {
method: "POST",
headers: { "content-type": "application/json" },
body: JSON.stringify({ model: request.model, input: "diagnostic" }),
signal: request.signal,
redirect: "error",
});
const payload = await response.json().catch(() => undefined) as {
embeddings?: unknown[];
} | undefined;
const vector = Array.isArray(payload?.embeddings) ? payload.embeddings[0] : undefined;
if (!response.ok || !Array.isArray(vector)) throw new Error("embedding probe failed");
return { available: true, dimensions: vector.length };
},
};
}
function configuredTimeout(value: number | undefined): number {
return Math.min(
Math.max(1, value ?? DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS),
MAX_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS,
);
}
function boundedTimeout(value: number | undefined, ceiling: number): number {
return Math.min(Math.max(1, value ?? ceiling), ceiling, MAX_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS);
}
async function withTimeout<T>(timeoutMs: number, operation: (signal: AbortSignal) => Promise<T>): Promise<T> {
const controller = new AbortController();
let timer: NodeJS.Timeout | undefined;
try {
return await new Promise<T>((resolve, reject) => {
timer = setTimeout(() => {
controller.abort();
reject(new Error("diagnostic timed out"));
}, timeoutMs);
void operation(controller.signal).then(resolve, reject);
});
} finally {
if (timer !== undefined) clearTimeout(timer);
controller.abort();
}
}
function sameResource(expected: DiagnosticResource, actual: DiagnosticResource): boolean {
return Object.entries(expected).every(([key, value]) => actual[key as keyof DiagnosticResource] === value);
}
function hasRequiredConnectorChecks(
result: ConnectorDiagnosticResult,
resource: DiagnosticResource,
): boolean {
return result.resolved
&& result.tlsVerified
&& result.authenticated
&& sameResource(resource, result.resource);
}
function diagnosticError(code: WorkspaceErrorCode, field?: string): Diagnostic {
return {
level: "error",
code,
...(field ? { field } : {}),
message: code === "binding_missing"
? "Installation binding is missing or invalid."
: code === "semantic_index_incompatible"
? "Semantic index metadata is incompatible with this workspace."
: code === "workspace_not_activatable"
? "This transport can be tested, but it is not available to runtime sessions."
: "Connector diagnostic failed.",
};
}
function bindingName(workspace: WorkspaceDescriptor, suffix: string): string {
const entry = buildInstallationContract(workspace).variables.find((variable) => (
variable.role === "DWH" && variable.suffix === suffix
));
if (!entry) throw new Error(`workspace contract is missing DWH_${suffix}`);
return entry.name;
}
function numericBinding(binding: Record<string, string>, name: string): number | undefined {
const value = Number(binding[name]);
return Number.isInteger(value) && value > 0 && value <= 65_535 ? value : undefined;
}
function requireSupportedDescriptor(workspace: unknown): void {
if (typeof workspace !== "object" || workspace === null) {
throw new Error("Workspace diagnoser supports only workspace schema version 3");
}
const metadata = Reflect.get(workspace, "workspace");
if (typeof metadata !== "object" || metadata === null
|| Reflect.get(metadata, "schema_version") !== 3) {
throw new Error("Workspace diagnoser supports only workspace schema version 3");
}
}
async function diagnoseValidatedWorkspace(
descriptor: WorkspaceDescriptor,
bindings: RuntimeBindings,
adapters: DiagnosticAdapters,
timeoutMs: number,
semanticRuntime: SemanticRuntimeConfig,
): Promise<ConnectorDiagnostics> {
const evidenceField = descriptor.evidence?.source.type === "http"
? "evidence.source.authentication"
: descriptor.evidence?.source.type === "s3"
? "evidence.source.credentials"
: undefined;
const evidenceDiagnostics = [...bindings.evidence.missing].sort().map((variable): Diagnostic => ({
...diagnosticError("binding_missing", evidenceField),
variable,
}));
const diagnostics = [
...[...bindings.dwh.missing].sort().map((field) => diagnosticError("binding_missing", field)),
...evidenceDiagnostics,
];
if (diagnostics.length > 0) return { activatable: false, diagnostics };
const dwhTimeout = boundedTimeout(descriptor.dwh.timeout_ms, timeoutMs);
let activatable = true;
const dwhValues = bindings.dwh.values;
const dwhField = (suffix: string) => bindingName(descriptor, suffix);
const dwhResource = { database: descriptor.dwh.database, schema: descriptor.dwh.schema };
let dwhRequest: ConnectorDiagnosticRequest | undefined;
if (bindings.dwh.transport === "rest_api") {
const diagnostic = descriptor.diagnostics?.dwh_rest;
const baseUrl = dwhValues[dwhField("BASE_URL")];
if (diagnostic && baseUrl) {
const credentialFile = diagnostic.auth === "none" ? undefined : dwhValues[dwhField("API_KEY_FILE")];
if (diagnostic.auth === "none" || credentialFile !== undefined) {
dwhRequest = {
role: "dwh",
transport: "rest_api",
baseUrl,
credentialFile,
tlsCaFile: dwhValues[dwhField("TLS_CA_FILE")],
resource: dwhResource,
timeoutMs: dwhTimeout,
signal: new AbortController().signal,
diagnostic,
};
}
}
} else if (bindings.dwh.transport === "postgres_direct") {
const host = dwhValues[dwhField("HOST")];
const port = numericBinding(dwhValues, dwhField("PORT"));
const user = dwhValues[dwhField("USER")];
const credentialFile = dwhValues[dwhField("PASSWORD_FILE")];
if (host && port && user && credentialFile) {
dwhRequest = {
role: "dwh",
transport: "postgres_direct",
host,
port,
user,
credentialFile,
tlsCaFile: dwhValues[dwhField("TLS_CA_FILE")],
resource: dwhResource,
timeoutMs: dwhTimeout,
signal: new AbortController().signal,
};
}
}
if (!dwhRequest) {
diagnostics.push(diagnosticError("workspace_not_activatable"));
return { activatable: false, diagnostics };
}
try {
const dwhResult = await withTimeout(dwhTimeout, (signal) => adapters.probeConnector({
...dwhRequest,
signal,
timeoutMs: dwhTimeout,
}));
if (!hasRequiredConnectorChecks(dwhResult, dwhRequest.resource)) {
diagnostics.push(diagnosticError("connector_unavailable"));
activatable = false;
}
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
activatable = false;
}
try {
const vector = await withTimeout(timeoutMs, (signal) => adapters.inspectQdrant({
baseUrl: semanticRuntime.internalQdrantUrl,
collection: descriptor.semantic_index.vector_store.collection,
timeoutMs,
signal,
}));
const expected = descriptor.semantic_index.vector_store;
if (vector.collection !== expected.collection
|| vector.dimensions !== expected.dimensions
|| vector.distance !== expected.distance) {
diagnostics.push(diagnosticError("semantic_index_incompatible"));
activatable = false;
}
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
activatable = false;
}
try {
const embedding = await withTimeout(timeoutMs, (signal) => adapters.probeEmbedding({
baseUrl: semanticRuntime.internalEmbeddingUrl,
model: semanticRuntime.internalEmbeddingModel,
timeoutMs,
signal,
}));
if (semanticRuntime.internalEmbeddingModel !== descriptor.semantic_index.embedding.model
|| semanticRuntime.internalEmbeddingDimensions !== descriptor.semantic_index.embedding.dimensions
|| !embedding.available
|| embedding.dimensions !== descriptor.semantic_index.embedding.dimensions) {
diagnostics.push(diagnosticError("semantic_index_incompatible"));
activatable = false;
}
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
activatable = false;
}
return {
activatable,
diagnostics: diagnostics.length > 0
? diagnostics
: [{
level: "info",
code: "binding_ok",
message: "Installation bindings and diagnostics succeeded.",
}],
};
}
const DEFAULT_SEMANTIC_RUNTIME: SemanticRuntimeConfig = {
internalQdrantUrl: "http://qdrant:6333",
internalEmbeddingUrl: "http://embedding:11434",
internalEmbeddingModel: "qwen3-embedding:0.6b",
internalEmbeddingDimensions: 1024,
};
export function createWorkspaceDiagnoser(
adapters: DiagnosticAdapters,
options: { timeoutMs?: number; semanticRuntime?: SemanticRuntimeConfig } = {},
) {
const timeoutMs = configuredTimeout(options.timeoutMs);
const semanticRuntime = options.semanticRuntime ?? DEFAULT_SEMANTIC_RUNTIME;
return async function diagnose(
workspace: WorkspaceDescriptor,
bindings: RuntimeBindings,
_options: { writeProbe: boolean },
): Promise<ConnectorDiagnostics> {
requireSupportedDescriptor(workspace);
const descriptor = validateWorkspaceDescriptor(workspace);
return await diagnoseValidatedWorkspace(descriptor, bindings, adapters, timeoutMs, semanticRuntime);
};
}
export function createProductionWorkspaceDiagnoser(
timeoutMs: number,
adapters: DiagnosticAdapters = createConcreteDiagnosticAdapters(),
semanticRuntime: SemanticRuntimeConfig = DEFAULT_SEMANTIC_RUNTIME,
) {
return createWorkspaceDiagnoser(adapters, { timeoutMs, semanticRuntime });
}
export const diagnoseWorkspace = createProductionWorkspaceDiagnoser(
DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS,
);