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

854 lines
36 KiB
TypeScript

import { randomUUID } from "node:crypto";
import { readFile, realpath } from "node:fs/promises";
import { createConnection } from "node:net";
import { createServer } from "node:net";
import { once } from "node:events";
import { spawn } from "node:child_process";
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,
validateCanonicalWorkspace,
type CanonicalWorkspace,
type RestDiagnosticRequest,
type WorkspaceDescriptor,
} from "./schema.js";
import type { WorkspaceErrorCode } from "./types.js";
export interface Diagnostic {
level: "error" | "warning" | "info";
code: WorkspaceErrorCode | "binding_ok";
field?: string;
message: string;
}
export interface WorkspaceDiagnostics {
activatable: boolean;
diagnostics: Diagnostic[];
}
type ConnectorRole = "dwh" | "vector";
interface DiagnosticResource {
database?: string;
schema?: string;
collection?: string;
}
type RestConnectorDiagnostic = RestDiagnosticRequest & {
response?: Record<string, string>;
};
export interface ConnectorDiagnosticRequest {
role: ConnectorRole;
transport: "postgres_direct" | "pgvector_direct" | "rest_api" | "ssh_tunnel";
host?: string;
port?: number;
baseUrl?: string;
user?: string;
credentialFile?: string;
tlsCaFile?: string;
resource: DiagnosticResource;
timeoutMs: number;
signal: AbortSignal;
diagnostic?: RestConnectorDiagnostic;
}
export interface ConnectorDiagnosticResult {
resolved: boolean;
tlsVerified: boolean;
authenticated: boolean;
resource: DiagnosticResource;
}
export interface SshTunnelRequest {
sshHost: string;
sshPort: number;
sshUser: string;
privateKeyFile: string;
knownHostsFile: string;
targetHost: string;
targetPort: number;
localHost: "127.0.0.1";
localPort: 0;
timeoutMs: number;
signal: AbortSignal;
}
export interface LoopbackTunnel {
host: "127.0.0.1";
port: number;
}
export interface VectorDiagnosticRequest {
transport?: "pgvector_direct" | "rest_api" | "ssh_tunnel";
baseUrl?: string;
credentialFile?: string;
tlsCaFile?: string;
diagnostic?: RestDiagnosticRequest & {
response: { collection: string; dimensions: string; distance: string };
};
collection: string;
dimensions?: number;
distance?: "cosine" | "l2" | "inner_product";
host?: string;
port?: number;
user?: string;
resource?: DiagnosticResource;
timeoutMs: number;
signal: AbortSignal;
}
export interface VectorDiagnosticResult {
collection?: string;
dimensions?: number;
distance?: "cosine" | "l2" | "inner_product";
}
export interface EmbeddingDiagnosticRequest {
baseUrl: string;
credentialFile?: string;
tlsCaFile?: string;
model: string;
timeoutMs: number;
signal: AbortSignal;
diagnostic?: RestDiagnosticRequest & { response: { model: string; dimensions: string } };
}
export interface EmbeddingDiagnosticResult {
available: boolean;
dimensions?: number;
}
export interface WriteDiagnosticRecordRequest {
collection: string;
id: string;
dimensions: number;
timeoutMs: number;
signal: AbortSignal;
credentialFile?: string;
tlsCaFile?: string;
baseUrl?: string;
diagnostic?: RestDiagnosticRequest & {
response: { operation: string };
};
}
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; signal: AbortSignal;
}): Promise<DatabaseDiagnosticClient>;
}
export interface SshProcessFactory {
start(request: SshTunnelRequest, args: readonly string[]): Promise<{
tunnel: LoopbackTunnel;
close(): Promise<void>;
}>;
}
export interface ConcreteDiagnosticAdapterDependencies {
directProtocol?: DirectProtocolFactory;
sshProcess?: SshProcessFactory;
databaseClient?: DatabaseDiagnosticClientFactory;
sshSpawn?: (args: readonly string[]) => { kill(signal?: NodeJS.Signals): boolean; once?(event: "error" | "exit", listener: (...args: any[]) => void): unknown; stderr?: { on(event: "data", listener: (data: Buffer | string) => void): unknown; off?(event: "data", listener: (data: Buffer | string) => void): unknown } };
reserveLoopbackPort?: () => Promise<number>;
sshForwardConfirmed?: (tunnel: LoopbackTunnel, signal: AbortSignal) => Promise<void>;
}
/**
* Adapters own protocol-specific I/O. They receive only binding file paths, never secret
* contents, and return metadata only; response bodies must stay inside the adapter.
*/
export interface DiagnosticAdapters {
probeConnector(request: ConnectorDiagnosticRequest): Promise<ConnectorDiagnosticResult>;
withSshTunnel<T>(
request: SshTunnelRequest,
probe: (tunnel: LoopbackTunnel) => Promise<T>,
): Promise<T>;
inspectVector(request: VectorDiagnosticRequest): Promise<VectorDiagnosticResult>;
probeEmbedding(request: EmbeddingDiagnosticRequest): Promise<EmbeddingDiagnosticResult>;
writeDiagnosticRecord(request: WriteDiagnosticRecordRequest): Promise<void>;
removeDiagnosticRecord(request: WriteDiagnosticRecordRequest): Promise<void>;
}
export const DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS = 5_000;
async function connectTcp(host: string, port: number, signal: AbortSignal): Promise<void> {
const socket = createConnection({ host, port });
const abort = () => socket.destroy();
signal.addEventListener("abort", abort, { once: true });
try {
await Promise.race([once(socket, "connect"), once(socket, "error").then(([error]) => Promise.reject(error))]);
} finally {
signal.removeEventListener("abort", abort);
socket.destroy();
}
}
async function secretPresent(file: string): Promise<boolean> {
return (await readFile(file, "utf8")).trim().length > 0;
}
async function sameSecretFile(first: string, second: string): Promise<boolean> {
try {
return await realpath(first) === await realpath(second);
} catch {
return first === second;
}
}
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 };
}
async function reserveLoopbackPort(): Promise<number> {
const server = createServer();
await new Promise<void>((resolve, reject) => {
server.once("error", reject);
server.listen(0, "127.0.0.1", resolve);
});
try {
const address = server.address();
if (!address || typeof address === "string") throw new Error("SSH tunnel port unavailable");
return address.port;
} finally {
await new Promise<void>((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
}
/**
* 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 spawnSsh = dependencies.sshSpawn ?? ((args: readonly string[]) => spawn("ssh", [...args], { stdio: ["ignore", "ignore", "pipe"] }));
const reserveSshPort = dependencies.reserveLoopbackPort ?? reserveLoopbackPort;
const sshProcess = dependencies.sshProcess ?? {
async start(request: SshTunnelRequest, args: readonly string[]) {
const port = await reserveSshPort();
const resolvedArgs = args.map((argument) => argument === `127.0.0.1:0:${request.targetHost}:${request.targetPort}`
? `127.0.0.1:${port}:${request.targetHost}:${request.targetPort}` : argument);
const child = spawnSsh(resolvedArgs);
const abort = () => { child.kill("SIGTERM"); };
request.signal.addEventListener("abort", abort, { once: true });
const tunnel = { host: "127.0.0.1" as const, port };
try {
await withTimeout(request.timeoutMs, async (signal) => {
await Promise.race([
dependencies.sshForwardConfirmed
? dependencies.sshForwardConfirmed(tunnel, signal)
: new Promise<void>((resolve, reject) => {
let stderrBuffer = "";
const confirmation = new RegExp(`Local forwarding listening on 127\\.0\\.0\\.1 port ${port}\\.?`);
const confirm = (data: Buffer | string) => {
stderrBuffer = `${stderrBuffer}${data.toString()}`.slice(-4096);
const lines = stderrBuffer.split(/\r?\n/);
stderrBuffer = lines.pop() ?? "";
if (lines.some((line) => confirmation.test(line))) {
child.stderr?.off?.("data", confirm);
resolve();
}
};
if (!child.stderr) return reject(new Error("SSH tunnel readiness failed"));
child.stderr.on("data", confirm);
signal.addEventListener("abort", () => reject(new Error("SSH tunnel readiness failed")), { once: true });
}),
new Promise<never>((_resolve, reject) => {
child.once?.("error", () => reject(new Error("SSH tunnel readiness failed")));
child.once?.("exit", () => reject(new Error("SSH tunnel readiness failed")));
}),
]);
});
} catch (error) {
request.signal.removeEventListener("abort", abort);
child.kill("SIGTERM");
throw error;
}
return {
tunnel,
async close() {
request.signal.removeEventListener("abort", abort);
const exited = child.once
? new Promise<void>((resolve) => child.once?.("exit", resolve))
: Promise.resolve();
child.kill("SIGTERM");
await withTimeout(request.timeoutMs, () => exited).catch(() => undefined);
},
};
},
};
const databaseClient = dependencies.databaseClient ?? {
async connect(request: { host: string; port: number; database: string; user: string; credentialFile: string; tlsCaFile?: 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"), rejectUnauthorized: true }
: { 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, 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 (request.diagnostic && "response" in request.diagnostic) {
const payload = await response.json().catch(() => undefined) as Record<string, unknown> | undefined;
const declared = request.diagnostic.response as { database?: string; schema?: string };
if (!payload || (declared.database && payload[declared.database] !== request.resource.database)
|| (declared.schema && payload[declared.schema] !== request.resource.schema)) throw new Error("REST probe failed");
}
return {
resolved: true,
tlsVerified: new URL(request.baseUrl).protocol === "https:",
authenticated: true,
resource: request.resource,
};
}
return await directProtocol.probe(request);
},
async withSshTunnel(request, probe) {
// The image supplies OpenSSH for the registry's SSH implementation. This adapter refuses
// an unverified host rather than falling back to an unsafe SSH option; the route-level
// tunnel owner supplies the process lifecycle in the next registry task.
if (!request.knownHostsFile || !(await secretPresent(request.privateKeyFile))) {
throw new Error("SSH probe failed");
}
const args = [
"-N", "-v", "-o", "BatchMode=yes", "-o", "ExitOnForwardFailure=yes", "-o", "StrictHostKeyChecking=yes",
"-o", `UserKnownHostsFile=${request.knownHostsFile}`, "-i", request.privateKeyFile,
"-p", String(request.sshPort), "-L", `127.0.0.1:0:${request.targetHost}:${request.targetPort}`,
`${request.sshUser}@${request.sshHost}`,
];
const tunnel = await sshProcess.start(request, args);
try {
return await probe(tunnel.tunnel);
} finally {
await tunnel.close().catch(() => undefined);
}
},
async inspectVector(request) {
if (request.transport === "pgvector_direct" || request.transport === "ssh_tunnel") {
const resource = request.resource;
if (!request.host || !request.port || !request.user || !request.credentialFile
|| !resource?.database || !resource.schema || !(await secretPresent(request.credentialFile))) {
throw new Error("vector metadata adapter is unavailable");
}
const client = await databaseClient.connect({
host: request.host, port: request.port, database: resource.database, user: request.user,
credentialFile: request.credentialFile, tlsCaFile: request.tlsCaFile, signal: request.signal,
});
try {
const metadata = await client.query(
"SELECT a.atttypmod - 4 AS dimensions, CASE WHEN pg_get_indexdef(i.indexrelid) LIKE '%vector_cosine_ops%' THEN 'cosine' WHEN pg_get_indexdef(i.indexrelid) LIKE '%vector_l2_ops%' THEN 'l2' WHEN pg_get_indexdef(i.indexrelid) LIKE '%vector_ip_ops%' THEN 'inner_product' END AS distance FROM pg_attribute a JOIN pg_class c ON c.oid = a.attrelid JOIN pg_namespace n ON n.oid = c.relnamespace JOIN pg_index i ON i.indrelid = c.oid AND a.attnum = ANY(i.indkey) WHERE n.nspname = $1 AND c.relname = $2 AND a.attnum > 0 AND NOT a.attisdropped AND a.atttypid = (SELECT oid FROM pg_type WHERE typname = 'vector') ORDER BY i.indexrelid LIMIT 1",
[resource.schema, request.collection],
);
const row = metadata.rows[0];
if (!row || !Number.isInteger(row.dimensions) || (row.distance !== "cosine" && row.distance !== "l2" && row.distance !== "inner_product")) throw new Error("vector metadata adapter is unavailable");
return { collection: request.collection, dimensions: row.dimensions as number, distance: row.distance as VectorDiagnosticResult["distance"] };
} finally {
await client.end().catch(() => undefined);
}
}
if (request.transport !== "rest_api" || !request.baseUrl || !request.diagnostic || request.tlsCaFile) {
throw new Error("vector metadata adapter is unavailable");
}
const response = await fetch(resolveDiagnosticUrl(request.baseUrl, request.diagnostic.path).toString(), {
method: request.diagnostic.method,
headers: await restHeaders(request.diagnostic, request.credentialFile),
signal: request.signal,
redirect: "error",
});
const payload = await response.json().catch(() => undefined) as Record<string, unknown> | undefined;
const fields = request.diagnostic.response;
if (!response.ok || !payload || typeof payload[fields.collection] !== "string"
|| !Number.isInteger(payload[fields.dimensions]) || typeof payload[fields.distance] !== "string") {
throw new Error("vector metadata adapter is unavailable");
}
return {
collection: payload[fields.collection] as string,
dimensions: payload[fields.dimensions] as number,
distance: payload[fields.distance] as VectorDiagnosticResult["distance"],
};
},
async probeEmbedding(request) {
if (!request.diagnostic || request.tlsCaFile) throw new Error("embedding probe failed");
const response = await fetch(resolveDiagnosticUrl(request.baseUrl, request.diagnostic.path).toString(), {
method: request.diagnostic.method,
headers: await restHeaders(request.diagnostic, request.credentialFile),
signal: request.signal,
redirect: "error",
});
const payload = await response.json().catch(() => undefined) as Record<string, unknown> | undefined;
if (!response.ok || !payload || payload[request.diagnostic.response.model] !== request.model
|| !Number.isInteger(payload[request.diagnostic.response.dimensions])) throw new Error("embedding probe failed");
return { available: true, dimensions: payload[request.diagnostic.response.dimensions] as number };
},
async writeDiagnosticRecord(request) {
if (!request.baseUrl || !request.diagnostic || request.tlsCaFile) throw new Error("vector write adapter is unavailable");
const response = await fetch(resolveDiagnosticUrl(request.baseUrl, request.diagnostic.path).toString(), {
method: request.diagnostic.method,
headers: {
...await restHeaders(request.diagnostic, request.credentialFile),
"content-type": "application/json",
},
body: JSON.stringify({ operation: "create", id: request.id, collection: request.collection, dimensions: request.dimensions }),
signal: request.signal,
redirect: "error",
});
const payload = await response.json().catch(() => undefined) as Record<string, unknown> | undefined;
if (!response.ok || !payload || payload[request.diagnostic.response.operation] !== "create") {
throw new Error("vector write adapter is unavailable");
}
},
async removeDiagnosticRecord(request) {
if (!request.baseUrl || !request.diagnostic || request.tlsCaFile) throw new Error("vector write adapter is unavailable");
const response = await fetch(resolveDiagnosticUrl(request.baseUrl, request.diagnostic.path).toString(), {
method: request.diagnostic.method,
headers: {
...await restHeaders(request.diagnostic, request.credentialFile),
"content-type": "application/json",
},
body: JSON.stringify({ operation: "remove", id: request.id, collection: request.collection }),
signal: request.signal,
redirect: "error",
});
const payload = await response.json().catch(() => undefined) as Record<string, unknown> | undefined;
if (!response.ok || !payload || payload[request.diagnostic.response.operation] !== "remove") {
throw new Error("vector write adapter is unavailable");
}
},
};
}
export function createProductionWorkspaceDiagnoser(
timeoutMs: number,
adapters: DiagnosticAdapters = createConcreteDiagnosticAdapters(),
) {
return createWorkspaceDiagnoser(adapters, { timeoutMs });
}
function boundedTimeout(value: number | undefined, fallback: number): number {
const selected = value ?? fallback;
return Math.min(Math.max(1, selected), fallback, 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."
: "Connector diagnostic failed.",
};
}
function bindingName(
workspace: CanonicalWorkspace,
role: "DWH" | "VECTOR" | "VECTOR_WRITER" | "EMBEDDING",
suffix: string,
): string {
const entry = buildInstallationContract(workspace).variables.find((variable) => (
variable.role === role && variable.suffix === suffix
));
if (!entry) throw new Error(`workspace contract is missing ${role}_${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 diagnosticsForMissingBindings(
workspace: CanonicalWorkspace,
bindings: RuntimeBindings,
): Diagnostic[] {
const missing = new Set([
...bindings.dwh.missing,
...bindings.vector.missing,
...bindings.embedding.missing,
]);
const knownHosts = [
bindings.dwh.transport === "ssh_tunnel" ? bindingName(workspace, "DWH", "SSH_KNOWN_HOSTS_FILE") : undefined,
bindings.vector.transport === "ssh_tunnel" ? bindingName(workspace, "VECTOR", "SSH_KNOWN_HOSTS_FILE") : undefined,
].filter((field): field is string => field !== undefined);
const ordered = [...new Set([...knownHosts.filter((field) => missing.has(field)), ...[...missing].sort()])];
return ordered.map((field) => diagnosticError("binding_missing", field));
}
function connectorRequest(
workspace: CanonicalWorkspace,
role: ConnectorRole,
bindings: RuntimeBindings,
timeoutMs: number,
): ConnectorDiagnosticRequest | SshTunnelRequest | undefined {
const binding = role === "dwh" ? bindings.dwh : bindings.vector;
const contractRole = role === "dwh" ? "DWH" : "VECTOR";
const values = binding.values;
const resource: DiagnosticResource = role === "dwh"
? { database: workspace.dwh.database, schema: workspace.dwh.schema }
: {
database: workspace.semantic_index.vector_store.database,
schema: workspace.semantic_index.vector_store.schema,
collection: workspace.semantic_index.vector_store.collection,
};
const field = (suffix: string) => bindingName(workspace, contractRole, suffix);
const credentialFile = values[field(binding.transport === "rest_api" ? "API_KEY_FILE" : "PASSWORD_FILE")];
if (credentialFile === undefined) return undefined;
if (binding.transport === "rest_api") {
const baseUrl = values[field("BASE_URL")];
const diagnostic = role === "dwh"
? workspace.diagnostics?.dwh_rest
: workspace.diagnostics?.vector_rest?.metadata;
if (baseUrl === undefined || diagnostic === undefined) return undefined;
return {
role,
transport: "rest_api",
baseUrl,
credentialFile,
tlsCaFile: values[field("TLS_CA_FILE")],
resource,
timeoutMs,
signal: new AbortController().signal,
diagnostic,
};
}
if (binding.transport === "ssh_tunnel") {
const sshHost = values[field("SSH_HOST")];
const sshPort = numericBinding(values, field("SSH_PORT"));
const sshUser = values[field("SSH_USER")];
const privateKeyFile = values[field("SSH_PRIVATE_KEY_FILE")];
const knownHostsFile = values[field("SSH_KNOWN_HOSTS_FILE")];
const targetHost = values[field("SSH_TARGET_HOST")];
const targetPort = numericBinding(values, field("SSH_TARGET_PORT"));
if (!sshHost || !sshPort || !sshUser || !privateKeyFile || !knownHostsFile || !targetHost || !targetPort) return undefined;
return {
sshHost, sshPort, sshUser, privateKeyFile, knownHostsFile, targetHost, targetPort,
localHost: "127.0.0.1", localPort: 0, timeoutMs, signal: new AbortController().signal,
};
}
const host = values[field("HOST")];
const port = numericBinding(values, field("PORT"));
const user = values[field("USER")];
if (!host || !port || !user) return undefined;
return {
role,
transport: binding.transport,
host,
port,
user,
credentialFile,
tlsCaFile: values[field("TLS_CA_FILE")],
resource,
timeoutMs,
signal: new AbortController().signal,
};
}
function tunnelProbeRequest(
workspace: CanonicalWorkspace,
role: ConnectorRole,
bindings: RuntimeBindings,
timeoutMs: number,
tunnel: LoopbackTunnel,
signal: AbortSignal,
): ConnectorDiagnosticRequest {
const binding = role === "dwh" ? bindings.dwh : bindings.vector;
const contractRole = role === "dwh" ? "DWH" : "VECTOR";
const password = binding.values[bindingName(workspace, contractRole, "PASSWORD_FILE")];
const user = binding.values[bindingName(workspace, contractRole, "USER")];
if (!password || !user) throw new Error("missing SSH connector credentials");
return {
role,
transport: "ssh_tunnel",
host: tunnel.host,
port: tunnel.port,
user,
credentialFile: password,
tlsCaFile: binding.values[bindingName(workspace, contractRole, "TLS_CA_FILE")],
resource: role === "dwh"
? { database: workspace.dwh.database, schema: workspace.dwh.schema }
: {
database: workspace.semantic_index.vector_store.database,
schema: workspace.semantic_index.vector_store.schema,
collection: workspace.semantic_index.vector_store.collection,
},
timeoutMs,
signal,
};
}
export function createWorkspaceDiagnoser(
adapters: DiagnosticAdapters,
options: { timeoutMs?: number } = {},
) {
const fallbackTimeout = boundedTimeout(options.timeoutMs, DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS);
return async function diagnoseWorkspace(
workspace: WorkspaceDescriptor,
bindings: RuntimeBindings,
options: { writeProbe: boolean },
): Promise<WorkspaceDiagnostics> {
const canonical = validateCanonicalWorkspace(workspace);
const diagnostics = diagnosticsForMissingBindings(canonical, bindings);
if (diagnostics.length > 0) return { activatable: false, diagnostics };
const dwhTimeout = boundedTimeout(canonical.dwh.timeout_ms, fallbackTimeout);
const vectorTimeout = boundedTimeout(canonical.semantic_index.vector_store.timeout_ms, fallbackTimeout);
const embeddingTimeout = boundedTimeout(canonical.semantic_index.embedding.timeout_ms, fallbackTimeout);
let tunneledVectorMetadata: VectorDiagnosticResult | undefined;
for (const role of ["dwh", "vector"] as const) {
const timeoutMs = role === "dwh" ? dwhTimeout : vectorTimeout;
const request = connectorRequest(canonical, role, bindings, timeoutMs);
if (!request) {
diagnostics.push(diagnosticError("binding_missing"));
continue;
}
try {
const result = "sshHost" in request
? await withTimeout(timeoutMs, (signal) => adapters.withSshTunnel(
{ ...request, signal },
async (tunnel) => {
const tunneledRequest = tunnelProbeRequest(canonical, role, bindings, timeoutMs, tunnel, signal);
const connector = await adapters.probeConnector(tunneledRequest);
if (role === "vector") {
tunneledVectorMetadata = await adapters.inspectVector({
transport: "ssh_tunnel", host: tunneledRequest.host, port: tunneledRequest.port,
user: tunneledRequest.user, credentialFile: tunneledRequest.credentialFile,
tlsCaFile: tunneledRequest.tlsCaFile, resource: tunneledRequest.resource,
collection: canonical.semantic_index.vector_store.collection,
dimensions: canonical.semantic_index.vector_store.dimensions,
distance: canonical.semantic_index.vector_store.distance,
timeoutMs: vectorTimeout, signal,
});
}
return connector;
},
))
: await withTimeout(timeoutMs, (signal) => adapters.probeConnector({ ...request, signal }));
const resource = role === "dwh"
? { database: canonical.dwh.database, schema: canonical.dwh.schema }
: {
database: canonical.semantic_index.vector_store.database,
schema: canonical.semantic_index.vector_store.schema,
collection: canonical.semantic_index.vector_store.collection,
};
if (!hasRequiredConnectorChecks(result, resource)) diagnostics.push(diagnosticError("connector_unavailable"));
else diagnostics.push({ level: "info", code: "binding_ok", message: `${role === "dwh" ? "DWH" : "Vector"} binding diagnostic passed.` });
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
}
}
if (!diagnostics.some((diagnostic) => diagnostic.level === "error")) {
try {
const vectorBinding = bindings.vector;
const vectorRest = canonical.diagnostics?.vector_rest?.metadata;
const vectorDirect = connectorRequest(canonical, "vector", bindings, vectorTimeout);
const directVectorRequest = vectorDirect && !("sshHost" in vectorDirect) ? vectorDirect : undefined;
const vector = tunneledVectorMetadata ?? await withTimeout(vectorTimeout, (signal) => adapters.inspectVector({
transport: vectorBinding.transport === "pgvector_direct" || vectorBinding.transport === "rest_api"
|| vectorBinding.transport === "ssh_tunnel" ? vectorBinding.transport : undefined,
baseUrl: vectorBinding.values[bindingName(canonical, "VECTOR", "BASE_URL")],
credentialFile: vectorBinding.values[bindingName(canonical, "VECTOR", "API_KEY_FILE")],
tlsCaFile: vectorBinding.values[bindingName(canonical, "VECTOR", "TLS_CA_FILE")],
diagnostic: vectorRest,
collection: canonical.semantic_index.vector_store.collection,
dimensions: canonical.semantic_index.vector_store.dimensions,
distance: canonical.semantic_index.vector_store.distance,
...(directVectorRequest ? {
host: directVectorRequest.host,
port: directVectorRequest.port,
user: directVectorRequest.user,
credentialFile: directVectorRequest.credentialFile,
tlsCaFile: directVectorRequest.tlsCaFile,
resource: directVectorRequest.resource,
} : {}),
timeoutMs: vectorTimeout,
signal,
}));
const expected = canonical.semantic_index.vector_store;
if (
vector.collection !== expected.collection
|| vector.dimensions !== expected.dimensions
|| vector.distance !== expected.distance
) diagnostics.push(diagnosticError("semantic_index_incompatible"));
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
}
}
if (!diagnostics.some((diagnostic) => diagnostic.level === "error")) {
try {
const embedding = await withTimeout(embeddingTimeout, (signal) => adapters.probeEmbedding({
baseUrl: bindings.embedding.values[bindingName(canonical, "EMBEDDING", "BASE_URL")] ?? "",
credentialFile: bindings.embedding.values[bindingName(canonical, "EMBEDDING", "API_KEY_FILE")],
tlsCaFile: bindings.embedding.values[bindingName(canonical, "EMBEDDING", "TLS_CA_FILE")],
model: canonical.semantic_index.embedding.model,
timeoutMs: embeddingTimeout,
signal,
diagnostic: canonical.diagnostics?.embedding,
}));
if (!embedding.available || embedding.dimensions !== canonical.semantic_index.embedding.dimensions) {
diagnostics.push(diagnosticError("semantic_index_incompatible"));
}
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
}
}
if (
options.writeProbe
&& canonical.semantic_index.vector_writer
&& canonical.diagnostics?.vector_rest?.reversible_probe
&& bindings.vector.transport === "rest_api"
&& !diagnostics.some((diagnostic) => diagnostic.level === "error")
) {
const credentialFile = bindings.vectorWriter.values[bindingName(canonical, "VECTOR_WRITER", "API_KEY_FILE")];
if (!credentialFile) return { activatable: true, diagnostics };
const readerCredentialFile = bindings.vector.values[bindingName(canonical, "VECTOR", "API_KEY_FILE")];
if (readerCredentialFile && await sameSecretFile(credentialFile, readerCredentialFile)) {
diagnostics.push(diagnosticError("binding_missing", bindingName(canonical, "VECTOR_WRITER", "API_KEY_FILE")));
return { activatable: false, diagnostics };
}
const request: WriteDiagnosticRecordRequest = {
collection: canonical.semantic_index.vector_store.collection,
id: `diagnostic:${randomUUID()}`,
dimensions: canonical.semantic_index.vector_store.dimensions,
timeoutMs: vectorTimeout,
signal: new AbortController().signal,
credentialFile,
tlsCaFile: bindings.vector.values[bindingName(canonical, "VECTOR", "TLS_CA_FILE")],
baseUrl: bindings.vector.values[bindingName(canonical, "VECTOR", "BASE_URL")],
diagnostic: canonical.diagnostics.vector_rest.reversible_probe,
};
let writeStarted = false;
let cleanupAttempted = false;
let cleanupFailed = false;
try {
writeStarted = true;
await withTimeout(vectorTimeout, (signal) => adapters.writeDiagnosticRecord({ ...request, signal }));
cleanupAttempted = true;
await withTimeout(vectorTimeout, (signal) => adapters.removeDiagnosticRecord({ ...request, signal }));
} catch {
cleanupFailed = true;
} finally {
if (writeStarted && (!cleanupAttempted || cleanupFailed)) {
try {
await withTimeout(vectorTimeout, (signal) => adapters.removeDiagnosticRecord({ ...request, signal }));
} catch {
cleanupFailed = true;
}
}
}
if (cleanupFailed) {
diagnostics.push(diagnosticError("connector_unavailable"));
}
}
return {
activatable: !diagnostics.some((diagnostic) => diagnostic.level === "error"),
diagnostics,
};
};
}
export const diagnoseWorkspace = createProductionWorkspaceDiagnoser(
DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS,
);