fix: harden workspace diagnostics probes
This commit is contained in:
@@ -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),
|
||||
|
||||
@@ -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,19 +120,96 @@ export interface DiagnosticAdapters {
|
||||
|
||||
export const DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS = 5_000;
|
||||
|
||||
function unavailableAdapters(): DiagnosticAdapters {
|
||||
const unavailable = async (): Promise<never> => {
|
||||
throw new Error("diagnostic adapter unavailable");
|
||||
};
|
||||
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;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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;
|
||||
@@ -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,
|
||||
);
|
||||
|
||||
@@ -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<void>(() => 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");
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user