189 lines
7.8 KiB
TypeScript
189 lines
7.8 KiB
TypeScript
import type { FastifyInstance } from "fastify";
|
|
import type { PiProcessManager } from "../pi/pi-process-manager.js";
|
|
import type { ThtRunner } from "../tht/tht-runner.js";
|
|
import type { SseHub } from "../sse/sse-hub.js";
|
|
import type { Settings } from "../settings/settings-store.js";
|
|
import { getUser } from "../auth/auth.js";
|
|
import type { ReadinessManager } from "../runtime/readiness-manager.js";
|
|
|
|
export function sessionRoutes(
|
|
app: FastifyInstance,
|
|
d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: () => Settings; readiness: ReadinessManager },
|
|
) {
|
|
const info = (id: string, text: string, level = "info") =>
|
|
d.hub.publish(id, "info", { type: "info", level, text });
|
|
|
|
const bindRuntime = (id: string, rt: ReturnType<PiProcessManager["createFor"]>) =>
|
|
rt.bridge.onClientEvent((e) => {
|
|
if (e.type === "system_event" && e.event === "session_failed") {
|
|
void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
|
|
}
|
|
d.hub.publish(id, e.type, e);
|
|
});
|
|
|
|
const bootstrap = (
|
|
id: string,
|
|
rt: ReturnType<PiProcessManager["createFor"]>,
|
|
configure: Promise<void>,
|
|
retrieval: Promise<void> | null,
|
|
start: () => void,
|
|
) => {
|
|
void (async () => {
|
|
try {
|
|
if (retrieval) info(id, "Preparing retrieval context");
|
|
await Promise.all([configure, retrieval]);
|
|
info(id, "Starting model");
|
|
start();
|
|
} catch (error) {
|
|
d.mgr.teardown(id);
|
|
const text = error instanceof Error ? error.message : String(error);
|
|
void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
|
|
rt.bridge.emitClientEvent({ type: "info", level: "error", text: `Session bootstrap failed: ${text}` });
|
|
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) => {
|
|
const workspace = d.getSettings().workspace ?? "";
|
|
void d.readiness.ensure(workspace).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 s = d.getSettings();
|
|
const ensure = await d.readiness.ensure(s.workspace ?? "");
|
|
if (!ensure.ok) return reply.code(503).send({ error: ensure.error ?? "Ollama/embeddings non disponibili" });
|
|
// Settings (global) supply workspace/provider/model/thinking. The new-question
|
|
// form sends only the question text. `workspace` selects the tht `-c <config>`.
|
|
const { id } = await d.tht.sessionNew({
|
|
question: b.question,
|
|
name: b.name,
|
|
workspace: s.workspace,
|
|
provider: s.provider,
|
|
model: s.model,
|
|
thinking: s.thinking,
|
|
});
|
|
const options = {
|
|
provider: s.provider,
|
|
model: s.model,
|
|
thinking: s.thinking,
|
|
author: getUser(req).id,
|
|
question: b.question,
|
|
};
|
|
const rt = d.mgr.createFor(id, options);
|
|
bindRuntime(id, rt);
|
|
info(id, "Session created");
|
|
bootstrap(
|
|
id, rt, d.mgr.configure(rt, options),
|
|
d.tht.searchPack(b.question, id, s.workspace),
|
|
() => d.mgr.start(id, rt, options),
|
|
);
|
|
return { id };
|
|
});
|
|
app.get("/sessions", async () => d.tht.sessionList(d.getSettings().workspace));
|
|
app.get("/sessions/:id", async (req) => d.tht.sessionShow((req.params as any).id, d.getSettings().workspace));
|
|
app.post("/sessions/:id/response", async (req, reply) => {
|
|
const id = (req.params as any).id;
|
|
const rt = d.mgr.get(id);
|
|
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
|
|
rt.bridge.respond((req.body as any).ui_response);
|
|
return reply.code(204).send();
|
|
});
|
|
app.post("/sessions/:id/steer", async (req, reply) => {
|
|
const rt = d.mgr.get((req.params as any).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 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 });
|
|
}
|
|
d.mgr.teardown(id);
|
|
d.hub.clear(id);
|
|
}
|
|
const manifest = (await d.tht.sessionShow(id, d.getSettings().workspace)) as { status?: string; archived?: boolean } | null;
|
|
if (manifest?.status === "finalized" || manifest?.archived) {
|
|
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
|
|
}
|
|
const settings = d.getSettings();
|
|
const ensure = await d.readiness.ensure(settings.workspace ?? "");
|
|
if (!ensure.ok) return reply.code(503).send({ error: ensure.error ?? "Ollama/embeddings non disponibili" });
|
|
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: getUser(req).id,
|
|
mode: "resume" as const,
|
|
};
|
|
await d.tht.reopenSession(id, settings.workspace);
|
|
const rt = d.mgr.createFor(id, options);
|
|
bindRuntime(id, rt);
|
|
info(id, "Resuming session");
|
|
bootstrap(id, rt, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
|
|
return reply.code(200).send({ id });
|
|
});
|
|
app.post("/sessions/:id/close", async (req) => {
|
|
const id = (req.params as { id: string }).id;
|
|
try {
|
|
await d.tht.closeSession(id, d.getSettings().workspace);
|
|
} finally {
|
|
d.mgr.teardown(id);
|
|
d.hub.clear(id);
|
|
}
|
|
return { closed: true };
|
|
});
|
|
app.get("/sessions/:id/events", (req, reply) => {
|
|
const id = (req.params as any).id;
|
|
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) => reply.raw.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
|
|
const off = d.hub.subscribe(id, send, rt?.bridge.pendingWidget() ?? null);
|
|
req.raw.on("close", off);
|
|
});
|
|
app.post("/sessions/:id/rename", async (req, reply) => {
|
|
await d.tht.setName((req.params as any).id, (req.body as any).name);
|
|
return reply.code(204).send();
|
|
});
|
|
app.post("/sessions/:id/group", async (req, reply) => {
|
|
await d.tht.setGroup((req.params as any).id, (req.body as any).group);
|
|
return reply.code(204).send();
|
|
});
|
|
app.post("/sessions/:id/archive", async (req, reply) => {
|
|
await d.tht.archive((req.params as any).id);
|
|
return reply.code(204).send();
|
|
});
|
|
app.post("/sessions/:id/unarchive", async (req, reply) => {
|
|
await d.tht.unarchive((req.params as any).id);
|
|
return reply.code(204).send();
|
|
});
|
|
app.delete("/sessions/:id", async (req, reply) => {
|
|
const id = (req.params as any).id;
|
|
d.mgr.teardown(id); // drop any live runtime before deleting on disk
|
|
await d.tht.deleteSession(id, d.getSettings().workspace);
|
|
return reply.code(204).send();
|
|
});
|
|
app.get("/sessions/:id/documents", async (req) => d.tht.documents((req.params as any).id));
|
|
}
|