fix: retire active pgvector artifacts

This commit is contained in:
2026-08-08 19:27:05 +02:00
parent 4e3fecbe8e
commit a6dbe1d023
16 changed files with 630 additions and 292 deletions
+6 -1
View File
@@ -69,7 +69,12 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
const hub = deps?.hub ?? new SseHub();
const workspaceRegistry = deps?.workspaceRegistry ?? new WorkspaceRegistry(config.workspaceRegistry);
const workspaceDiagnoser = deps?.workspaceDiagnoser
?? createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs);
?? createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs, undefined, {
internalQdrantUrl: config.internalQdrantUrl,
internalEmbeddingUrl: config.internalEmbeddingUrl,
internalEmbeddingModel: config.internalEmbeddingModel,
internalEmbeddingDimensions: config.internalEmbeddingDimensions,
});
const workspaceRuntimeSupport = deps?.workspaceRuntimeSupport ?? ((workspace: WorkspaceDescriptor) => (
supportsSessionRuntime(resolveRuntimeBindings(
workspace,
+230 -2
View File
@@ -16,6 +16,7 @@ import {
type WorkspaceDescriptor,
} from "./schema.js";
import type { WorkspaceErrorCode } from "./types.js";
import type { SemanticRuntimeConfig } from "./runtime-renderer.js";
export interface Diagnostic {
level: "error" | "warning" | "info";
@@ -395,6 +396,26 @@ export function createConcreteDiagnosticAdapters(
}
},
async inspectVector(request) {
if (request.transport === "rest_api" && request.baseUrl && request.diagnostic === undefined) {
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") {
throw new Error("vector metadata adapter is unavailable");
}
return {
collection: request.collection,
dimensions: size as number,
distance: distance.toLowerCase() as VectorDiagnosticResult["distance"],
};
}
if (request.transport === "pgvector_direct" || request.transport === "ssh_tunnel") {
const resource = request.resource;
if (!request.host || !request.port || !request.user || !request.credentialFile
@@ -440,6 +461,21 @@ export function createConcreteDiagnosticAdapters(
};
},
async probeEmbedding(request) {
if (!request.diagnostic && !request.tlsCaFile) {
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 };
}
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,
@@ -492,8 +528,31 @@ export function createConcreteDiagnosticAdapters(
export function createProductionWorkspaceDiagnoser(
timeoutMs: number,
adapters: DiagnosticAdapters = createConcreteDiagnosticAdapters(),
semanticRuntime: SemanticRuntimeConfig = {
internalQdrantUrl: "http://qdrant:6333",
internalEmbeddingUrl: "http://embedding:11434",
internalEmbeddingModel: "qwen3-embedding:0.6b",
internalEmbeddingDimensions: 1024,
},
) {
return createWorkspaceDiagnoser(adapters, { timeoutMs });
const legacyDiagnoser = createWorkspaceDiagnoser(adapters, { timeoutMs });
return async (
workspace: WorkspaceDescriptor,
bindings: RuntimeBindings,
options: { writeProbe: boolean },
): Promise<WorkspaceDiagnostics> => {
const descriptor = validateWorkspaceDescriptor(workspace);
if (descriptor.workspace.schema_version !== 3) {
return await legacyDiagnoser(descriptor, bindings, options);
}
return await diagnoseSchemaV3Workspace(
descriptor as Extract<WorkspaceDescriptor, { workspace: { schema_version: 3 } }>,
bindings,
adapters,
timeoutMs,
semanticRuntime,
);
};
}
function boundedTimeout(value: number | undefined, fallback: number): number {
@@ -526,6 +585,22 @@ function hasRequiredConnectorChecks(result: ConnectorDiagnosticResult, resource:
return result.resolved && result.tlsVerified && result.authenticated && sameResource(resource, result.resource);
}
function hasMatchingVectorMetadata(
actual: VectorDiagnosticResult,
expected: { collection: string; dimensions: number; distance: "cosine" | "l2" | "inner_product" },
): boolean {
return actual.collection === expected.collection
&& actual.dimensions === expected.dimensions
&& actual.distance === expected.distance;
}
function hasMatchingEmbeddingMetadata(
actual: EmbeddingDiagnosticResult,
expected: { dimensions: number },
): boolean {
return actual.available && actual.dimensions === expected.dimensions;
}
function diagnosticError(code: WorkspaceErrorCode, field?: string): Diagnostic {
return {
level: "error",
@@ -542,7 +617,7 @@ function diagnosticError(code: WorkspaceErrorCode, field?: string): Diagnostic {
}
function bindingName(
workspace: WorkspaceV2,
workspace: WorkspaceDescriptor,
role: "DWH" | "VECTOR" | "VECTOR_WRITER" | "EMBEDDING",
suffix: string,
): string {
@@ -558,6 +633,159 @@ function numericBinding(binding: Record<string, string>, name: string): number |
return Number.isInteger(value) && value > 0 && value <= 65_535 ? value : undefined;
}
async function diagnoseSchemaV3Workspace(
descriptor: Extract<WorkspaceDescriptor, { workspace: { schema_version: 3 } }>,
bindings: RuntimeBindings,
adapters: DiagnosticAdapters,
timeoutMs: number,
semanticRuntime: SemanticRuntimeConfig,
): Promise<WorkspaceDiagnostics> {
const diagnostics = [...bindings.dwh.missing]
.sort()
.map((field) => diagnosticError("binding_missing", field));
if (diagnostics.length > 0) {
return { activatable: false, diagnostics };
}
const dwhTimeout = boundedTimeout(descriptor.dwh.timeout_ms, timeoutMs);
const vectorTimeout = timeoutMs;
const embeddingTimeout = timeoutMs;
let activatable = true;
const dwhValues = bindings.dwh.values;
const dwhField = (suffix: string) => bindingName(descriptor, "DWH", suffix);
const dwhResource = { database: descriptor.dwh.database, schema: descriptor.dwh.schema };
let dwhRequest: ConnectorDiagnosticRequest | SshTunnelRequest | 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,
};
}
} else {
const sshHost = dwhValues[dwhField("SSH_HOST")];
const sshPort = numericBinding(dwhValues, dwhField("SSH_PORT"));
const sshUser = dwhValues[dwhField("SSH_USER")];
const privateKeyFile = dwhValues[dwhField("SSH_PRIVATE_KEY_FILE")];
const knownHostsFile = dwhValues[dwhField("SSH_KNOWN_HOSTS_FILE")];
const targetHost = dwhValues[dwhField("SSH_TARGET_HOST")];
const targetPort = numericBinding(dwhValues, dwhField("SSH_TARGET_PORT"));
if (sshHost && sshPort && sshUser && privateKeyFile && knownHostsFile && targetHost && targetPort) {
dwhRequest = {
sshHost,
sshPort,
sshUser,
privateKeyFile,
knownHostsFile,
targetHost,
targetPort,
localHost: "127.0.0.1",
localPort: 0,
timeoutMs: dwhTimeout,
signal: new AbortController().signal,
};
}
}
if (!dwhRequest || "sshHost" in 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(vectorTimeout, (signal) => adapters.inspectVector({
transport: "rest_api",
baseUrl: semanticRuntime.internalQdrantUrl,
collection: descriptor.semantic_index.vector_store.collection,
dimensions: descriptor.semantic_index.vector_store.dimensions,
distance: descriptor.semantic_index.vector_store.distance,
timeoutMs: vectorTimeout,
signal,
}));
if (!hasMatchingVectorMetadata(vector, descriptor.semantic_index.vector_store)) {
diagnostics.push(diagnosticError("semantic_index_incompatible"));
activatable = false;
}
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
activatable = false;
}
try {
const embedding = await withTimeout(embeddingTimeout, (signal) => adapters.probeEmbedding({
baseUrl: semanticRuntime.internalEmbeddingUrl,
model: semanticRuntime.internalEmbeddingModel,
timeoutMs: embeddingTimeout,
signal,
}));
if (
semanticRuntime.internalEmbeddingModel !== descriptor.semantic_index.embedding.model
|| semanticRuntime.internalEmbeddingDimensions !== descriptor.semantic_index.embedding.dimensions
|| !hasMatchingEmbeddingMetadata(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." }],
};
}
function diagnosticsForMissingBindings(
workspace: WorkspaceV2,
bindings: RuntimeBindings,
@@ -160,6 +160,38 @@ const writerBindings: RuntimeBindings = {
},
};
const workspaceV3 = parseWorkspaceYaml(`workspace:
schema_version: 3
id: psd-clinical
name: Policlinico San Donato
language: it
dwh:
engine: postgres
database: warehouse
schema: datawarehouse
timeout_ms: 8000
supported_transports: [postgres_direct, rest_api]
semantic_index:
vector_store:
engine: qdrant
collection: psd-clinical
dimensions: 1024
distance: cosine
embedding:
provider: ollama_internal
model: qwen3-embedding:0.6b
dimensions: 1024
llm_policy:
allowed: [zai/glm-5.2]
`);
const bindingsV3: RuntimeBindings = {
dwh: bindings.dwh,
vector: { transport: "rest_api", missing: [], values: {} },
vectorWriter: { transport: "rest_api", missing: [], values: {} },
embedding: { transport: "rest_api", missing: [], values: {} },
};
function successfulAdapters(overrides: Partial<DiagnosticAdapters> = {}): DiagnosticAdapters {
return {
probeConnector: vi.fn(async (request) => ({
@@ -843,6 +875,63 @@ test("constructs the production diagnoser with the configured timeout and inject
expect(adapters.probeEmbedding).toHaveBeenCalledWith(expect.objectContaining({ timeoutMs: 1234 }));
});
test("diagnoses a schema-v3 workspace through internal Qdrant and embedding config without workspace semantic bindings", async () => {
const adapters = successfulAdapters({
inspectVector: vi.fn(async () => ({
collection: "psd-clinical",
dimensions: 1024,
distance: "cosine",
})),
probeEmbedding: vi.fn(async () => ({ available: true, dimensions: 1024 })),
});
const result = await createProductionWorkspaceDiagnoser(1234, adapters, {
internalQdrantUrl: "http://qdrant:6333",
internalEmbeddingUrl: "http://embedding:11434",
internalEmbeddingModel: "qwen3-embedding:0.6b",
internalEmbeddingDimensions: 1024,
})(workspaceV3, bindingsV3, { writeProbe: false });
expect(result.activatable).toBe(true);
expect(adapters.inspectVector).toHaveBeenCalledWith(expect.objectContaining({
transport: "rest_api",
baseUrl: "http://qdrant:6333",
collection: "psd-clinical",
dimensions: 1024,
distance: "cosine",
timeoutMs: 1234,
}));
expect(adapters.probeEmbedding).toHaveBeenCalledWith(expect.objectContaining({
baseUrl: "http://embedding:11434",
model: "qwen3-embedding:0.6b",
timeoutMs: 1234,
}));
expect(adapters.probeConnector).toHaveBeenCalledTimes(1);
});
test("fails closed for schema-v3 when internal semantic diagnostics do not match descriptor identity", async () => {
const adapters = successfulAdapters({
inspectVector: vi.fn(async () => ({
collection: "wrong-collection",
dimensions: 1024,
distance: "cosine",
})),
probeEmbedding: vi.fn(async () => ({ available: true, dimensions: 1024 })),
});
const result = await createProductionWorkspaceDiagnoser(1234, adapters, {
internalQdrantUrl: "http://qdrant:6333",
internalEmbeddingUrl: "http://embedding:11434",
internalEmbeddingModel: "qwen3-embedding:0.6b",
internalEmbeddingDimensions: 1024,
})(workspaceV3, bindingsV3, { writeProbe: false });
expect(result.activatable).toBe(false);
expect(result.diagnostics).toContainEqual(expect.objectContaining({
code: "semantic_index_incompatible",
}));
});
test("retries bounded cleanup after a write-probe removal times out", async () => {
const adapters = successfulAdapters({
removeDiagnosticRecord: vi.fn(() => new Promise<void>(() => undefined)),