feat: pre-check DWH reachability before creating a session (local dev only)
New session now refuses to spawn a Pi runtime that would only die in bootstrap
retrieval when the DWH/vector host is unreachable (e.g. a dropped VPN). Before
`session new`, POST /sessions probes the DWH via `tht db ping`; if it is down it
returns 503 {code:"dwh_unreachable"} with a clear message and creates nothing.
- Gated behind the THT_DWH_PRECHECK flag (default off), enabled only by the local
dev launcher (run-stack.sh) — containers/CI never pay the probe, and existing
tests that don't set it are unaffected.
- ThtRunner.dbPing() runs `tht db ping` with a 10s timeout (run() gains an optional
timeout that SIGKILLs a hung child).
- Frontend: apiFetch throws a typed ApiError (status + parsed payload); the new-
session composer shows the specific alert on `dwh_unreachable` instead of the
generic retry hint, keeping the question for retry.
Verified live on an isolated backend (precheck on + broken DWH host → 503
dwh_unreachable, no session created) and via unit tests (backend 228, frontend 308).
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
+1
-1
@@ -92,7 +92,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
|
|||||||
app.get("/health", async () => ({ status: "ok" }));
|
app.get("/health", async () => ({ status: "ok" }));
|
||||||
app.get("/me", async (req) => getPrincipal(req));
|
app.get("/me", async (req) => getPrincipal(req));
|
||||||
sessionRoutes(app, {
|
sessionRoutes(app, {
|
||||||
mgr, tht: tht as ThtRunner, hub, getSettings, readiness,
|
mgr, tht: tht as ThtRunner, hub, getSettings, readiness, dwhPrecheck: config.dwhPrecheck,
|
||||||
});
|
});
|
||||||
sqlRoutes(app, { tht: tht as ThtRunner, getSettings });
|
sqlRoutes(app, { tht: tht as ThtRunner, getSettings });
|
||||||
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
|
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
|
||||||
|
|||||||
@@ -16,6 +16,12 @@ export interface AppConfig {
|
|||||||
secretsFile?: string;
|
secretsFile?: string;
|
||||||
secretFiles: Readonly<Record<string, string | undefined>>;
|
secretFiles: Readonly<Record<string, string | undefined>>;
|
||||||
modelApiKeyFile?: string;
|
modelApiKeyFile?: string;
|
||||||
|
/**
|
||||||
|
* Local-only: when true, POST /sessions probes DWH reachability (`tht db ping`) and
|
||||||
|
* refuses to create a session if it is down. Off by default so containers/CI never
|
||||||
|
* pay the probe; the local dev launcher (run-stack.sh) opts in via THT_DWH_PRECHECK.
|
||||||
|
*/
|
||||||
|
dwhPrecheck: boolean;
|
||||||
}
|
}
|
||||||
export function loadConfig(env: Record<string, string | undefined>): AppConfig {
|
export function loadConfig(env: Record<string, string | undefined>): AppConfig {
|
||||||
const authMode = env.AUTH_MODE ?? "none";
|
const authMode = env.AUTH_MODE ?? "none";
|
||||||
@@ -99,5 +105,6 @@ export function loadConfig(env: Record<string, string | undefined>): AppConfig {
|
|||||||
secretsFile,
|
secretsFile,
|
||||||
secretFiles,
|
secretFiles,
|
||||||
modelApiKeyFile,
|
modelApiKeyFile,
|
||||||
|
dwhPrecheck: env.THT_DWH_PRECHECK === "true" || env.THT_DWH_PRECHECK === "1",
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,6 +13,8 @@ const READINESS_FAILURE_MESSAGE =
|
|||||||
"Session services are not ready. Check configuration and connectivity, then try again.";
|
"Session services are not ready. Check configuration and connectivity, then try again.";
|
||||||
const RESUME_FAILURE_MESSAGE =
|
const RESUME_FAILURE_MESSAGE =
|
||||||
"Session could not be resumed. Check configuration and connectivity, then try again.";
|
"Session could not be resumed. Check configuration and connectivity, then try again.";
|
||||||
|
const DWH_UNREACHABLE_MESSAGE =
|
||||||
|
"Cannot start a session: the data warehouse is unreachable. Check the VPN connection and try again.";
|
||||||
|
|
||||||
function eventCursor(...values: unknown[]): number {
|
function eventCursor(...values: unknown[]): number {
|
||||||
let cursor = 0;
|
let cursor = 0;
|
||||||
@@ -30,6 +32,8 @@ export function sessionRoutes(
|
|||||||
mgr: PiProcessManager; tht: ThtRunner; hub: SseHub;
|
mgr: PiProcessManager; tht: ThtRunner; hub: SseHub;
|
||||||
getSettings: (principal: PrincipalContext) => Promise<Settings>;
|
getSettings: (principal: PrincipalContext) => Promise<Settings>;
|
||||||
readiness: ReadinessManager;
|
readiness: ReadinessManager;
|
||||||
|
/** Local-only guard: probe DWH reachability before creating a session (run-stack.sh). */
|
||||||
|
dwhPrecheck?: boolean;
|
||||||
},
|
},
|
||||||
) {
|
) {
|
||||||
const lifecycleTails = new Map<string, Promise<void>>();
|
const lifecycleTails = new Map<string, Promise<void>>();
|
||||||
@@ -181,6 +185,16 @@ export function sessionRoutes(
|
|||||||
const runner = runnerFor(principal);
|
const runner = runnerFor(principal);
|
||||||
const ensure = await d.readiness.ensure(s.workspace ?? "", principal);
|
const ensure = await d.readiness.ensure(s.workspace ?? "", principal);
|
||||||
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
||||||
|
// Local-only: verify the DWH is reachable BEFORE creating the session, so a dropped
|
||||||
|
// VPN surfaces as an up-front alert instead of a session that spawns Pi and then dies
|
||||||
|
// in bootstrap retrieval. `code` lets the client show a specific message.
|
||||||
|
if (d.dwhPrecheck) {
|
||||||
|
const ping = await runner.dbPing(s.workspace);
|
||||||
|
if (!ping.ok) {
|
||||||
|
console.error(`[dwh-precheck] refusing new session — DWH unreachable: ${ping.detail}`);
|
||||||
|
return reply.code(503).send({ error: DWH_UNREACHABLE_MESSAGE, code: "dwh_unreachable" });
|
||||||
|
}
|
||||||
|
}
|
||||||
// Settings (global) supply workspace/provider/model/thinking. The new-question
|
// Settings (global) supply workspace/provider/model/thinking. The new-question
|
||||||
// form sends only the question text. `workspace` selects the tht `-c <config>`.
|
// form sends only the question text. `workspace` selects the tht `-c <config>`.
|
||||||
let id: string;
|
let id: string;
|
||||||
|
|||||||
@@ -64,7 +64,9 @@ export class ThtRunner {
|
|||||||
return [...args, ...this.configArg(workspace)];
|
return [...args, ...this.configArg(workspace)];
|
||||||
}
|
}
|
||||||
|
|
||||||
run(args: string[], workspace?: string): Promise<{ code: number; stdout: string; stderr: string }> {
|
run(
|
||||||
|
args: string[], workspace?: string, timeoutMs?: number,
|
||||||
|
): Promise<{ code: number; stdout: string; stderr: string }> {
|
||||||
return new Promise((resolve) => {
|
return new Promise((resolve) => {
|
||||||
const env: NodeJS.ProcessEnv = { ...process.env };
|
const env: NodeJS.ProcessEnv = { ...process.env };
|
||||||
delete env.THT_DATA_ROOT;
|
delete env.THT_DATA_ROOT;
|
||||||
@@ -78,11 +80,19 @@ export class ThtRunner {
|
|||||||
let stdout = "";
|
let stdout = "";
|
||||||
let stderr = "";
|
let stderr = "";
|
||||||
let settled = false;
|
let settled = false;
|
||||||
|
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||||
const finish = (result: { code: number; stdout: string; stderr: string }) => {
|
const finish = (result: { code: number; stdout: string; stderr: string }) => {
|
||||||
if (settled) return;
|
if (settled) return;
|
||||||
settled = true;
|
settled = true;
|
||||||
|
if (timer) clearTimeout(timer);
|
||||||
resolve(result);
|
resolve(result);
|
||||||
};
|
};
|
||||||
|
if (timeoutMs !== undefined) {
|
||||||
|
timer = setTimeout(() => {
|
||||||
|
try { ch.kill("SIGKILL"); } catch { /* already gone */ }
|
||||||
|
finish({ code: 124, stdout, stderr: stderr || `timed out after ${timeoutMs}ms` });
|
||||||
|
}, timeoutMs);
|
||||||
|
}
|
||||||
ch.stdout.on("data", (d: Buffer) => (stdout += d));
|
ch.stdout.on("data", (d: Buffer) => (stdout += d));
|
||||||
ch.stderr.on("data", (d: Buffer) => (stderr += d));
|
ch.stderr.on("data", (d: Buffer) => (stderr += d));
|
||||||
ch.on("error", (error) => finish({ code: 1, stdout, stderr: stderr || error.message }));
|
ch.on("error", (error) => finish({ code: 1, stdout, stderr: stderr || error.message }));
|
||||||
@@ -129,6 +139,17 @@ export class ThtRunner {
|
|||||||
if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`);
|
if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Probe DWH reachability via `tht db ping` (a REST health check). Never throws —
|
||||||
|
* returns ok=false with the failure detail so the caller can refuse a new session
|
||||||
|
* cleanly. Timed out to keep POST /sessions responsive when the host hangs.
|
||||||
|
*/
|
||||||
|
async dbPing(workspace?: string): Promise<{ ok: boolean; detail: string }> {
|
||||||
|
const { code, stdout, stderr } = await this.run(["db", "ping"], workspace, 10_000);
|
||||||
|
if (code === 0) return { ok: true, detail: stdout.trim() };
|
||||||
|
return { ok: false, detail: (stderr || stdout).trim() };
|
||||||
|
}
|
||||||
|
|
||||||
sessionList(workspace?: string) {
|
sessionList(workspace?: string) {
|
||||||
return this.json<SessionRow[]>(["session", "list", "--json"], workspace);
|
return this.json<SessionRow[]>(["session", "list", "--json"], workspace);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -178,6 +178,63 @@ test("POST /sessions usa i settings (workspace/provider/model/thinking) e crea+a
|
|||||||
unlinkSync(modelKey);
|
unlinkSync(modelKey);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("POST /sessions refuses to create a session when the local DWH precheck fails", async () => {
|
||||||
|
let sessionNewCalls = 0;
|
||||||
|
let pinged = 0;
|
||||||
|
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness", THT_DWH_PRECHECK: "1" }), {
|
||||||
|
thtRunner: {
|
||||||
|
dbPing: async () => { pinged += 1; return { ok: false, detail: "DWH non accessibile" }; },
|
||||||
|
sessionNew: async () => { sessionNewCalls += 1; return { id: "s1" }; },
|
||||||
|
searchPack: async () => {},
|
||||||
|
} as any,
|
||||||
|
readiness: { ensure: async () => ({ ok: true }) } as any,
|
||||||
|
getSettings: () => ({ workspace: "w" }) as any,
|
||||||
|
});
|
||||||
|
const res = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
||||||
|
expect(res.statusCode).toBe(503);
|
||||||
|
expect(res.json()).toMatchObject({ code: "dwh_unreachable" });
|
||||||
|
expect(pinged).toBe(1);
|
||||||
|
expect(sessionNewCalls).toBe(0); // session must not be created
|
||||||
|
});
|
||||||
|
|
||||||
|
test("POST /sessions proceeds past a passing DWH precheck", async () => {
|
||||||
|
let pinged = 0;
|
||||||
|
let created = 0;
|
||||||
|
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness", THT_DWH_PRECHECK: "1" }), {
|
||||||
|
thtRunner: {
|
||||||
|
dbPing: async () => { pinged += 1; return { ok: true, detail: "OK" }; },
|
||||||
|
sessionNew: async () => { created += 1; return { id: "s1" }; },
|
||||||
|
searchPack: async () => {},
|
||||||
|
sessionList: async () => [{ id: "s1" }],
|
||||||
|
} as any,
|
||||||
|
readiness: { ensure: async () => ({ ok: true }) } as any,
|
||||||
|
getSettings: () => ({ workspace: "w" }) as any,
|
||||||
|
spawnFn: () => nodeSpawn("node", [FAKE, SCRIPT]) as any,
|
||||||
|
});
|
||||||
|
const res = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
||||||
|
expect(res.statusCode).toBe(200);
|
||||||
|
expect(pinged).toBe(1);
|
||||||
|
expect(created).toBe(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("POST /sessions skips the DWH precheck when the flag is off (default)", async () => {
|
||||||
|
let pinged = 0;
|
||||||
|
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), {
|
||||||
|
thtRunner: {
|
||||||
|
dbPing: async () => { pinged += 1; return { ok: false, detail: "should never run" }; },
|
||||||
|
sessionNew: async () => ({ id: "s1" }),
|
||||||
|
searchPack: async () => {},
|
||||||
|
sessionList: async () => [{ id: "s1" }],
|
||||||
|
} as any,
|
||||||
|
readiness: { ensure: async () => ({ ok: true }) } as any,
|
||||||
|
getSettings: () => ({ workspace: "w" }) as any,
|
||||||
|
spawnFn: () => nodeSpawn("node", [FAKE, SCRIPT]) as any,
|
||||||
|
});
|
||||||
|
const res = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
||||||
|
expect(res.statusCode).toBe(200);
|
||||||
|
expect(pinged).toBe(0); // probe never runs without the flag
|
||||||
|
});
|
||||||
|
|
||||||
test("POST /sessions configura Pi con il thinking globale selezionato", async () => {
|
test("POST /sessions configura Pi con il thinking globale selezionato", async () => {
|
||||||
let configured: any;
|
let configured: any;
|
||||||
const bridge = { onClientEvent: () => {}, emitClientEvent: () => {} };
|
const bridge = { onClientEvent: () => {}, emitClientEvent: () => {} };
|
||||||
|
|||||||
@@ -1,5 +1,17 @@
|
|||||||
import { backendBaseUrl as BASE, joinBackendPath } from "./runtime-config";
|
import { backendBaseUrl as BASE, joinBackendPath } from "./runtime-config";
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Error thrown for non-2xx responses. `.message` stays "<status> <body>" for
|
||||||
|
* backward compatibility; `.status` and `.payload` (parsed JSON body, if any)
|
||||||
|
* let callers branch on a specific failure — e.g. a `code: "dwh_unreachable"`.
|
||||||
|
*/
|
||||||
|
export class ApiError extends Error {
|
||||||
|
constructor(readonly status: number, readonly bodyText: string, readonly payload: unknown) {
|
||||||
|
super(`${status} ${bodyText}`);
|
||||||
|
this.name = "ApiError";
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export async function apiFetch<T>(path: string, init?: RequestInit): Promise<T> {
|
export async function apiFetch<T>(path: string, init?: RequestInit): Promise<T> {
|
||||||
// Only declare a JSON content-type when we actually send a body. Body-less
|
// Only declare a JSON content-type when we actually send a body. Body-less
|
||||||
// POSTs (resume, close) would otherwise make Fastify reject the empty body
|
// POSTs (resume, close) would otherwise make Fastify reject the empty body
|
||||||
@@ -11,7 +23,12 @@ export async function apiFetch<T>(path: string, init?: RequestInit): Promise<T>
|
|||||||
headers["content-type"] = "application/json";
|
headers["content-type"] = "application/json";
|
||||||
}
|
}
|
||||||
const res = await fetch(joinBackendPath(BASE, path), { ...init, headers });
|
const res = await fetch(joinBackendPath(BASE, path), { ...init, headers });
|
||||||
if (!res.ok) throw new Error(`${res.status} ${await res.text().catch(() => "")}`);
|
if (!res.ok) {
|
||||||
|
const bodyText = await res.text().catch(() => "");
|
||||||
|
let payload: unknown;
|
||||||
|
try { payload = bodyText ? JSON.parse(bodyText) : undefined; } catch { payload = undefined; }
|
||||||
|
throw new ApiError(res.status, bodyText, payload);
|
||||||
|
}
|
||||||
if (res.status === 204) return undefined as T;
|
if (res.status === 204) return undefined as T;
|
||||||
// Accepted fire-and-forget endpoints may legitimately return 202 with no
|
// Accepted fire-and-forget endpoints may legitimately return 202 with no
|
||||||
// representation. Keep apiFetch useful for both 202 and 204 contracts.
|
// representation. Keep apiFetch useful for both 202 and 204 contracts.
|
||||||
|
|||||||
@@ -95,6 +95,32 @@ test("a failed create restores the landing view and preserves the question for r
|
|||||||
expect(FakeEventSource.instances).toHaveLength(0);
|
expect(FakeEventSource.instances).toHaveLength(0);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("a DWH-unreachable precheck shows a specific alert and preserves the question", async () => {
|
||||||
|
server.use(
|
||||||
|
http.post("http://localhost:8787/sessions", () =>
|
||||||
|
HttpResponse.json(
|
||||||
|
{
|
||||||
|
error: "Cannot start a session: the data warehouse is unreachable. Check the VPN connection and try again.",
|
||||||
|
code: "dwh_unreachable",
|
||||||
|
},
|
||||||
|
{ status: 503 },
|
||||||
|
)),
|
||||||
|
);
|
||||||
|
renderShell();
|
||||||
|
|
||||||
|
const composer = screen.getByRole("textbox", { name: /new question/i });
|
||||||
|
await userEvent.type(composer, "quanti pazienti?");
|
||||||
|
await userEvent.click(screen.getByRole("button", { name: /send/i }));
|
||||||
|
|
||||||
|
// Specific alert, not the generic "failed to create session" hint.
|
||||||
|
expect(await screen.findByText(/data warehouse is unreachable/i)).toBeInTheDocument();
|
||||||
|
expect(screen.queryByText(/failed to create session/i)).not.toBeInTheDocument();
|
||||||
|
// No session was created: landing view stays and the question is kept for retry.
|
||||||
|
expect(screen.getByText(/type your question/i)).toBeInTheDocument();
|
||||||
|
expect(composer).toHaveValue("quanti pazienti?");
|
||||||
|
expect(FakeEventSource.instances).toHaveLength(0);
|
||||||
|
});
|
||||||
|
|
||||||
test("model selector shows the three Pi-enabled models and persists the selected provider", async () => {
|
test("model selector shows the three Pi-enabled models and persists the selected provider", async () => {
|
||||||
let saved: unknown;
|
let saved: unknown;
|
||||||
server.use(
|
server.use(
|
||||||
|
|||||||
@@ -387,10 +387,10 @@ export function AppShell() {
|
|||||||
refresh();
|
refresh();
|
||||||
}
|
}
|
||||||
|
|
||||||
function failSessionCreation() {
|
function failSessionCreation(message?: string) {
|
||||||
setCreatingSession(false);
|
setCreatingSession(false);
|
||||||
resetSession();
|
resetSession();
|
||||||
toast.error("Failed to create session. Your question is ready to retry.");
|
toast.error(message ?? "Failed to create session. Your question is ready to retry.");
|
||||||
}
|
}
|
||||||
|
|
||||||
async function stopSession() {
|
async function stopSession() {
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ import { useEffect, useRef, useState } from "react";
|
|||||||
import { CornerDownLeft } from "lucide-react";
|
import { CornerDownLeft } from "lucide-react";
|
||||||
import { useQuery, useQueryClient } from "@tanstack/react-query";
|
import { useQuery, useQueryClient } from "@tanstack/react-query";
|
||||||
import { postSteer, createSession } from "../api/sessions";
|
import { postSteer, createSession } from "../api/sessions";
|
||||||
|
import { ApiError } from "../api/client";
|
||||||
import { getSettings, putSettings, type Settings } from "../api/settings";
|
import { getSettings, putSettings, type Settings } from "../api/settings";
|
||||||
import { listWorkspaces } from "../api/workspaces";
|
import { listWorkspaces } from "../api/workspaces";
|
||||||
import { listModels } from "../api/models";
|
import { listModels } from "../api/models";
|
||||||
@@ -18,7 +19,7 @@ interface Props {
|
|||||||
/** Paint the provisional session view before the create request completes. */
|
/** Paint the provisional session view before the create request completes. */
|
||||||
onSessionCreating?: (question: string) => void;
|
onSessionCreating?: (question: string) => void;
|
||||||
/** Revert the provisional view when session creation fails. */
|
/** Revert the provisional view when session creation fails. */
|
||||||
onSessionCreateFailed?: () => void;
|
onSessionCreateFailed?: (message?: string) => void;
|
||||||
/** Interrupt the running session, persisting its state. */
|
/** Interrupt the running session, persisting its state. */
|
||||||
onStop?: () => void;
|
onStop?: () => void;
|
||||||
/** Lets the parent focus the composer (e.g. on "New session"). */
|
/** Lets the parent focus the composer (e.g. on "New session"). */
|
||||||
@@ -82,9 +83,15 @@ export function SteerInput({
|
|||||||
const { id } = await createSession({ question: trimmed });
|
const { id } = await createSession({ question: trimmed });
|
||||||
onSessionCreated?.(id);
|
onSessionCreated?.(id);
|
||||||
setText("");
|
setText("");
|
||||||
} catch {
|
} catch (error) {
|
||||||
// Keep the question in the composer so retrying does not require retyping.
|
// Keep the question in the composer so retrying does not require retyping.
|
||||||
onSessionCreateFailed?.();
|
// A DWH-unreachable precheck (local dev, VPN down) carries a specific code;
|
||||||
|
// surface its message as the alert instead of the generic retry hint.
|
||||||
|
const payload = error instanceof ApiError
|
||||||
|
? (error.payload as { code?: string; error?: string } | undefined)
|
||||||
|
: undefined;
|
||||||
|
const alert = payload?.code === "dwh_unreachable" ? payload.error : undefined;
|
||||||
|
onSessionCreateFailed?.(alert);
|
||||||
} finally {
|
} finally {
|
||||||
setBusy(false);
|
setBusy(false);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -50,6 +50,7 @@ trap cleanup EXIT INT TERM
|
|||||||
THT_BIN="$THT_BIN" \
|
THT_BIN="$THT_BIN" \
|
||||||
PI_BIN="pi" \
|
PI_BIN="pi" \
|
||||||
AUTH_MODE="none" \
|
AUTH_MODE="none" \
|
||||||
|
THT_DWH_PRECHECK="1" \
|
||||||
npm run dev
|
npm run dev
|
||||||
) &
|
) &
|
||||||
pids+=($!)
|
pids+=($!)
|
||||||
|
|||||||
Reference in New Issue
Block a user