feat(backend): session routes + SSE + response/steer wiring
Make buildApp injectable (deps.thtRunner + deps.spawnFn); create
src/routes/sessions.ts with all 7 session routes under authPreHandler;
wire bridge.onClientEvent→hub.publish before returning {id} from POST
/sessions; SSE subscribes with pendingWidget safety net.
18/18 tests pass, tsc clean.
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,42 @@
|
||||
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 { getUser } from "../auth/auth.js";
|
||||
|
||||
export function sessionRoutes(app: FastifyInstance, d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub }) {
|
||||
app.post("/sessions", async (req, reply) => {
|
||||
const b = req.body as any;
|
||||
const { id } = await d.tht.sessionNew({ question: b.question, provider: b.provider, model: b.model, thinking: b.thinking, name: b.name });
|
||||
const rt = await d.mgr.spawnFor(id, { provider: b.provider, model: b.model, thinking: b.thinking, author: getUser(req).id });
|
||||
rt.bridge.onClientEvent((e) => d.hub.publish(id, e.type, e));
|
||||
return { id };
|
||||
});
|
||||
app.get("/sessions", async () => d.tht.sessionList());
|
||||
app.get("/sessions/:id", async (req) => d.tht.sessionShow((req.params as any).id));
|
||||
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/close", async (req) => {
|
||||
d.mgr.teardown((req.params as any).id);
|
||||
return { closed: true };
|
||||
});
|
||||
app.get("/sessions/:id/events", (req, reply) => {
|
||||
const id = (req.params as any).id;
|
||||
const rt = d.mgr.get(id);
|
||||
reply.raw.writeHead(200, { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", Connection: "keep-alive" });
|
||||
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);
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user