Files
ThothII/backend/src/routes/sessions.ts
T

552 lines
25 KiB
TypeScript

import type { FastifyInstance } from "fastify";
import type { PiProcessManager } from "../pi/pi-process-manager.js";
import type { SessionRow, ThtRunner } from "../tht/tht-runner.js";
import type { SseHub } from "../sse/sse-hub.js";
import type { Settings } from "../settings/settings-store.js";
import { getPrincipal } from "../auth/auth.js";
import type { PrincipalContext } from "../auth/principal.js";
import type { ReadinessManager } from "../runtime/readiness-manager.js";
import type { ListModelsFn } from "./meta.js";
const BOOTSTRAP_FAILURE_MESSAGE =
"Session startup failed. Check configuration and connectivity, then Resume the session.";
const READINESS_FAILURE_MESSAGE =
"Session services are not ready. Check configuration and connectivity, then try again.";
const RESUME_FAILURE_MESSAGE =
"Session could not be resumed. Check configuration and connectivity, then try again.";
const DWH_UNREACHABLE_MESSAGE =
"Cannot start a session: the database is unreachable. Check the VPN connection and try again.";
const MODEL_UNAVAILABLE_MESSAGE =
"Selected model is unavailable. Check Pi authentication and model settings, then try again.";
export function sessionRoutes(
app: FastifyInstance,
d: {
mgr: PiProcessManager; tht: ThtRunner; hub: SseHub;
getSettings: (principal: PrincipalContext) => Promise<Settings>;
readiness: ReadinessManager;
listModels: ListModelsFn;
/** Local-only guard: probe DWH reachability before creating a session (run-stack.sh). */
dwhPrecheck?: boolean;
},
) {
const lifecycleTails = new Map<string, Promise<void>>();
const boundRuntimes = new Map<
string,
ReturnType<PiProcessManager["createFor"]>
>();
const runtimeManifestReaders = new WeakMap<
ReturnType<PiProcessManager["createFor"]>,
() => Promise<any>
>();
const failurePersistenceClaimed = new WeakSet<
ReturnType<PiProcessManager["createFor"]>
>();
const withSessionLifecycle = async <T>(id: string, work: () => Promise<T>): Promise<T> => {
const previous = lifecycleTails.get(id) ?? Promise.resolve();
let release!: () => void;
const gate = new Promise<void>((resolve) => { release = resolve; });
const tail = previous.then(() => gate);
lifecycleTails.set(id, tail);
await previous;
try {
return await work();
} finally {
release();
if (lifecycleTails.get(id) === tail) lifecycleTails.delete(id);
}
};
const info = (id: string, text: string, level = "info") =>
d.hub.publish(id, "info", { type: "info", level, text });
const runnerFor = (principal: PrincipalContext): any => {
const runner = d.tht as any;
return typeof runner.withPrincipal === "function" ? runner.withPrincipal(principal) : runner;
};
const isNotFound = (error: unknown) =>
/not found|non trovata|inesistente|404/i.test(error instanceof Error ? error.message : String(error));
/** RLS makes a foreign session indistinguishable from a missing one. */
const authorize = async (principal: PrincipalContext, id: string, workspace?: string): Promise<any | undefined> => {
try {
const runner = runnerFor(principal);
// Dependency-injected runners in legacy route tests may model only the mutation under
// test. Production ThtRunner always exposes sessionShow; keep that test seam harmless.
if (typeof runner.sessionShow !== "function") return {};
const manifest = await runner.sessionShow(id, workspace);
return manifest ?? undefined;
} catch (error) {
if (isNotFound(error)) return undefined;
throw error;
}
};
const storageFailure = (reply: any) => reply.code(503).send({ error: "session storage is unavailable" });
const releaseIfFinalized = async (
id: string, rt: ReturnType<PiProcessManager["createFor"]>,
): Promise<boolean> => {
if (d.mgr.get(id) !== rt || boundRuntimes.get(id) !== rt) return false;
const readManifest = runtimeManifestReaders.get(rt);
if (!readManifest || (await readManifest())?.status !== "finalized") return false;
if (!d.mgr.teardownIfCurrent(id, rt)) return false;
if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
return true;
};
const releaseFinalizedRuntimes = async (): Promise<void> => {
await Promise.all([...boundRuntimes.entries()].map(([id, rt]) =>
withSessionLifecycle(id, () => releaseIfFinalized(id, rt)).catch((error: unknown) => {
console.error(`[session:${id}] stale runtime cleanup failed:`, error);
}),
));
};
const bindRuntime = (
id: string, rt: ReturnType<PiProcessManager["createFor"]>, runner: any, workspace?: string,
) => {
const previous = boundRuntimes.get(id);
boundRuntimes.set(id, rt);
if (typeof runner.sessionShow === "function") {
runtimeManifestReaders.set(rt, () => runner.sessionShow(id, workspace));
}
try {
rt.bridge.onClientEvent((e) => {
// Child termination is asynchronous. Ignore queued events from a runtime once a newer
// identity is bound or the session is explicitly closed/deleted. The active identity
// remains bound after an unexpected exit so its public failure events still reach SSE.
if (boundRuntimes.get(id) !== rt) return;
if (e.type === "system_event" && e.event === "session_failed") {
if (!failurePersistenceClaimed.has(rt)) {
failurePersistenceClaimed.add(rt);
void withSessionLifecycle(id, async () => {
// agent_end can release the old binding before this queued work acquires the lock.
// Undefined means no replacement; a different identity means Resume won and the
// old failure must not touch its manifest.
const bound = boundRuntimes.get(id);
if (bound !== undefined && bound !== rt) return;
// Best-effort by design (crash path), but a storage outage must be
// visible server-side: the manifest stays open and resume remains legal.
await runner.failSession(id, workspace).catch((error: unknown) => {
console.error(`[session:${id}] failSession persistence failed:`, error);
});
}).catch(() => undefined);
}
}
d.hub.publish(id, e.type, e);
if (e.type === "system_event" && e.event === "agent_end") {
if (d.mgr.get(id) !== rt && boundRuntimes.get(id) === rt) {
boundRuntimes.delete(id);
} else if (typeof runner.sessionShow === "function") {
void withSessionLifecycle(id, () => releaseIfFinalized(id, rt)).catch((error: unknown) => {
console.error(`[session:${id}] terminal runtime cleanup failed:`, error);
});
}
}
});
} catch (error) {
if (previous) boundRuntimes.set(id, previous);
else if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
throw error;
}
};
const bootstrap = (
id: string,
rt: ReturnType<PiProcessManager["createFor"]>,
runner: any,
workspace: string | undefined,
configure: Promise<void>,
retrieval: Promise<void> | null,
start: () => void,
) => {
void (async () => {
try {
if (retrieval) info(id, "Preparing retrieval context");
await Promise.all([configure, retrieval]);
if (d.mgr.get(id) !== rt) return;
info(id, "Starting model");
if (d.mgr.get(id) !== rt) return;
start();
} catch (error) {
// The client only ever sees the generic BOOTSTRAP_FAILURE_MESSAGE; log the real
// cause server-side so failures (e.g. an unreachable DWH/vector host behind a
// dropped VPN) are diagnosable from the backend console instead of silent.
console.error(`[pi:${id}] bootstrap failed:`, error instanceof Error ? error.message : error);
void withSessionLifecycle(id, async () => {
// A bootstrap continuation can settle after Close/Delete or after a replacement was
// installed. Claim only the runtime identity that actually failed; holding the same
// lifecycle lock through persistence prevents a Resume from becoming that failure's
// accidental target.
if (d.mgr.get(id) !== rt || !d.mgr.teardownIfCurrent(id, rt)) return;
if (!failurePersistenceClaimed.has(rt)) {
failurePersistenceClaimed.add(rt);
await runner.failSession(id, workspace).catch(() => undefined);
}
rt.bridge.emitClientEvent({ type: "info", level: "error", text: BOOTSTRAP_FAILURE_MESSAGE });
rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" });
rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" });
});
}
})();
};
app.post("/runtime/prewarm", async (req, reply) => {
let settings: Settings;
try { settings = await d.getSettings(getPrincipal(req)); } catch { return storageFailure(reply); }
const workspace = settings.workspace ?? "";
void d.readiness.ensure(workspace, getPrincipal(req)).catch(() => undefined);
return reply.code(202).send({ status: "warming" });
});
app.post("/sessions", async (req, reply) => {
const b = req.body as { question: string; name?: string };
const principal = getPrincipal(req);
let s: Settings;
try { s = await d.getSettings(principal); } catch { return storageFailure(reply); }
const runner = runnerFor(principal);
await releaseFinalizedRuntimes();
const ensure = await d.readiness.ensure(s.workspace ?? "", principal);
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" });
}
}
if (s.provider && s.model) {
let available: Awaited<ReturnType<ListModelsFn>>;
try {
available = await d.listModels();
} catch {
return reply.code(503).send({
error: MODEL_UNAVAILABLE_MESSAGE,
code: "model_unavailable",
});
}
const selectedAvailable = available.some(
(candidate) => candidate.provider === s.provider && candidate.id === s.model,
);
if (!selectedAvailable) {
return reply.code(503).send({
error: MODEL_UNAVAILABLE_MESSAGE,
code: "model_unavailable",
});
}
}
// Settings (global) supply workspace/provider/model/thinking. The new-question
// form sends only the question text. `workspace` selects the tht `-c <config>`.
let id: string;
try {
({ id } = await runner.sessionNew({
question: b.question, name: b.name, workspace: s.workspace,
provider: s.provider, model: s.model, thinking: s.thinking,
}));
} catch { return storageFailure(reply); }
const options = {
provider: s.provider,
model: s.model,
thinking: s.thinking,
author: principal.displayName ?? principal.subject,
principal,
question: b.question,
};
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
try {
rt = d.mgr.createFor(id, options);
bindRuntime(id, rt, runner, s.workspace);
} catch (error) {
if (rt) d.mgr.teardownIfCurrent(id, rt);
console.error(
`[pi:${id}] runtime construction failed:`,
error instanceof Error ? error.message : "unknown error",
);
await runner.failSession(id, s.workspace).catch((persistenceError: unknown) => {
console.error(`[session:${id}] failSession persistence failed:`, persistenceError);
});
return reply.code(503).send({ error: BOOTSTRAP_FAILURE_MESSAGE });
}
info(id, "Session created");
bootstrap(
id, rt, runner, s.workspace, d.mgr.configure(rt, options),
runner.searchPack(b.question, id, s.workspace),
() => d.mgr.start(id, rt, options),
);
return { id };
});
app.get("/sessions", async (req, reply) => {
const principal = getPrincipal(req);
const scope = (req.query as { scope?: string }).scope ?? "mine";
if (scope !== "mine" && scope !== "all") return reply.code(400).send({ error: "scope must be mine or all" });
if (scope === "all" && !principal.isAdmin) return reply.code(403).send({ error: "admin scope required" });
try {
const settings = await d.getSettings(principal);
// Admin RLS is deliberately disabled for a normal 'mine' listing.
const scopedPrincipal = scope === "mine" ? { ...principal, isAdmin: false } : principal;
const list: SessionRow[] = await runnerFor(scopedPrincipal).sessionList(settings.workspace);
// Annotate each row with whether a live Pi runtime is currently bound. The client
// opens an `active` session straight into its live view (reconnecting to its pending
// gate), while a cold session keeps its explicit Resume affordance — so a mere click
// never spawns a runtime.
return list.map((row) => ({ ...row, active: d.mgr.get(row.id) !== undefined }));
} catch { return storageFailure(reply); }
});
app.get("/sessions/:id", async (req, reply) => {
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
const manifest = await authorize(principal, (req.params as any).id, settings.workspace);
return manifest ?? reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
});
app.post("/sessions/:id/response", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
const rt = d.mgr.get(id);
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
if (!rt.bridge.respond((req.body as any).ui_response)) {
return reply.code(409).send({ error: "risposta non corrispondente al gate in attesa" });
}
return reply.code(204).send();
});
app.post("/sessions/:id/steer", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
const rt = d.mgr.get(id);
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
rt.bridge.steer((req.body as any).text);
return reply.code(204).send();
});
app.post("/sessions/:id/resume", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
return withSessionLifecycle(id, async () => {
let settings: Settings;
let manifest: any;
try {
settings = await d.getSettings(principal);
manifest = await authorize(principal, id, settings.workspace);
} catch { return storageFailure(reply); }
if (!manifest) return reply.code(404).send({ error: "session not found" });
const runner = runnerFor(principal);
// Read-only contract FIRST: a finalized/archived session must refuse resume even
// when a lingering runtime still looks active — the manifest is the truth.
if (manifest?.status === "finalized" || manifest?.archived) {
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
}
// This check belongs inside the per-session lock: a preceding cold Resume may have
// installed a running runtime while this request was waiting.
const existing = d.mgr.get(id);
if (existing) {
const state = existing.bridge.turnState();
if (state === "running" || state === "waiting") {
return reply.code(200).send({ id, alreadyActive: true });
}
}
const ensure = await d.readiness.ensure(settings.workspace ?? "", principal);
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
const saved = manifest as { provider?: string; model?: string; thinking?: string } | null;
const options = {
provider: saved?.provider,
model: saved?.model,
thinking: saved?.thinking ?? settings.thinking,
author: principal.displayName ?? principal.subject,
principal,
mode: "resume" as const,
};
// Reopening is validation, not the transport commit point. Keep the old hub intact if
// persistence cannot be reopened.
try {
await runner.reopenSession(id, settings.workspace);
} catch {
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
}
// Re-check immediately before the replacement commit. A lifecycle operation that ran
// before this request acquired the lock may have changed or removed the runtime.
const current = d.mgr.get(id);
if (current) {
const state = current.bridge.turnState();
if (state === "running" || state === "waiting") {
return reply.code(200).send({ id, alreadyActive: true });
}
}
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
try {
if (current) {
if (boundRuntimes.get(id) === current) boundRuntimes.delete(id);
d.mgr.teardownIfCurrent(id, current);
}
rt = d.mgr.createFor(id, options);
bindRuntime(id, rt, runner, settings.workspace);
} catch {
// A created-but-unbound runtime is not usable. The old hub remains attached because
// clear() has not happened yet.
if (rt) {
if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
d.mgr.teardownIfCurrent(id, rt);
}
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
}
// Commit the replacement only after reopen + runtime creation/binding succeeded, and
// immediately before the first event produced by the new Resume.
d.hub.clear(id);
info(id, "Resuming session");
bootstrap(id, rt, runner, settings.workspace, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
return reply.code(200).send({ id, alreadyActive: false });
});
});
app.post("/sessions/:id/close", async (req, reply) => {
const id = (req.params as { id: string }).id;
const principal = getPrincipal(req);
return withSessionLifecycle(id, async () => {
let settings: Settings;
try {
settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
// Invalidate the live generation before persistence can yield. Otherwise its deferred
// bootstrap may start Pi while Close is already in progress.
const current = d.mgr.get(id);
boundRuntimes.delete(id);
if (current) d.mgr.teardownIfCurrent(id, current);
try {
await runnerFor(principal).closeSession(id, settings.workspace);
} catch {
return storageFailure(reply);
} finally {
// clear, NOT forget: a closed session can be reopened, and the per-session seq
// monotonicity is what keeps a browser's old cursor detectable. The buffer is
// emptied here; only delete discards the id counter.
d.hub.clear(id);
}
return { closed: true };
});
});
app.get("/sessions/:id/events", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
const rt = d.mgr.get(id);
// Add CORS headers manually: reply.raw.writeHead bypasses Fastify's onSend hooks
// (where @fastify/cors injects headers), so we must set them explicitly here.
const origin = (req.headers.origin as string | undefined) ?? "*";
reply.raw.writeHead(200, {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
Connection: "keep-alive",
"Access-Control-Allow-Origin": origin,
"Access-Control-Allow-Credentials": "true",
});
// Send the handshake immediately. Without this, Node waits for the first event body and
// proxies/clients cannot establish an idle SSE subscription or inspect its headers.
reply.raw.flushHeaders();
const send = (event: string, data: object, eventId: string) =>
reply.raw.write(`id: ${eventId}\nevent: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
const off = d.hub.subscribe(id, send, {
// Raw cursor candidates: the hub parses "<generation>:<seq>" and treats any
// other-generation (or legacy numeric) cursor as stale → replay from the start.
after: [req.headers["last-event-id"], (req.query as { lastEventId?: unknown }).lastEventId],
pending: rt?.bridge.pendingWidget() ?? null,
// clear()/forget() end every old transport so native EventSource reconnects with its
// Last-Event-ID instead of remaining attached to a subscriber callback that no longer exists.
close: () => {
if (!reply.raw.writableEnded) reply.raw.end();
},
});
req.raw.on("close", off);
});
app.post("/sessions/:id/rename", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).setName(id, (req.body as any).name, settings.workspace);
} catch { return storageFailure(reply); }
return reply.code(204).send();
});
app.post("/sessions/:id/group", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).setGroup(id, (req.body as any).group, settings.workspace);
} catch { return storageFailure(reply); }
return reply.code(204).send();
});
app.post("/sessions/:id/archive", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).archive(id, settings.workspace);
} catch { return storageFailure(reply); }
return reply.code(204).send();
});
app.post("/sessions/:id/unarchive", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
await runnerFor(principal).unarchive(id, settings.workspace);
} catch { return storageFailure(reply); }
return reply.code(204).send();
});
app.delete("/sessions/:id", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
return withSessionLifecycle(id, async () => {
let settings: Settings;
try {
settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
} catch { return storageFailure(reply); }
const current = d.mgr.get(id);
boundRuntimes.delete(id);
if (current) d.mgr.teardownIfCurrent(id, current);
try {
await runnerFor(principal).deleteSession(id, settings.workspace);
} catch {
return storageFailure(reply);
}
d.hub.forget(id);
return reply.code(204).send();
});
});
app.get("/sessions/:id/documents", async (req, reply) => {
const id = (req.params as any).id;
const principal = getPrincipal(req);
try {
const settings = await d.getSettings(principal);
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
return await runnerFor(principal).documents(id, settings.workspace);
} catch { return storageFailure(reply); }
});
}