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:
7
.changeset/ce-turn-admission-json-parse.md
Normal file
7
.changeset/ce-turn-admission-json-parse.md
Normal 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.
|
||||||
@@ -8,6 +8,7 @@ import {
|
|||||||
type InteractiveAgentResult,
|
type InteractiveAgentResult,
|
||||||
type InteractiveAgentSession,
|
type InteractiveAgentSession,
|
||||||
} from "../interactive-ai-session.js";
|
} from "../interactive-ai-session.js";
|
||||||
|
import { interactiveSessionLog } from "../logger.js";
|
||||||
import type { AgentRuntime, AgentRuntimeOptions, AgentSessionResult } from "../agent-runtime.js";
|
import type { AgentRuntime, AgentRuntimeOptions, AgentSessionResult } from "../agent-runtime.js";
|
||||||
import type { AgentSession } from "@earendil-works/pi-coding-agent";
|
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 });
|
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 () => {
|
it("retries twice on unparseable output then surfaces an error event (no hang), logging each failed parse", async () => {
|
||||||
// First turn: garbage. Reformat retry: still garbage. → error.
|
const warnSpy = vi.spyOn(interactiveSessionLog, "warn");
|
||||||
const scripted = makeScriptedAgent(["not json at all", "still not json"]);
|
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), {
|
const { session } = await createInteractiveAiSessionWith(factoryFor(scripted.session), {
|
||||||
cwd: "/tmp",
|
cwd: "/tmp",
|
||||||
systemPrompt: "protocol",
|
systemPrompt: "protocol",
|
||||||
|
defaultProvider: "openai",
|
||||||
|
defaultModelId: "gpt-5",
|
||||||
});
|
});
|
||||||
|
|
||||||
await session.prompt("start");
|
await session.prompt("start");
|
||||||
@@ -291,14 +296,38 @@ describe("interactive-ai-session seam", () => {
|
|||||||
expect(ev.type).toBe("error");
|
expect(ev.type).toBe("error");
|
||||||
expect(ev.type === "error" && ev.data.message).toMatch(/parse/i);
|
expect(ev.type === "error" && ev.data.message).toMatch(/parse/i);
|
||||||
|
|
||||||
// The reformat-retry prompt was actually sent (2 prompts: initial + retry).
|
// Both reformat-retry prompts were actually sent (3 prompts: initial + 2 retries).
|
||||||
expect(scripted.promptCalls().length).toBe(2);
|
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.
|
// Terminal: nextEvent keeps returning the error, never hangs.
|
||||||
expect((await session.nextEvent()).type).toBe("error");
|
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 question: PlanningQuestion = { id: "q1", type: "text", question: "?" };
|
||||||
const scripted = makeScriptedAgent(["garbage", q(question)]);
|
const scripted = makeScriptedAgent(["garbage", q(question)]);
|
||||||
const { session } = await createInteractiveAiSessionWith(factoryFor(scripted.session), {
|
const { session } = await createInteractiveAiSessionWith(factoryFor(scripted.session), {
|
||||||
@@ -311,6 +340,20 @@ describe("interactive-ai-session seam", () => {
|
|||||||
expect(ev.type).toBe("question");
|
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 () => {
|
it("surfaces agent prompt errors as an error event without throwing", async () => {
|
||||||
const throwing: InteractiveAgentSession = {
|
const throwing: InteractiveAgentSession = {
|
||||||
prompt: vi.fn(async () => {
|
prompt: vi.fn(async () => {
|
||||||
|
|||||||
@@ -23,6 +23,7 @@ import type {
|
|||||||
} from "@fusion/core";
|
} from "@fusion/core";
|
||||||
import type { AgentRuntime } from "./agent-runtime.js";
|
import type { AgentRuntime } from "./agent-runtime.js";
|
||||||
import { askAcpOnce } from "./cli-agent-ask.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). */
|
/** Minimal shape of an agent session we depend on (subset of pi's AgentSession). */
|
||||||
export interface InteractiveAgentSession {
|
export interface InteractiveAgentSession {
|
||||||
@@ -55,8 +56,12 @@ export type PlanningExecutorSelection =
|
|||||||
| { kind: "model" }
|
| { kind: "model" }
|
||||||
| { kind: "cli-agent"; runtime: AgentRuntime };
|
| { 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 =
|
const REFORMAT_PROMPT =
|
||||||
"Your previous response could not be parsed as JSON. " +
|
"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";
|
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.
|
* Build the interactive session over an injected agent factory.
|
||||||
* Exported for direct (deterministic, fake-agent) testing.
|
* Exported for direct (deterministic, fake-agent) testing.
|
||||||
@@ -384,6 +396,9 @@ export async function createInteractiveAiSessionWith(
|
|||||||
break;
|
break;
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
lastError = err instanceof Error ? err : new Error(String(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) {
|
if (attempt < MAX_PARSE_RETRIES) {
|
||||||
try {
|
try {
|
||||||
await agent.prompt(REFORMAT_PROMPT);
|
await agent.prompt(REFORMAT_PROMPT);
|
||||||
@@ -398,6 +413,9 @@ export async function createInteractiveAiSessionWith(
|
|||||||
|
|
||||||
if (!parsed) {
|
if (!parsed) {
|
||||||
state = "error";
|
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 = {
|
const ev: InteractiveAiSessionEvent = {
|
||||||
type: "error",
|
type: "error",
|
||||||
data: { message: `Failed to parse agent response: ${lastError?.message ?? "Unknown error"}`, cause: lastError },
|
data: { message: `Failed to parse agent response: ${lastError?.message ?? "Unknown error"}`, cause: lastError },
|
||||||
|
|||||||
@@ -132,6 +132,9 @@ export const autopilotLog = createLogger("autopilot");
|
|||||||
/** Logger for the heartbeat execution subsystem. */
|
/** Logger for the heartbeat execution subsystem. */
|
||||||
export const heartbeatLog = createLogger("heartbeat");
|
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. */
|
/** Logger for remote node runtime/client subsystems. */
|
||||||
export const remoteNodeLog = createLogger("remote-node");
|
export const remoteNodeLog = createLogger("remote-node");
|
||||||
|
|
||||||
|
|||||||
@@ -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();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -1,5 +1,5 @@
|
|||||||
import type { PluginContext, PluginRouteDefinition, PluginRouteResponse } from "@fusion/core";
|
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 { recoverStaleSessionsForContext } from "../session/session-recovery.js";
|
||||||
import { asCeSessionStatus } from "../session/session-store.js";
|
import { asCeSessionStatus } from "../session/session-store.js";
|
||||||
import { getCePipelineStore } from "../sync/pipeline-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 });
|
const result = await orch.resume(id, { detach: true });
|
||||||
return { status: 200, body: { session: result.session } };
|
return { status: 200, body: { session: result.session } };
|
||||||
} catch (err) {
|
} 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) } };
|
return { status: 404, body: { error: err instanceof Error ? err.message : String(err) } };
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -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).
|
* 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>>();
|
private readonly progressPersistence = new Map<string, Promise<void>>();
|
||||||
/** Process-local terminal state used only when detached failure persistence itself fails. */
|
/** Process-local terminal state used only when detached failure persistence itself fails. */
|
||||||
private readonly detachedFailureFallbacks = new Map<string, { session: CeSession; durableFingerprint: string }>();
|
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) {
|
constructor(deps: OrchestratorDeps) {
|
||||||
this.ctx = deps.ctx;
|
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.
|
* 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.createWithPlanHandoffClaimAsync(sessionInput, handoffArtifactPath)
|
||||||
: await this.store.createAsync(sessionInput);
|
: await this.store.createAsync(sessionInput);
|
||||||
await this.store.appendHistoryAsync(session.id, { role: "user", text: opts.openingMessage, at: new Date().toISOString() });
|
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) {
|
if (opts.detach) {
|
||||||
let accepted = await this.requireSession(session.id);
|
let accepted: CeSession;
|
||||||
const turn = Promise.resolve().then(() => this.runOpeningTurn(
|
try {
|
||||||
session.id,
|
accepted = await this.requireSession(session.id);
|
||||||
stage,
|
} catch (err) {
|
||||||
opts.openingMessage,
|
releaseTurn();
|
||||||
(active) => { accepted = active; },
|
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);
|
this.detachTurn(() => accepted, "opening turn", turn);
|
||||||
return { session: accepted };
|
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. */
|
/** Resolve the newest durable same-project Brainstorm handoff accepted for in-place Plan enrichment. */
|
||||||
@@ -613,35 +671,49 @@ export class CeOrchestrator {
|
|||||||
response: unknown,
|
response: unknown,
|
||||||
opts: { detach?: boolean } = {},
|
opts: { detach?: boolean } = {},
|
||||||
): Promise<CeStepResult> {
|
): Promise<CeStepResult> {
|
||||||
const session = await this.requireSession(sessionId);
|
// Turn admission FIRST, synchronously — a re-submitted answer from a
|
||||||
if (session.status !== "awaiting_input") {
|
// remounted mobile view must be rejected before it can touch persisted
|
||||||
throw new Error(`Session ${sessionId} is not awaiting input (status=${session.status}).`);
|
// state or the live handle (see CeTurnInProgressError).
|
||||||
}
|
const releaseTurn = this.reserveTurn(sessionId);
|
||||||
// Validate the questionId BEFORE mutating any persisted state. A stale/wrong
|
// Once the detached background turn is armed, its promise-finally owns the
|
||||||
// questionId must NOT clear `currentQuestion` or flip status to active —
|
// release; until then (validation throws) the finally below owns it.
|
||||||
// doing so would destroy the recovery anchor while the seam rejects the
|
let releaseHandedOff = false;
|
||||||
// mismatch, leaving the DB diverged from the live session. Reject cleanly and
|
try {
|
||||||
// leave `currentQuestion`/status intact so the session stays answerable.
|
const session = await this.requireSession(sessionId);
|
||||||
if (questionId !== session.currentQuestion?.id) {
|
if (session.status !== "awaiting_input") {
|
||||||
throw new Error(
|
throw new Error(`Session ${sessionId} is not awaiting input (status=${session.status}).`);
|
||||||
`Session ${sessionId} is awaiting question "${session.currentQuestion?.id ?? "(none)"}", not "${questionId}".`,
|
}
|
||||||
);
|
// Validate the questionId BEFORE mutating any persisted state. A stale/wrong
|
||||||
}
|
// questionId must NOT clear `currentQuestion` or flip status to active —
|
||||||
const live = this.live.get(sessionId);
|
// doing so would destroy the recovery anchor while the seam rejects the
|
||||||
if (!live && !this.factory) {
|
// mismatch, leaving the DB diverged from the live session. Reject cleanly and
|
||||||
throw new Error(INTERACTIVE_AI_UNAVAILABLE_MESSAGE);
|
// 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 (opts.detach) {
|
||||||
// If the process lost its live handle, rehydration can take time. Mirror
|
// If the process lost its live handle, rehydration can take time. Mirror
|
||||||
// resume(detach): mark the row active immediately while the background
|
// resume(detach): mark the row active immediately while the background
|
||||||
// turn re-creates the handle and converges through persisted state.
|
// 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 accepted = await this.store.updateAsync(sessionId, { status: "active", currentQuestion: null, error: null }) ?? session;
|
||||||
const turn = Promise.resolve().then(() => this.runAnswerTurn(accepted, questionId, response));
|
const turn = Promise.resolve()
|
||||||
this.detachTurn(() => accepted, "answer turn", turn);
|
.then(() => this.runAnswerTurn(accepted, questionId, response))
|
||||||
return { session: accepted };
|
.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> {
|
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.
|
* left `interrupted` with a clear error explaining it can't be continued here.
|
||||||
*/
|
*/
|
||||||
async resume(sessionId: string, opts: { detach?: boolean } = {}): Promise<CeStepResult> {
|
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.
|
// Terminal / already-answerable-with-a-live-handle cases need no rehydration.
|
||||||
if (session.status === "completed") return { session };
|
if (session.status === "completed") return { session };
|
||||||
if (session.status === "awaiting_input" && this.live.has(sessionId)) {
|
if (session.status === "awaiting_input" && this.live.has(sessionId)) {
|
||||||
return { session }; // already live + answerable.
|
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) };
|
|
||||||
}
|
}
|
||||||
const next = await this.store.updateAsync(sessionId, { status: "awaiting_input", error: null }) ?? session;
|
|
||||||
return { session: next };
|
|
||||||
})();
|
|
||||||
|
|
||||||
if (opts.detach) {
|
// No pending question → nothing to re-prime to. Mark active so the caller can
|
||||||
// Rehydration replays the conversation against the live model and can be
|
// re-run the turn with fresh input (retry for `error`, resume for others).
|
||||||
// slow; the route posture marks the session active and converges via
|
if (!session.currentQuestion) {
|
||||||
// push/poll.
|
const next = await this.store.updateAsync(sessionId, { status: "active", error: null }) ?? session;
|
||||||
const next = await this.store.updateAsync(sessionId, { status: "active", error: null }) ?? session;
|
return { session: next };
|
||||||
this.detachTurn(() => next, "rehydration", rehydration);
|
}
|
||||||
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).
|
// before disposeLive clears the transient buffers (same as runTurn failure).
|
||||||
const interrupted = await this.interruptSession(sessionId, new Error("Cancelled by user"));
|
const interrupted = await this.interruptSession(sessionId, new Error("Cancelled by user"));
|
||||||
this.disposeLive(sessionId);
|
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;
|
return interrupted;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -885,6 +972,9 @@ export class CeOrchestrator {
|
|||||||
async discard(sessionId: string): Promise<boolean> {
|
async discard(sessionId: string): Promise<boolean> {
|
||||||
await this.drainProgressPersistence(sessionId);
|
await this.drainProgressPersistence(sessionId);
|
||||||
this.disposeLive(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);
|
const deleted = await this.store.deleteAsync(sessionId);
|
||||||
if (deleted) this.detachedFailureFallbacks.delete(sessionId);
|
if (deleted) this.detachedFailureFallbacks.delete(sessionId);
|
||||||
return deleted;
|
return deleted;
|
||||||
|
|||||||
Reference in New Issue
Block a user