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 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 () => {
|
||||
|
||||
@@ -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 },
|
||||
|
||||
@@ -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");
|
||||
|
||||
|
||||
@@ -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 { 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) } };
|
||||
}
|
||||
},
|
||||
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user