diff --git a/.changeset/ce-turn-admission-json-parse.md b/.changeset/ce-turn-admission-json-parse.md new file mode 100644 index 0000000000..80c7d8b81a --- /dev/null +++ b/.changeset/ce-turn-admission-json-parse.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Fix Compound Engineering sessions dying with "AI returned no valid JSON" when turns race; add retry and diagnostics. +category: fix +dev: CE orchestrator now enforces synchronous single-turn admission per session (concurrent answer/resume gets `CeTurnInProgressError`, HTTP 409) so a re-entered mobile view cannot displace the in-flight turn's live agent. The interactive AI session seam gains a second reformat retry and logs bounded raw-response snippets with provider/model via `interactiveSessionLog` on every parse failure. diff --git a/packages/engine/src/__tests__/interactive-ai-session.test.ts b/packages/engine/src/__tests__/interactive-ai-session.test.ts index 87ebd6930c..08910c0a49 100644 --- a/packages/engine/src/__tests__/interactive-ai-session.test.ts +++ b/packages/engine/src/__tests__/interactive-ai-session.test.ts @@ -8,6 +8,7 @@ import { type InteractiveAgentResult, type InteractiveAgentSession, } from "../interactive-ai-session.js"; +import { interactiveSessionLog } from "../logger.js"; import type { AgentRuntime, AgentRuntimeOptions, AgentSessionResult } from "../agent-runtime.js"; import type { AgentSession } from "@earendil-works/pi-coding-agent"; @@ -278,12 +279,16 @@ describe("interactive-ai-session seam", () => { expect(JSON.parse(lastPrompt)).toMatchObject({ type: "answer", questionId: question.id, response: answer }); }); - it("retries once on unparseable output then surfaces an error event (no hang)", async () => { - // First turn: garbage. Reformat retry: still garbage. → error. - const scripted = makeScriptedAgent(["not json at all", "still not json"]); + it("retries twice on unparseable output then surfaces an error event (no hang), logging each failed parse", async () => { + const warnSpy = vi.spyOn(interactiveSessionLog, "warn"); + const errorSpy = vi.spyOn(interactiveSessionLog, "error"); + // First turn: garbage. Both reformat retries: still garbage. → error. + const scripted = makeScriptedAgent(["not json at all", "still not json", "nope"]); const { session } = await createInteractiveAiSessionWith(factoryFor(scripted.session), { cwd: "/tmp", systemPrompt: "protocol", + defaultProvider: "openai", + defaultModelId: "gpt-5", }); await session.prompt("start"); @@ -291,14 +296,38 @@ describe("interactive-ai-session seam", () => { expect(ev.type).toBe("error"); expect(ev.type === "error" && ev.data.message).toMatch(/parse/i); - // The reformat-retry prompt was actually sent (2 prompts: initial + retry). - expect(scripted.promptCalls().length).toBe(2); + // Both reformat-retry prompts were actually sent (3 prompts: initial + 2 retries). + expect(scripted.promptCalls().length).toBe(3); + + // Every failed parse attempt logged a bounded raw-response snippet with the + // resolved provider/model, and the terminal give-up logged at error level. + expect(warnSpy).toHaveBeenCalledTimes(3); + expect(warnSpy.mock.calls[0][0]).toContain("provider=openai, model=gpt-5"); + expect(warnSpy.mock.calls[0][0]).toContain('"not json at all"'); + expect(errorSpy).toHaveBeenCalledTimes(1); + expect(errorSpy.mock.calls[0][0]).toContain("giving up after 3 parse attempts"); // Terminal: nextEvent keeps returning the error, never hangs. expect((await session.nextEvent()).type).toBe("error"); + warnSpy.mockRestore(); + errorSpy.mockRestore(); }); - it("recovers when the reformat retry produces valid JSON", async () => { + it("flags an empty assistant message distinctly in the parse-failure log", async () => { + const warnSpy = vi.spyOn(interactiveSessionLog, "warn"); + const scripted = makeScriptedAgent(["", "", ""]); + const { session } = await createInteractiveAiSessionWith(factoryFor(scripted.session), { + cwd: "/tmp", + systemPrompt: "protocol", + }); + + await session.prompt("start"); + expect((await session.nextEvent()).type).toBe("error"); + expect(warnSpy.mock.calls[0][0]).toContain("empty assistant message"); + warnSpy.mockRestore(); + }); + + it("recovers when the first reformat retry produces valid JSON", async () => { const question: PlanningQuestion = { id: "q1", type: "text", question: "?" }; const scripted = makeScriptedAgent(["garbage", q(question)]); const { session } = await createInteractiveAiSessionWith(factoryFor(scripted.session), { @@ -311,6 +340,20 @@ describe("interactive-ai-session seam", () => { expect(ev.type).toBe("question"); }); + it("recovers when only the second reformat retry produces valid JSON", async () => { + const question: PlanningQuestion = { id: "q1", type: "text", question: "?" }; + const scripted = makeScriptedAgent(["garbage", "still garbage", q(question)]); + const { session } = await createInteractiveAiSessionWith(factoryFor(scripted.session), { + cwd: "/tmp", + systemPrompt: "protocol", + }); + + await session.prompt("start"); + const ev = await session.nextEvent(); + expect(ev.type).toBe("question"); + expect(scripted.promptCalls().length).toBe(3); + }); + it("surfaces agent prompt errors as an error event without throwing", async () => { const throwing: InteractiveAgentSession = { prompt: vi.fn(async () => { diff --git a/packages/engine/src/interactive-ai-session.ts b/packages/engine/src/interactive-ai-session.ts index dda12de0f4..2be6601e2f 100644 --- a/packages/engine/src/interactive-ai-session.ts +++ b/packages/engine/src/interactive-ai-session.ts @@ -23,6 +23,7 @@ import type { } from "@fusion/core"; import type { AgentRuntime } from "./agent-runtime.js"; import { askAcpOnce } from "./cli-agent-ask.js"; +import { interactiveSessionLog } from "./logger.js"; /** Minimal shape of an agent session we depend on (subset of pi's AgentSession). */ export interface InteractiveAgentSession { @@ -55,8 +56,12 @@ export type PlanningExecutorSelection = | { kind: "model" } | { kind: "cli-agent"; runtime: AgentRuntime }; -/** One bounded reformat retry, matching planning.ts's MAX_PARSE_RETRIES. */ -const MAX_PARSE_RETRIES = 1; +/* +FNXC:InteractiveSessionParse 2026-07-23-10:40: +A CE Strategy session in the field died with "Failed to parse agent response: AI returned no valid JSON" on the opening turn (non-Anthropic default models comply less reliably with the JSON-only protocol). +Two bounded reformat retries instead of one give a non-compliant model a second corrective nudge before the session goes terminal, and every failed parse attempt logs a bounded snippet of the raw assistant text plus the resolved provider/model so support can distinguish model non-compliance from an empty message (the disposed-agent race signature) without a repro. +*/ +const MAX_PARSE_RETRIES = 2; const REFORMAT_PROMPT = "Your previous response could not be parsed as JSON. " + @@ -337,6 +342,13 @@ function extractLastAssistantText(session: InteractiveAgentSession): string { type LoopState = "idle" | "awaiting_input" | "complete" | "error"; +/** Bounded, quoted preview of an unparseable assistant response for diagnostics. */ +function describeUnparseableResponse(text: string): string { + if (!text.trim()) return "(empty assistant message — possible disposed/displaced agent or a tool-call-only final turn)"; + const preview = text.length > 400 ? `${text.slice(0, 400)}…` : text; + return JSON.stringify(preview); +} + /** * Build the interactive session over an injected agent factory. * Exported for direct (deterministic, fake-agent) testing. @@ -384,6 +396,9 @@ export async function createInteractiveAiSessionWith( break; } catch (err) { lastError = err instanceof Error ? err : new Error(String(err)); + interactiveSessionLog.warn( + `parse attempt ${attempt + 1}/${MAX_PARSE_RETRIES + 1} failed (provider=${options.defaultProvider ?? "default"}, model=${options.defaultModelId ?? "default"}): ${lastError.message} — response: ${describeUnparseableResponse(responseText)}`, + ); if (attempt < MAX_PARSE_RETRIES) { try { await agent.prompt(REFORMAT_PROMPT); @@ -398,6 +413,9 @@ export async function createInteractiveAiSessionWith( if (!parsed) { state = "error"; + interactiveSessionLog.error( + `giving up after ${MAX_PARSE_RETRIES + 1} parse attempts (provider=${options.defaultProvider ?? "default"}, model=${options.defaultModelId ?? "default"}): ${lastError?.message ?? "Unknown error"} — last response: ${describeUnparseableResponse(responseText)}`, + ); const ev: InteractiveAiSessionEvent = { type: "error", data: { message: `Failed to parse agent response: ${lastError?.message ?? "Unknown error"}`, cause: lastError }, diff --git a/packages/engine/src/logger.ts b/packages/engine/src/logger.ts index 3253d7b94a..2c08964b4a 100644 --- a/packages/engine/src/logger.ts +++ b/packages/engine/src/logger.ts @@ -132,6 +132,9 @@ export const autopilotLog = createLogger("autopilot"); /** Logger for the heartbeat execution subsystem. */ export const heartbeatLog = createLogger("heartbeat"); +/** Logger for the interactive AI session seam (planning / CE stage sessions). */ +export const interactiveSessionLog = createLogger("interactive-session"); + /** Logger for remote node runtime/client subsystems. */ export const remoteNodeLog = createLogger("remote-node"); diff --git a/plugins/fusion-plugin-compound-engineering/src/__tests__/orchestrator-turn-admission.test.ts b/plugins/fusion-plugin-compound-engineering/src/__tests__/orchestrator-turn-admission.test.ts new file mode 100644 index 0000000000..b750360bf7 --- /dev/null +++ b/plugins/fusion-plugin-compound-engineering/src/__tests__/orchestrator-turn-admission.test.ts @@ -0,0 +1,192 @@ +import { afterEach, beforeEach, expect, it, vi } from "vitest"; +import type { InteractiveAiSession, InteractiveAiSessionEvent, PlanningQuestion } from "@fusion/core"; +import { CeOrchestrator, CeTurnInProgressError } from "../session/orchestrator.js"; +import { makeHarness, pgDescribe, type TestHarness } from "./_harness.js"; + +/* +FNXC:CompoundEngineeringConcurrency 2026-07-23-11:10: +Port of the planning turn-admission invariant (FNXC:PlanningTurnAdmission, packages/dashboard/src/planning.ts 2026-07-22). +Field incident: a mobile CE Strategy session died with "Failed to parse agent response: AI returned no valid JSON" — the planning analogue was a re-entered view re-submitting a turn, displacing the in-flight turn's live agent, which then read an empty assistant message. +Invariant under test: at most one turn is admitted per CE session at a time; a concurrent answer/resume is rejected with CeTurnInProgressError and the in-flight turn settles unharmed; the reservation is released when the turn settles (including detached turns) and by explicit cancel(). +*/ + +const QUESTION: PlanningQuestion = { + id: "q1", + type: "text", + question: "Direction?", +}; + +let h: TestHarness; + +beforeEach(async () => { + h = await makeHarness(); +}); + +afterEach(() => { + h.close(); +}); + +/** + * Scripted session whose ANSWER turn blocks until the test releases it: turn 1 + * (prompt) yields the question; the answer turn parks on a gate, then yields + * complete. Lets the test hold a turn in flight deterministically. + */ +function gatedAnswerSession(): { session: InteractiveAiSession; releaseAnswer: () => void; answerCalls: () => number } { + let cursor = -1; + let release: (() => void) | undefined; + const gate = new Promise((resolve) => { + release = resolve; + }); + let answers = 0; + const session: InteractiveAiSession = { + prompt: vi.fn(async () => { + cursor++; + }), + answer: vi.fn(async () => { + answers++; + await gate; + cursor++; + }), + nextEvent: vi.fn(async (): Promise => { + return cursor === 0 + ? { type: "question", data: QUESTION } + : { type: "complete", data: { artifact: "# Done\n" } }; + }), + dispose: vi.fn(), + }; + return { session, releaseAnswer: () => release?.(), answerCalls: () => answers }; +} + +pgDescribe("turn admission (single in-flight turn per session)", () => { + it("rejects a concurrent answer with CeTurnInProgressError and the in-flight turn settles unharmed", async () => { + const gated = gatedAnswerSession(); + const orch = new CeOrchestrator({ + ctx: h.ctx, + createInteractiveAiSession: vi.fn(async () => ({ session: gated.session })), + projectRoot: h.projectRoot, + }); + + const started = await orch.start("brainstorm", { openingMessage: "kick off" }); + expect(started.session.status).toBe("awaiting_input"); + const id = started.session.id; + + // First answer holds the turn slot (its agent turn is gated open). + const first = orch.answer(id, "q1", "north"); + // Give the first entry its synchronous reservation before racing it. + await Promise.resolve(); + + // A re-submitted answer (remounted view) is rejected, NOT admitted. + await expect(orch.answer(id, "q1", "north")).rejects.toThrow(CeTurnInProgressError); + // A racing resume is rejected the same way. + await expect(orch.resume(id)).rejects.toThrow(CeTurnInProgressError); + + // The surviving turn was never displaced: it settles cleanly. + gated.releaseAnswer(); + const settled = await first; + expect(settled.event?.type).toBe("complete"); + expect(settled.session.status).toBe("completed"); + // Exactly ONE answer reached the live agent. + expect(gated.answerCalls()).toBe(1); + }); + + it("releases the reservation when the turn settles (terminal session accepts resume again)", async () => { + const gated = gatedAnswerSession(); + const orch = new CeOrchestrator({ + ctx: h.ctx, + createInteractiveAiSession: vi.fn(async () => ({ session: gated.session })), + projectRoot: h.projectRoot, + }); + + const started = await orch.start("brainstorm", { openingMessage: "kick off" }); + const id = started.session.id; + const first = orch.answer(id, "q1", "north"); + await Promise.resolve(); + gated.releaseAnswer(); + await first; + + // Reservation released → resume is admitted (completed → no-op, no conflict throw). + const resumed = await orch.resume(id); + expect(resumed.session.status).toBe("completed"); + }); + + it("holds the reservation across a DETACHED answer turn and releases on settle", async () => { + const gated = gatedAnswerSession(); + const orch = new CeOrchestrator({ + ctx: h.ctx, + createInteractiveAiSession: vi.fn(async () => ({ session: gated.session })), + projectRoot: h.projectRoot, + }); + + const started = await orch.start("brainstorm", { openingMessage: "kick off" }); + const id = started.session.id; + + // Detached answer returns immediately with the turn running in background… + const accepted = await orch.answer(id, "q1", "north", { detach: true }); + expect(accepted.session.status).toBe("active"); + + // …but the turn slot stays reserved for the whole background turn. + await expect(orch.answer(id, "q1", "north", { detach: true })).rejects.toThrow(CeTurnInProgressError); + await expect(orch.resume(id, { detach: true })).rejects.toThrow(CeTurnInProgressError); + + gated.releaseAnswer(); + await vi.waitFor(async () => { + const state = await orch.getState(id); + expect(state?.status).toBe("completed"); + }); + // Background settle released the slot. + const resumed = await orch.resume(id); + expect(resumed.session.status).toBe("completed"); + expect(gated.answerCalls()).toBe(1); + }); + + it("a rejected concurrent answer does not clear the awaiting question (recovery anchor intact)", async () => { + const gated = gatedAnswerSession(); + const orch = new CeOrchestrator({ + ctx: h.ctx, + createInteractiveAiSession: vi.fn(async () => ({ session: gated.session })), + projectRoot: h.projectRoot, + }); + + const started = await orch.start("brainstorm", { openingMessage: "kick off" }); + const id = started.session.id; + const first = orch.answer(id, "q1", "north", { detach: true }); + await first; + + // The rejected duplicate must not have mutated persisted state mid-turn. + await expect(orch.answer(id, "q1", "dupe", { detach: true })).rejects.toThrow(CeTurnInProgressError); + const state = await orch.getState(id); + expect(state?.status).toBe("active"); // still the first turn's accepted state + + gated.releaseAnswer(); + await vi.waitFor(async () => { + expect((await orch.getState(id))?.status).toBe("completed"); + }); + }); + + it("cancel() clears the reservation so a fresh turn is admitted after explicit teardown", async () => { + const gated = gatedAnswerSession(); + const orch = new CeOrchestrator({ + ctx: h.ctx, + createInteractiveAiSession: vi.fn(async () => ({ session: gated.session })), + projectRoot: h.projectRoot, + }); + + const started = await orch.start("brainstorm", { openingMessage: "kick off" }); + const id = started.session.id; + const inFlight = orch.answer(id, "q1", "north", { detach: true }); + await inFlight; + await expect(orch.resume(id)).rejects.toThrow(CeTurnInProgressError); + + // Explicit teardown: cancel interrupts the session AND frees the slot. + const cancelled = await orch.cancel(id); + expect(cancelled?.status).toBe("interrupted"); + // Resume is admitted again (no CeTurnInProgressError): the interrupted + // session (answer already cleared currentQuestion) resumes to the "active" + // retry posture. + const resumed = await orch.resume(id); + expect(resumed.session.status).toBe("active"); + + // Unblock the zombie turn so the test run doesn't leak a pending promise. + gated.releaseAnswer(); + }); +}); diff --git a/plugins/fusion-plugin-compound-engineering/src/routes/session-routes.ts b/plugins/fusion-plugin-compound-engineering/src/routes/session-routes.ts index 3541fac57d..cc408140f1 100644 --- a/plugins/fusion-plugin-compound-engineering/src/routes/session-routes.ts +++ b/plugins/fusion-plugin-compound-engineering/src/routes/session-routes.ts @@ -1,5 +1,5 @@ import type { PluginContext, PluginRouteDefinition, PluginRouteResponse } from "@fusion/core"; -import { CeOrchestrator } from "../session/orchestrator.js"; +import { CeOrchestrator, CeTurnInProgressError } from "../session/orchestrator.js"; import { recoverStaleSessionsForContext } from "../session/session-recovery.js"; import { asCeSessionStatus } from "../session/session-store.js"; import { getCePipelineStore } from "../sync/pipeline-store.js"; @@ -102,6 +102,13 @@ export function createSessionRoutes(): PluginRouteDefinition[] { const result = await orch.resume(id, { detach: true }); return { status: 200, body: { session: result.session } }; } catch (err) { + /* + * FNXC:CompoundEngineeringConcurrency 2026-07-23-11:05: + * A resume racing an in-flight turn is a retryable conflict (409), not a missing session (404) — the mobile client re-enters the view and re-fires resume while the previous turn is still settling. + */ + if (err instanceof CeTurnInProgressError) { + return { status: 409, body: { error: err.message } }; + } return { status: 404, body: { error: err instanceof Error ? err.message : String(err) } }; } }, diff --git a/plugins/fusion-plugin-compound-engineering/src/session/orchestrator.ts b/plugins/fusion-plugin-compound-engineering/src/session/orchestrator.ts index 7dedb997bc..1f3794f8cb 100644 --- a/plugins/fusion-plugin-compound-engineering/src/session/orchestrator.ts +++ b/plugins/fusion-plugin-compound-engineering/src/session/orchestrator.ts @@ -109,6 +109,21 @@ export class CeTurnTimeoutError extends Error { } } +/* +FNXC:CompoundEngineeringConcurrency 2026-07-23-10:55: +Port of the planning turn-admission invariant (see FNXC:PlanningTurnAdmission in packages/dashboard/src/planning.ts, 2026-07-22). +Mobile view unmount/remount can re-submit an answer or resume while the previous turn is still running; without a guard the second entry rehydrates or re-prompts the SAME live handle, displacing the first turn — the displaced/disposed agent then yields an empty assistant message and the surviving turn dies with "Failed to parse agent response: AI returned no valid JSON". +Invariant: at most one turn (opening / answer / resume-rehydration) may be admitted per CE session at any time, reserved SYNCHRONOUSLY before any await and held until the turn fully settles (including detached background turns). A concurrent entry is rejected with CeTurnInProgressError (routes map it to HTTP 409) instead of displacing a healthy in-flight turn. cancel()/discard() intentionally bypass the guard — they are explicit teardown, and the displaced turn settles as interrupted through the existing no-silent-loss path. +*/ +export class CeTurnInProgressError extends Error { + constructor(sessionId: string) { + super( + `A turn is already in progress for CE session ${sessionId}. Wait for it to settle (watch the live output) instead of re-submitting.`, + ); + this.name = "CeTurnInProgressError"; + } +} + /** * Which session backend a CE stage runs on (U9 seam). * @@ -292,6 +307,14 @@ export class CeOrchestrator { private readonly progressPersistence = new Map>(); /** Process-local terminal state used only when detached failure persistence itself fails. */ private readonly detachedFailureFallbacks = new Map(); + /** + * Sessions with an admitted in-flight turn (see CeTurnInProgressError), + * keyed to a per-reservation token so a stale release (from a turn that + * cancel()/discard() force-cleared) can never drop a NEWER turn's + * reservation. Membership is toggled SYNCHRONOUSLY — no await between check + * and add — so two overlapping entries cannot both pass the guard. + */ + private readonly turnReservations = new Map(); constructor(deps: OrchestratorDeps) { this.ctx = deps.ctx; @@ -491,6 +514,25 @@ export class CeOrchestrator { }; } + /** + * Synchronously reserve the session's single turn slot; returns the + * idempotent release fn. Throws CeTurnInProgressError when a turn is already + * admitted. Callers MUST release when the turn fully settles (finally on the + * sync path, promise-finally on the detached path). + */ + private reserveTurn(sessionId: string): () => void { + if (this.turnReservations.has(sessionId)) { + throw new CeTurnInProgressError(sessionId); + } + const token = Symbol(sessionId); + this.turnReservations.set(sessionId, token); + return () => { + if (this.turnReservations.get(sessionId) === token) { + this.turnReservations.delete(sessionId); + } + }; + } + /** * Start a fresh session for a registered stage and run the opening turn. * @@ -533,18 +575,34 @@ export class CeOrchestrator { ? await this.store.createWithPlanHandoffClaimAsync(sessionInput, handoffArtifactPath) : await this.store.createAsync(sessionInput); await this.store.appendHistoryAsync(session.id, { role: "user", text: opts.openingMessage, at: new Date().toISOString() }); + // The row is brand-new, so nothing can hold this reservation yet — but the + // opening turn must HOLD it so an answer/resume racing in mid-turn is + // rejected instead of displacing the live handle. + const releaseTurn = this.reserveTurn(session.id); if (opts.detach) { - let accepted = await this.requireSession(session.id); - const turn = Promise.resolve().then(() => this.runOpeningTurn( - session.id, - stage, - opts.openingMessage, - (active) => { accepted = active; }, - )); + let accepted: CeSession; + try { + accepted = await this.requireSession(session.id); + } catch (err) { + releaseTurn(); + throw err; + } + const turn = Promise.resolve() + .then(() => this.runOpeningTurn( + session.id, + stage, + opts.openingMessage, + (active) => { accepted = active; }, + )) + .finally(releaseTurn); this.detachTurn(() => accepted, "opening turn", turn); return { session: accepted }; } - return this.runOpeningTurn(session.id, stage, opts.openingMessage); + try { + return await this.runOpeningTurn(session.id, stage, opts.openingMessage); + } finally { + releaseTurn(); + } } /** Resolve the newest durable same-project Brainstorm handoff accepted for in-place Plan enrichment. */ @@ -613,35 +671,49 @@ export class CeOrchestrator { response: unknown, opts: { detach?: boolean } = {}, ): Promise { - const session = await this.requireSession(sessionId); - if (session.status !== "awaiting_input") { - throw new Error(`Session ${sessionId} is not awaiting input (status=${session.status}).`); - } - // Validate the questionId BEFORE mutating any persisted state. A stale/wrong - // questionId must NOT clear `currentQuestion` or flip status to active — - // doing so would destroy the recovery anchor while the seam rejects the - // mismatch, leaving the DB diverged from the live session. Reject cleanly and - // leave `currentQuestion`/status intact so the session stays answerable. - if (questionId !== session.currentQuestion?.id) { - throw new Error( - `Session ${sessionId} is awaiting question "${session.currentQuestion?.id ?? "(none)"}", not "${questionId}".`, - ); - } - const live = this.live.get(sessionId); - if (!live && !this.factory) { - throw new Error(INTERACTIVE_AI_UNAVAILABLE_MESSAGE); - } + // Turn admission FIRST, synchronously — a re-submitted answer from a + // remounted mobile view must be rejected before it can touch persisted + // state or the live handle (see CeTurnInProgressError). + const releaseTurn = this.reserveTurn(sessionId); + // Once the detached background turn is armed, its promise-finally owns the + // release; until then (validation throws) the finally below owns it. + let releaseHandedOff = false; + try { + const session = await this.requireSession(sessionId); + if (session.status !== "awaiting_input") { + throw new Error(`Session ${sessionId} is not awaiting input (status=${session.status}).`); + } + // Validate the questionId BEFORE mutating any persisted state. A stale/wrong + // questionId must NOT clear `currentQuestion` or flip status to active — + // doing so would destroy the recovery anchor while the seam rejects the + // mismatch, leaving the DB diverged from the live session. Reject cleanly and + // leave `currentQuestion`/status intact so the session stays answerable. + if (questionId !== session.currentQuestion?.id) { + throw new Error( + `Session ${sessionId} is awaiting question "${session.currentQuestion?.id ?? "(none)"}", not "${questionId}".`, + ); + } + const live = this.live.get(sessionId); + if (!live && !this.factory) { + throw new Error(INTERACTIVE_AI_UNAVAILABLE_MESSAGE); + } - if (opts.detach) { - // If the process lost its live handle, rehydration can take time. Mirror - // resume(detach): mark the row active immediately while the background - // turn re-creates the handle and converges through persisted state. - const accepted = await this.store.updateAsync(sessionId, { status: "active", currentQuestion: null, error: null }) ?? session; - const turn = Promise.resolve().then(() => this.runAnswerTurn(accepted, questionId, response)); - this.detachTurn(() => accepted, "answer turn", turn); - return { session: accepted }; + if (opts.detach) { + // If the process lost its live handle, rehydration can take time. Mirror + // resume(detach): mark the row active immediately while the background + // turn re-creates the handle and converges through persisted state. + const accepted = await this.store.updateAsync(sessionId, { status: "active", currentQuestion: null, error: null }) ?? session; + const turn = Promise.resolve() + .then(() => this.runAnswerTurn(accepted, questionId, response)) + .finally(releaseTurn); + releaseHandedOff = true; + this.detachTurn(() => accepted, "answer turn", turn); + return { session: accepted }; + } + return await this.runAnswerTurn(session, questionId, response); + } finally { + if (!releaseHandedOff) releaseTurn(); } - return this.runAnswerTurn(session, questionId, response); } private async runAnswerTurn(session: CeSession, questionId: string, response: unknown): Promise { @@ -691,62 +763,72 @@ export class CeOrchestrator { * left `interrupted` with a clear error explaining it can't be continued here. */ async resume(sessionId: string, opts: { detach?: boolean } = {}): Promise { - const session = await this.requireSession(sessionId); + // Turn admission FIRST, synchronously — a resume racing an in-flight + // opening/answer turn (mobile re-entry, double-tapped Resume) must not + // rehydrate a second live handle over the one the running turn is using. + const releaseTurn = this.reserveTurn(sessionId); + let releaseHandedOff = false; + try { + const session = await this.requireSession(sessionId); - // Terminal / already-answerable-with-a-live-handle cases need no rehydration. - if (session.status === "completed") return { session }; - if (session.status === "awaiting_input" && this.live.has(sessionId)) { - return { session }; // already live + answerable. - } - - // No pending question → nothing to re-prime to. Mark active so the caller can - // re-run the turn with fresh input (retry for `error`, resume for others). - if (!session.currentQuestion) { - const next = await this.store.updateAsync(sessionId, { status: "active", error: null }) ?? session; - return { session: next }; - } - - // A live handle already exists (e.g. interrupted but not disposed) — just - // restore the answerable status. - if (this.live.has(sessionId)) { - const next = await this.store.updateAsync(sessionId, { status: "awaiting_input", error: null }) ?? session; - return { session: next }; - } - - // Rehydration path: re-create the live session and replay history back to the - // current question. - if (!this.factory) { - // Honest status: we cannot back an answerable state in this process, so do - // not pretend the session is resumable here. Surface a clear error. - const next = - await this.store.updateAsync(sessionId, { - status: "interrupted", - error: INTERACTIVE_AI_UNAVAILABLE_MESSAGE, - }) ?? session; - return { session: next }; - } - - const rehydration = (async (): Promise => { - try { - await this.rehydrate(session); - } catch (err) { - // Rehydration failed — keep progress, surface the failure, do not - // advertise an answerable status we can't back. - return { session: await this.interruptSession(sessionId, err) }; + // Terminal / already-answerable-with-a-live-handle cases need no rehydration. + if (session.status === "completed") return { session }; + if (session.status === "awaiting_input" && this.live.has(sessionId)) { + return { session }; // already live + answerable. } - const next = await this.store.updateAsync(sessionId, { status: "awaiting_input", error: null }) ?? session; - return { session: next }; - })(); - if (opts.detach) { - // Rehydration replays the conversation against the live model and can be - // slow; the route posture marks the session active and converges via - // push/poll. - const next = await this.store.updateAsync(sessionId, { status: "active", error: null }) ?? session; - this.detachTurn(() => next, "rehydration", rehydration); - return { session: next }; + // No pending question → nothing to re-prime to. Mark active so the caller can + // re-run the turn with fresh input (retry for `error`, resume for others). + if (!session.currentQuestion) { + const next = await this.store.updateAsync(sessionId, { status: "active", error: null }) ?? session; + return { session: next }; + } + + // A live handle already exists (e.g. interrupted but not disposed) — just + // restore the answerable status. + if (this.live.has(sessionId)) { + const next = await this.store.updateAsync(sessionId, { status: "awaiting_input", error: null }) ?? session; + return { session: next }; + } + + // Rehydration path: re-create the live session and replay history back to the + // current question. + if (!this.factory) { + // Honest status: we cannot back an answerable state in this process, so do + // not pretend the session is resumable here. Surface a clear error. + const next = + await this.store.updateAsync(sessionId, { + status: "interrupted", + error: INTERACTIVE_AI_UNAVAILABLE_MESSAGE, + }) ?? session; + return { session: next }; + } + + const rehydration = (async (): Promise => { + try { + await this.rehydrate(session); + } catch (err) { + // Rehydration failed — keep progress, surface the failure, do not + // advertise an answerable status we can't back. + return { session: await this.interruptSession(sessionId, err) }; + } + const next = await this.store.updateAsync(sessionId, { status: "awaiting_input", error: null }) ?? session; + return { session: next }; + })().finally(releaseTurn); + releaseHandedOff = true; + + if (opts.detach) { + // Rehydration replays the conversation against the live model and can be + // slow; the route posture marks the session active and converges via + // push/poll. + const next = await this.store.updateAsync(sessionId, { status: "active", error: null }) ?? session; + this.detachTurn(() => next, "rehydration", rehydration); + return { session: next }; + } + return await rehydration; + } finally { + if (!releaseHandedOff) releaseTurn(); } - return rehydration; } /** @@ -873,6 +955,11 @@ export class CeOrchestrator { // before disposeLive clears the transient buffers (same as runTurn failure). const interrupted = await this.interruptSession(sessionId, new Error("Cancelled by user")); this.disposeLive(sessionId); + // Explicit teardown clears the turn reservation: the displaced turn's agent + // is disposed (it settles as interrupted), and the session must accept a + // fresh turn even if that turn's own release was wedged (e.g. a hung + // rehydration, which runs without a watchdog). + this.turnReservations.delete(sessionId); return interrupted; } @@ -885,6 +972,9 @@ export class CeOrchestrator { async discard(sessionId: string): Promise { await this.drainProgressPersistence(sessionId); this.disposeLive(sessionId); + // Same reservation clear as cancel(): the row is going away; never leave a + // stale reservation behind for a reused orchestrator instance. + this.turnReservations.delete(sessionId); const deleted = await this.store.deleteAsync(sessionId); if (deleted) this.detachedFailureFallbacks.delete(sessionId); return deleted;