diff --git a/PROJECT_STATE.md b/PROJECT_STATE.md index 7ce3e4bc..7cb2930e 100644 --- a/PROJECT_STATE.md +++ b/PROJECT_STATE.md @@ -85,13 +85,15 @@ ThothII gira in Docker sul server co-locato, **embedded nel portale omics_portal - **Pi turns have an explicit lifecycle.** The bridge tracks `idle`, `running`, `waiting`, and `failed`; a reviewer gate is `waiting`, responses/steering return to `running`, and an - assistant provider error becomes `failed`. Provider error details are never forwarded to the - client; the UI receives a fixed sanitized recovery message. + assistant provider error or unexpected Pi child exit becomes `failed`. Provider error details + are never forwarded to the client; the UI receives a fixed sanitized recovery message. - **Resume preserves only active work.** `running`/`waiting` runtimes return as already active. - `idle`/`failed` runtimes are torn down, their SSE buffer/subscribers are cleared, and the normal - cold path restarts from the persisted provider/model/thinking with `/riprendi-sessione`. A - successful Resume of the currently selected session also closes and recreates its EventSource, - so the replacement runtime cannot be left behind an old same-ID stream. + Every validated cold path—including recovery after a child has already exited—clears stale SSE + state before reopening and restarts from persisted provider/model/thinking with + `/riprendi-sessione`; `idle`/`failed` runtimes are torn down at that point. Failed validation does + not detach the existing stream. A successful Resume of the currently selected session also + closes and recreates its EventSource, so the replacement runtime cannot be left behind an old + same-ID stream. - **Private Qwen routing is live.** Core is attached to both `omics_portal_omics_network` and external `localllm_default`; frontend remains only on the portal network. The mounted Pi profile resolves `local-qwen/qwen3.6-35b-a3b` at the sanitized base URL diff --git a/backend/src/app.ts b/backend/src/app.ts index 9ad984f8..d077b6eb 100644 --- a/backend/src/app.ts +++ b/backend/src/app.ts @@ -20,6 +20,7 @@ export interface BuildAppDeps { listModels?: ListModelsFn; getSettings?: () => Settings; readiness?: ReadinessManager; + hub?: SseHub; } export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance { @@ -39,7 +40,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc dataRoot: config.dataRoot, }); const mgr = deps?.mgr ?? new PiProcessManager(config, deps?.spawnFn ? { spawnFn: deps.spawnFn } : undefined); - const hub = new SseHub(); + const hub = deps?.hub ?? new SseHub(); const readiness = deps?.readiness ?? new ReadinessManager( tht as ThtRunner, Math.round(config.ollamaEnsureTimeoutMs / 1000), diff --git a/backend/src/bridge/session-bridge.ts b/backend/src/bridge/session-bridge.ts index e7834b73..175f319c 100644 --- a/backend/src/bridge/session-bridge.ts +++ b/backend/src/bridge/session-bridge.ts @@ -33,7 +33,7 @@ export class SessionBridge { m.message?.role === "assistant" && m.message?.stopReason === "error" ) { - this.state = "failed"; + this.markFailed(); this.fan({ type: "info", level: "error", @@ -69,6 +69,9 @@ export class SessionBridge { beginTurn(): void { this.state = "running"; } + /** Record a backend-detected failure; callers own any sanitized client message. */ + markFailed(): void { this.state = "failed"; } + onClientEvent(cb: (e: ClientEvent) => void): void { this.cbs.add(cb); } /** Eventi generati dal backend stesso (es. exit inatteso del child Pi), non da Pi. */ diff --git a/backend/src/pi/pi-process-manager.ts b/backend/src/pi/pi-process-manager.ts index 7f6af95b..6933f079 100644 --- a/backend/src/pi/pi-process-manager.ts +++ b/backend/src/pi/pi-process-manager.ts @@ -103,6 +103,7 @@ export class PiProcessManager { child.on("exit", (code) => { console.error(`[pi:${sessionId}] exited code=${code ?? "?"} mapped=${this.runtimes.get(sessionId) === rt}`); if (this.runtimes.get(sessionId) === rt) { + rt.bridge.markFailed(); this.runtimes.delete(sessionId); rt.bridge.emitClientEvent({ type: "info", diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index ab0c7295..08f5d019 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -106,8 +106,6 @@ export function sessionRoutes( 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) { @@ -124,6 +122,8 @@ export function sessionRoutes( author: getUser(req).id, mode: "resume" as const, }; + if (existing) d.mgr.teardown(id); + d.hub.clear(id); await d.tht.reopenSession(id, settings.workspace); const rt = d.mgr.createFor(id, options); bindRuntime(id, rt); diff --git a/backend/test/pi-process-manager.test.ts b/backend/test/pi-process-manager.test.ts index f10d6ce8..b6d88466 100644 --- a/backend/test/pi-process-manager.test.ts +++ b/backend/test/pi-process-manager.test.ts @@ -92,6 +92,7 @@ test("un exit INATTESO del child notifica il client (info error + agent_end)", a const seen: any[] = []; rt.bridge.onClientEvent((e) => seen.push(e)); child.emit("exit", 137); + expect(rt.bridge.turnState()).toBe("failed"); expect(seen).toEqual([ { type: "info", level: "error", text: expect.stringContaining("137") }, { type: "system_event", event: "session_failed" }, diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 9dd9bf4e..48cd03a9 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -137,18 +137,21 @@ test.each(["running", "waiting"])( "POST resume preserves a %s runtime", async (state) => { let tornDown = false; + let cleared = false; const existing = { bridge: { turnState: () => state } } as any; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr: { get: () => existing, teardown: () => { tornDown = true; }, } as any, + hub: { clear: () => { cleared = true; } } as any, thtRunner: {} as any, getSettings: () => ({ workspace: "psd" }) as any, }); const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); expect(response.json()).toEqual({ id: "s1", alreadyActive: true }); expect(tornDown).toBe(false); + expect(cleared).toBe(false); }, ); @@ -166,6 +169,10 @@ test.each(["idle", "failed"])( configure: async () => {}, start: () => order.push("start"), } as any, + hub: { + clear: (id: string) => order.push(`clear:${id}`), + publish: () => {}, + } as any, thtRunner: { sessionShow: async () => ({ status: "open", archived: false, @@ -179,10 +186,95 @@ test.each(["idle", "failed"])( const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); await new Promise((resolve) => setImmediate(resolve)); expect(response.json()).toEqual({ id: "s1" }); - expect(order).toEqual(["teardown:s1", "reopen", "create", "start"]); + expect(order).toEqual(["teardown:s1", "clear:s1", "reopen", "create", "start"]); }, ); +test("POST resume without a runtime clears stale SSE state before cold start", async () => { + const order: string[] = []; + let createOptions: any; + const newRuntime = { bridge: { onClientEvent: () => {}, emitClientEvent: () => {} } } as any; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => undefined, + createFor: (_id: string, options: any) => { + createOptions = options; + order.push("create"); + return newRuntime; + }, + configure: async () => {}, + start: () => order.push("start"), + } as any, + hub: { + clear: (id: string) => order.push(`clear:${id}`), + publish: () => {}, + } as any, + thtRunner: { + sessionShow: async () => ({ + status: "open", archived: false, + provider: "local-qwen", model: "qwen3.6-35b-a3b", thinking: "low", + }), + reopenSession: async () => order.push("reopen"), + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local", thinking: "medium" }) as any, + }); + + const response = await app.inject({ method: "POST", url: "/sessions/crashed/resume" }); + await new Promise((resolve) => setImmediate(resolve)); + + expect(response.json()).toEqual({ id: "crashed" }); + expect(order).toEqual(["clear:crashed", "reopen", "create", "start"]); + expect(createOptions).toMatchObject({ + provider: "local-qwen", model: "qwen3.6-35b-a3b", thinking: "low", mode: "resume", + }); +}); + +test("POST resume keeps an idle runtime stream attached when the manifest is read-only", async () => { + let tornDown = false; + let cleared = false; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => ({ bridge: { turnState: () => "idle" } }), + teardown: () => { tornDown = true; }, + } as any, + hub: { clear: () => { cleared = true; } } as any, + thtRunner: { + sessionShow: async () => ({ status: "finalized", archived: false }), + } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); + + expect(response.statusCode).toBe(409); + expect(tornDown).toBe(false); + expect(cleared).toBe(false); +}); + +test("POST resume keeps a failed runtime stream attached when readiness fails", async () => { + let tornDown = false; + let cleared = false; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => ({ bridge: { turnState: () => "failed" } }), + teardown: () => { tornDown = true; }, + } as any, + hub: { clear: () => { cleared = true; } } as any, + thtRunner: { + sessionShow: async () => ({ status: "open", archived: false }), + } as any, + readiness: { ensure: async () => ({ ok: false, error: "not ready" }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); + + expect(response.statusCode).toBe(503); + expect(tornDown).toBe(false); + expect(cleared).toBe(false); +}); + test("POST /sessions/:id/response inoltra al bridge (no error)", async () => { const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { diff --git a/backend/test/session-bridge.test.ts b/backend/test/session-bridge.test.ts index 7069a1e6..693d0763 100644 --- a/backend/test/session-bridge.test.ts +++ b/backend/test/session-bridge.test.ts @@ -68,6 +68,19 @@ test("assistant provider errors are sanitized and leave the turn failed", () => expect(JSON.stringify(seen)).not.toContain("DO_NOT_LEAK"); }); +test("markFailed records backend-detected failure without emitting raw detail", () => { + const { rpc } = fakeRpc(); + const bridge = new SessionBridge(rpc); + const seen: any[] = []; + bridge.onClientEvent((event) => seen.push(event)); + bridge.beginTurn(); + + bridge.markFailed(); + + expect(bridge.turnState()).toBe("failed"); + expect(seen).toEqual([]); +}); + test("reviewer wait and response transition waiting back to running", () => { const { rpc, fire } = fakeRpc(); const bridge = new SessionBridge(rpc); diff --git a/brain/codebase/workflow-ui-contracts.md b/brain/codebase/workflow-ui-contracts.md index e065f650..6b3ef000 100644 --- a/brain/codebase/workflow-ui-contracts.md +++ b/brain/codebase/workflow-ui-contracts.md @@ -14,10 +14,11 @@ - `SqlViewer`'s horizontal/vertical layout control is meaningful only with multiple SQL blocks; hide it for the single CTE result shown by `CteResultViewer`. - A Pi turn is `idle`, `running`, `waiting`, or `failed`. A reviewer gate/request moves it to - `waiting`; the reviewer response and steering move it back to `running`; provider failures - remain `failed` after `agent_end` and emit only the fixed sanitized recovery message. -- Resume preserves `running`/`waiting` runtimes. It replaces `idle`/`failed` runtimes, clears the - old SSE buffer/subscribers before cold resume, and reuses the persisted provider/model/thinking. + `waiting`; the reviewer response and steering move it back to `running`; provider failures and + unexpected Pi child exits mark it `failed` without forwarding raw failure detail. +- Resume preserves `running`/`waiting` runtimes. After manifest/readiness validation, every cold + path clears old SSE buffer/subscribers before reopen—including when a crashed child has already + left no runtime—and reuses persisted provider/model/thinking. Failed validation does not clear. - A successful Resume of the already selected session increments the stream generation so React closes the old EventSource and opens the same session URL again. Failed Resume must not reconnect. - SSE endpoints are intentionally keep-alive. Browser cleanup and one-off probes must explicitly diff --git a/compose.yaml b/compose.yaml index 1ffd5e30..92b9c700 100644 --- a/compose.yaml +++ b/compose.yaml @@ -4,7 +4,9 @@ # per il DNS usato dagli upstream nginx del portale. # NESSUNA porta host esposta: il backend è invisibile dall'esterno. # -# Prereq: il portale deve essere up (crea la rete): +# Prereq: devono esistere entrambe le reti esterne: il portale crea +# omics_portal_omics_network e lo stack vLLM crea localllm_default. +# Avvia il portale con: # cd /home/chirone/omics_portal && docker compose up -d # Poi: docker compose up -d --build name: thothii diff --git a/docs/superpowers/plans/2026-07-14-qwen-connectivity-resume-recovery.md b/docs/superpowers/plans/2026-07-14-qwen-connectivity-resume-recovery.md index 3f376fe3..81bc617e 100644 --- a/docs/superpowers/plans/2026-07-14-qwen-connectivity-resume-recovery.md +++ b/docs/superpowers/plans/2026-07-14-qwen-connectivity-resume-recovery.md @@ -637,6 +637,7 @@ async function waitForGate(id) { } } finally { clearTimeout(timer); + controller.abort(); } } diff --git a/frontend/src/shell/AppShell.session-mgmt.test.tsx b/frontend/src/shell/AppShell.session-mgmt.test.tsx index 5ae6fb2a..efe09953 100644 --- a/frontend/src/shell/AppShell.session-mgmt.test.tsx +++ b/frontend/src/shell/AppShell.session-mgmt.test.tsx @@ -116,6 +116,86 @@ test("resuming the active session reconnects its EventSource", async () => { expect(first.closed).toBe(true); }); +test("same-session Resume reconnects only after the deferred POST succeeds", async () => { + server.use( + http.post("http://localhost:8787/sessions/:id/resume", () => + new HttpResponse(null, { status: 204 })), + http.get("http://localhost:8787/sessions/:id", () => + HttpResponse.json({ id: "s1", status: "open", phase: 1 })), + ); + wrap(); + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); + const first = FakeEventSource.instances[0]; + + let releaseResume!: () => void; + let markStarted!: () => void; + const resumeStarted = new Promise((resolve) => { markStarted = resolve; }); + const resumeReleased = new Promise((resolve) => { releaseResume = resolve; }); + server.use(http.post("http://localhost:8787/sessions/:id/resume", async () => { + markStarted(); + await resumeReleased; + return new HttpResponse(null, { status: 204 }); + })); + + await userEvent.click(screen.getByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await resumeStarted; + expect(FakeEventSource.instances).toHaveLength(1); + expect(first.closed).toBe(false); + + releaseResume(); + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(2)); + expect(first.closed).toBe(true); +}); + +test("a failed same-session Resume does not create a replacement EventSource", async () => { + server.use( + http.post("http://localhost:8787/sessions/:id/resume", () => + new HttpResponse(null, { status: 204 })), + http.get("http://localhost:8787/sessions/:id", () => + HttpResponse.json({ id: "s1", status: "open", phase: 1 })), + ); + wrap(); + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); + + server.use(http.post("http://localhost:8787/sessions/:id/resume", () => + new HttpResponse(null, { status: 409 }))); + await userEvent.click(screen.getByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(screen.getByText("Domanda originale")).toBeInTheDocument()); + + expect(FakeEventSource.instances).toHaveLength(1); +}); + +test("resuming a different session opens exactly one new EventSource", async () => { + const other = { + ...LIST[0], id: "s3", question: "Attiva tre", group: null, + created_at: "2026-01-03T00:00:00Z", + }; + server.use( + http.get("http://localhost:8787/sessions", () => HttpResponse.json([LIST[0], other])), + http.post("http://localhost:8787/sessions/:id/resume", () => + new HttpResponse(null, { status: 204 })), + http.get("http://localhost:8787/sessions/:id", ({ params }) => + HttpResponse.json({ id: params.id, status: "open", phase: params.id === "s3" ? 3 : 1 })), + ); + wrap(); + await userEvent.click(await screen.findByText("Attiva uno")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); + + await userEvent.click(screen.getByText("Attiva tre")); + await userEvent.click(await screen.findByRole("button", { name: /resume/i })); + await waitFor(() => expect(useSessionStore.getState().currentPhase).toBe("F3")); + + expect(FakeEventSource.instances).toHaveLength(2); + expect(FakeEventSource.instances[1].url).toContain("/sessions/s3/events"); +}); + test("a failed resume keeps the panel open and does not activate the session", async () => { server.use(http.post("http://localhost:8787/sessions/:id/resume", () => new HttpResponse(null, { status: 409 }))); wrap(); diff --git a/scripts/test-qwen-network-config.sh b/scripts/test-qwen-network-config.sh index 1c42339d..a40352b1 100755 --- a/scripts/test-qwen-network-config.sh +++ b/scripts/test-qwen-network-config.sh @@ -8,7 +8,7 @@ mkdir -p "$tmp/deploy" cp compose.yaml "$tmp/compose.yaml" : >"$tmp/deploy/thothii.env" -docker compose --project-directory "$tmp" config --format json >"$tmp/config.json" +docker compose --project-directory "$tmp" -f "$tmp/compose.yaml" config --format json >"$tmp/config.json" node - "$tmp/config.json" <<'NODE' const fs = require("fs"); const config = JSON.parse(fs.readFileSync(process.argv[2], "utf8"));