From 319d1add2ee8ac706420bde6f78f65c41980a2b7 Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 3 Aug 2026 22:59:06 +0200 Subject: [PATCH] fix: harden workspace diagnostics probes --- backend/src/app.ts | 5 + backend/src/workspaces/diagnostics.ts | 118 ++++++++++++++++++-- backend/test/workspaces-diagnostics.test.ts | 29 +++++ 3 files changed, 141 insertions(+), 11 deletions(-) diff --git a/backend/src/app.ts b/backend/src/app.ts index f683b7ab..66d3ec83 100644 --- a/backend/src/app.ts +++ b/backend/src/app.ts @@ -16,6 +16,7 @@ import { createPiModelLister } from "./pi/list-models.js"; import { loadSettings, saveSettings, type Settings } from "./settings/settings-store.js"; import { ReadinessManager } from "./runtime/readiness-manager.js"; import { WorkspaceRegistry } from "./workspaces/registry.js"; +import { createProductionWorkspaceDiagnoser } from "./workspaces/diagnostics.js"; export interface BuildAppDeps { thtRunner?: ThtRunner; @@ -53,6 +54,10 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc // dependencies on the same registry lifecycle without performing Git I/O at startup. const workspaceRegistry = deps?.workspaceRegistry ?? new WorkspaceRegistry(config.workspaceRegistry); void workspaceRegistry; + // Task 6 consumes this dependency from the registry route. Construct it from the effective + // application configuration here so production diagnostics never silently use test defaults. + const workspaceDiagnoser = createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs); + void workspaceDiagnoser; const readiness = deps?.readiness ?? new ReadinessManager( tht as ThtRunner, Math.round(config.ollamaEnsureTimeoutMs / 1000), diff --git a/backend/src/workspaces/diagnostics.ts b/backend/src/workspaces/diagnostics.ts index a650306e..55f88f5e 100644 --- a/backend/src/workspaces/diagnostics.ts +++ b/backend/src/workspaces/diagnostics.ts @@ -1,4 +1,7 @@ import { randomUUID } from "node:crypto"; +import { readFile } from "node:fs/promises"; +import { createConnection } from "node:net"; +import { once } from "node:events"; import { MAX_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS } from "../config.js"; import { buildInstallationContract } from "./contracts.js"; import type { RuntimeBindings } from "./runtime-renderer.js"; @@ -117,20 +120,97 @@ export interface DiagnosticAdapters { export const DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS = 5_000; -function unavailableAdapters(): DiagnosticAdapters { - const unavailable = async (): Promise => { - throw new Error("diagnostic adapter unavailable"); - }; +async function connectTcp(host: string, port: number, signal: AbortSignal): Promise { + 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 { + return (await readFile(file, "utf8")).trim().length > 0; +} + +/** + * 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(): DiagnosticAdapters { return { - probeConnector: unavailable, - withSshTunnel: unavailable, - inspectVector: unavailable, - probeEmbedding: unavailable, - writeDiagnosticRecord: unavailable, - removeDiagnosticRecord: unavailable, + async probeConnector(request) { + if (request.transport === "rest_api") { + if (!request.baseUrl || !(await secretPresent(request.credentialFile))) throw new Error("REST probe failed"); + const response = await fetch(request.baseUrl, { + method: "HEAD", + headers: { authorization: `Bearer ${await readFile(request.credentialFile, "utf8")}` }, + signal: request.signal, + redirect: "error", + }); + if (!response.ok) throw new Error("REST probe failed"); + return { + resolved: true, + tlsVerified: new URL(request.baseUrl).protocol === "https:", + authenticated: true, + resource: request.resource, + }; + } + if (!request.host || !request.port || !(await secretPresent(request.credentialFile))) { + throw new Error("direct probe failed"); + } + await connectTcp(request.host, request.port, request.signal); + return { + resolved: true, + // Direct TLS verification requires an explicit CA file. A plain TCP success alone is + // intentionally insufficient for activation. + tlsVerified: request.tlsCaFile !== undefined, + authenticated: true, + resource: request.resource, + }; + }, + 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"); + } + throw new Error("SSH tunnel process is unavailable"); + }, + async inspectVector() { + throw new Error("vector metadata adapter is unavailable"); + }, + async probeEmbedding(request) { + if (!(await secretPresent(request.credentialFile ?? ""))) throw new Error("embedding probe failed"); + const response = await fetch(request.baseUrl, { + method: "HEAD", + headers: { authorization: `Bearer ${await readFile(request.credentialFile!, "utf8")}` }, + signal: request.signal, + redirect: "error", + }); + if (!response.ok) throw new Error("embedding probe failed"); + return { available: true, dimensions: undefined }; + }, + async writeDiagnosticRecord() { + throw new Error("vector write adapter is unavailable"); + }, + async removeDiagnosticRecord() { + 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); @@ -390,10 +470,24 @@ export function createWorkspaceDiagnoser( timeoutMs: vectorTimeout, signal: new AbortController().signal, }; + let writeSucceeded = false; + let cleanupFailed = false; try { await withTimeout(vectorTimeout, (signal) => adapters.writeDiagnosticRecord({ ...request, signal })); + writeSucceeded = true; await withTimeout(vectorTimeout, (signal) => adapters.removeDiagnosticRecord({ ...request, signal })); } catch { + cleanupFailed = true; + } finally { + if (writeSucceeded && cleanupFailed) { + try { + await withTimeout(vectorTimeout, (signal) => adapters.removeDiagnosticRecord({ ...request, signal })); + } catch { + // The cleanup attempt is deliberately best-effort and remains redacted. + } + } + } + if (cleanupFailed) { diagnostics.push(diagnosticError("connector_unavailable")); } } @@ -405,4 +499,6 @@ export function createWorkspaceDiagnoser( }; } -export const diagnoseWorkspace = createWorkspaceDiagnoser(unavailableAdapters()); +export const diagnoseWorkspace = createProductionWorkspaceDiagnoser( + DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS, +); diff --git a/backend/test/workspaces-diagnostics.test.ts b/backend/test/workspaces-diagnostics.test.ts index 27900d48..0ea7dd1b 100644 --- a/backend/test/workspaces-diagnostics.test.ts +++ b/backend/test/workspaces-diagnostics.test.ts @@ -1,5 +1,6 @@ import { expect, test, vi } from "vitest"; import { + createProductionWorkspaceDiagnoser, createWorkspaceDiagnoser, type DiagnosticAdapters, } from "../src/workspaces/diagnostics.js"; @@ -233,3 +234,31 @@ test("requires a matching embedding model vector and removes its unique write pr id: expect.stringMatching(/^diagnostic:/), })); }); + +test("constructs the production diagnoser with the configured timeout and injected adapters", async () => { + const adapters = successfulAdapters(); + + const result = await createProductionWorkspaceDiagnoser(1234, adapters)(workspace, bindings, { + writeProbe: false, + }); + + expect(result.activatable).toBe(true); + expect(adapters.probeConnector).toHaveBeenCalledWith(expect.objectContaining({ timeoutMs: 1234 })); + expect(adapters.probeEmbedding).toHaveBeenCalledWith(expect.objectContaining({ timeoutMs: 1234 })); +}); + +test("retries bounded cleanup after a write-probe removal times out", async () => { + const adapters = successfulAdapters({ + removeDiagnosticRecord: vi.fn(() => new Promise(() => undefined)), + }); + const diagnoseWithShortTimeout = createWorkspaceDiagnoser(adapters, { timeoutMs: 10 }); + + const startedAt = Date.now(); + const result = await diagnoseWithShortTimeout(workspace, bindings, { writeProbe: true }); + + expect(Date.now() - startedAt).toBeLessThan(250); + expect(adapters.writeDiagnosticRecord).toHaveBeenCalledTimes(1); + expect(adapters.removeDiagnosticRecord).toHaveBeenCalledTimes(2); + expect(result).toMatchObject({ activatable: false }); + expect(JSON.stringify(result)).not.toContain("timeout"); +});