fix: guard CE sessions against concurrent-turn displacement and harden JSON-protocol parsing

Port the planning turn-admission invariant (FNXC:PlanningTurnAdmission,
2026-07-22) into the Compound Engineering orchestrator: at most one turn
(opening/answer/resume-rehydration) is admitted per CE session, reserved
synchronously and held until the turn settles — a re-entered mobile view
re-submitting a turn now gets CeTurnInProgressError (HTTP 409) instead of
displacing the in-flight turn's live agent, which surfaced as "Failed to
parse agent response: AI returned no valid JSON". cancel()/discard()
force-clear the reservation; releases are token-scoped so a stale release
can't drop a newer turn's slot.

In the engine interactive-ai-session seam: bump the reformat retry from
one to two attempts (non-Anthropic default models comply less reliably
with the JSON-only protocol), and log every failed parse with a bounded
raw-response snippet plus resolved provider/model — including a distinct
empty-assistant-message marker — so support can diagnose these reports
without a repro.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
gsxdsm
2026-07-23 10:31:22 -07:00
parent 492d375e8f
commit 2499803c73
7 changed files with 456 additions and 96 deletions

View File

@@ -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.

View File

@@ -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 () => {

View File

@@ -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 },

View File

@@ -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");

View File

@@ -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<void>((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<InteractiveAiSessionEvent> => {
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();
});
});

View File

@@ -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) } };
}
},

View File

@@ -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<string, Promise<void>>();
/** Process-local terminal state used only when detached failure persistence itself fails. */
private readonly detachedFailureFallbacks = new Map<string, { session: CeSession; durableFingerprint: string }>();
/**
* 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<string, symbol>();
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<CeStepResult> {
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<CeStepResult> {
@@ -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<CeStepResult> {
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<CeStepResult> => {
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<CeStepResult> => {
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<boolean> {
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;