diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 86de25f8..94be96e1 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -35,6 +35,10 @@ export function sessionRoutes( string, ReturnType >(); + const runtimeManifestReaders = new WeakMap< + ReturnType, + () => Promise + >(); const failurePersistenceClaimed = new WeakSet< ReturnType >(); @@ -82,11 +86,33 @@ export function sessionRoutes( const storageFailure = (reply: any) => reply.code(503).send({ error: "session storage is unavailable" }); + const releaseIfFinalized = async ( + id: string, rt: ReturnType, + ): Promise => { + 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 => { + 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, 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 @@ -115,14 +141,7 @@ export function sessionRoutes( if (d.mgr.get(id) !== rt && boundRuntimes.get(id) === rt) { boundRuntimes.delete(id); } else if (typeof runner.sessionShow === "function") { - void withSessionLifecycle(id, async () => { - if (d.mgr.get(id) !== rt || boundRuntimes.get(id) !== rt) return; - const manifest = await runner.sessionShow(id, workspace); - if (manifest?.status !== "finalized") return; - if (d.mgr.teardownIfCurrent(id, rt) && boundRuntimes.get(id) === rt) { - boundRuntimes.delete(id); - } - }).catch((error: unknown) => { + void withSessionLifecycle(id, () => releaseIfFinalized(id, rt)).catch((error: unknown) => { console.error(`[session:${id}] terminal runtime cleanup failed:`, error); }); } @@ -189,6 +208,7 @@ export function sessionRoutes( 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 diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 5b7d503d..05fc72dc 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -181,6 +181,64 @@ test("POST /sessions usa i settings (workspace/provider/model/thinking) e crea+a unlinkSync(modelKey); }); +test("POST /sessions reclaims a finalized runtime before enforcing the Pi limit", async () => { + const runtimes = new Map(); + const statuses = new Map(); + const tornDown: string[] = []; + const order: string[] = []; + let nextId = 0; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: (id: string) => runtimes.get(id), + createFor: (id: string) => { + order.push(`create:${id}`); + if (runtimes.size >= 1) throw new Error("max Pi processes reached"); + const runtime = { bridge: { onClientEvent: () => {} } }; + runtimes.set(id, runtime); + return runtime; + }, + configure: async () => {}, + start: () => {}, + teardownIfCurrent: (id: string, expected: any) => { + if (runtimes.get(id) !== expected) return false; + runtimes.delete(id); + tornDown.push(id); + order.push(`teardown:${id}`); + return true; + }, + } as any, + thtRunner: { + sessionNew: async () => { + const id = `s${++nextId}`; + order.push(`new:${id}`); + statuses.set(id, "open"); + return { id }; + }, + sessionShow: async (id: string) => { + order.push(`show:${id}`); + return { status: statuses.get(id) }; + }, + searchPack: async () => {}, + failSession: async () => {}, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + expect((await app.inject({ method: "POST", url: "/sessions", payload: { question: "one" } })).statusCode) + .toBe(200); + statuses.set("s1", "finalized"); + order.length = 0; + + const second = await app.inject({ method: "POST", url: "/sessions", payload: { question: "two" } }); + + expect(second.statusCode).toBe(200); + expect(second.json()).toEqual({ id: "s2" }); + expect(tornDown).toEqual(["s1"]); + expect(order).toEqual(["show:s1", "teardown:s1", "new:s2", "create:s2"]); + expect(runtimes.has("s2")).toBe(true); +}); + test("POST /sessions refuses to create a session when the local DWH precheck fails", async () => { let sessionNewCalls = 0; let pinged = 0;