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"; 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."; export function sessionRoutes( app: FastifyInstance, d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: (principal: PrincipalContext) => Promise; readiness: ReadinessManager; /** Local-only guard: probe DWH reachability before creating a session (run-stack.sh). */ dwhPrecheck?: boolean; }, ) { const lifecycleTails = new Map>(); const boundRuntimes = new Map< string, ReturnType >(); const failurePersistenceClaimed = new WeakSet< ReturnType >(); const withSessionLifecycle = async (id: string, work: () => Promise): Promise => { const previous = lifecycleTails.get(id) ?? Promise.resolve(); let release!: () => void; const gate = new Promise((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 => { 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 bindRuntime = ( id: string, rt: ReturnType, runner: any, workspace?: string, ) => { const previous = boundRuntimes.get(id); boundRuntimes.set(id, rt); 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" && d.mgr.get(id) !== rt && boundRuntimes.get(id) === rt ) { boundRuntimes.delete(id); } }); } 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, runner: any, workspace: string | undefined, configure: Promise, retrieval: Promise | 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); 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" }); } } // Settings (global) supply workspace/provider/model/thinking. The new-question // form sends only the question text. `workspace` selects the tht `-c `. 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, }; const rt = d.mgr.createFor(id, options); bindRuntime(id, rt, runner, s.workspace); 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 | 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 ":" 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); } }); }