feat: diagnose workspace connector bindings
This commit is contained in:
@@ -0,0 +1,408 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { MAX_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS } from "../config.js";
|
||||
import { buildInstallationContract } from "./contracts.js";
|
||||
import type { RuntimeBindings } from "./runtime-renderer.js";
|
||||
import { validateCanonicalWorkspace, type CanonicalWorkspace } 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;
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
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 {
|
||||
collection: string;
|
||||
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;
|
||||
}
|
||||
|
||||
export interface EmbeddingDiagnosticResult {
|
||||
available: boolean;
|
||||
dimensions?: number;
|
||||
}
|
||||
|
||||
export interface WriteDiagnosticRecordRequest {
|
||||
collection: string;
|
||||
id: string;
|
||||
dimensions: number;
|
||||
timeoutMs: number;
|
||||
signal: AbortSignal;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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;
|
||||
|
||||
function unavailableAdapters(): DiagnosticAdapters {
|
||||
const unavailable = async (): Promise<never> => {
|
||||
throw new Error("diagnostic adapter unavailable");
|
||||
};
|
||||
return {
|
||||
probeConnector: unavailable,
|
||||
withSshTunnel: unavailable,
|
||||
inspectVector: unavailable,
|
||||
probeEmbedding: unavailable,
|
||||
writeDiagnosticRecord: unavailable,
|
||||
removeDiagnosticRecord: unavailable,
|
||||
};
|
||||
}
|
||||
|
||||
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" | "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 }
|
||||
: { 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")];
|
||||
if (baseUrl === undefined) return undefined;
|
||||
return {
|
||||
role,
|
||||
transport: "rest_api",
|
||||
baseUrl,
|
||||
credentialFile,
|
||||
tlsCaFile: values[field("TLS_CA_FILE")],
|
||||
resource,
|
||||
timeoutMs,
|
||||
signal: new AbortController().signal,
|
||||
};
|
||||
}
|
||||
|
||||
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 }
|
||||
: { 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: CanonicalWorkspace,
|
||||
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);
|
||||
|
||||
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 },
|
||||
(tunnel) => adapters.probeConnector(tunnelProbeRequest(
|
||||
canonical, role, bindings, timeoutMs, tunnel, signal,
|
||||
)),
|
||||
))
|
||||
: await withTimeout(timeoutMs, (signal) => adapters.probeConnector({ ...request, signal }));
|
||||
const resource = role === "dwh"
|
||||
? { database: canonical.dwh.database, schema: canonical.dwh.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 vector = await withTimeout(vectorTimeout, (signal) => adapters.inspectVector({
|
||||
collection: canonical.semantic_index.vector_store.collection,
|
||||
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,
|
||||
}));
|
||||
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 && !diagnostics.some((diagnostic) => diagnostic.level === "error")) {
|
||||
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,
|
||||
};
|
||||
try {
|
||||
await withTimeout(vectorTimeout, (signal) => adapters.writeDiagnosticRecord({ ...request, signal }));
|
||||
await withTimeout(vectorTimeout, (signal) => adapters.removeDiagnosticRecord({ ...request, signal }));
|
||||
} catch {
|
||||
diagnostics.push(diagnosticError("connector_unavailable"));
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
activatable: !diagnostics.some((diagnostic) => diagnostic.level === "error"),
|
||||
diagnostics,
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
export const diagnoseWorkspace = createWorkspaceDiagnoser(unavailableAdapters());
|
||||
Reference in New Issue
Block a user