FN-5802: prevent planning-session double-submit freeze
Eliminate the double-submit race so planning sessions no longer get stuck in generating. - add concurrency guards in planning session state transitions to reject overlapping submits - tighten planning route handling to avoid duplicate generation kicks for the same session - expand planning tests to cover race scenarios and ensure session recovery behavior Files changed: packages/dashboard/src/__tests__/planning.test.ts | 141 +++++++++++++++++++++ packages/dashboard/src/planning.ts | 92 +++++++++++--- .../src/routes/register-planning-subtask-routes.ts | 4 +- 3 files changed, 219 insertions(+), 18 deletions(-) Fusion-Task-Id: FN-5802 Fusion-Task-Lineage: 0ae4d04a-319f-4c93-9855-214d9292f259
This commit is contained in:
@@ -27,11 +27,14 @@ import {
|
||||
__setCreateFnAgent,
|
||||
__setPlanningDiagnostics,
|
||||
__setPlanningNtfyHelpers,
|
||||
__getActiveGenerationForTests,
|
||||
__runGenerationWithTimeoutForTests,
|
||||
rehydrateFromStore,
|
||||
setAiSessionStore,
|
||||
RateLimitError,
|
||||
SessionNotFoundError,
|
||||
InvalidSessionStateError,
|
||||
GenerationInProgressError,
|
||||
parseAgentResponse,
|
||||
buildDepthPromptSuffix,
|
||||
generateSubtasksFromPlanning,
|
||||
@@ -835,6 +838,43 @@ describe("planning module", () => {
|
||||
});
|
||||
|
||||
describe("submitResponse", () => {
|
||||
it("rejects overlapping submit for same question and keeps one history entry", async () => {
|
||||
const mockIp = getUniqueIp();
|
||||
const { sessionId } = await createSession(mockIp, initialPlan, MOCK_TASK_STORE, TEST_ROOT_DIR);
|
||||
const session = getSession(sessionId);
|
||||
expect(session?.currentQuestion?.id).toBe("q-scope");
|
||||
expect(session?.agent).toBeDefined();
|
||||
|
||||
let releasePrompt: (() => void) | undefined;
|
||||
const promptMock = vi.fn(
|
||||
(_message: string, options?: { signal?: AbortSignal }) =>
|
||||
new Promise<void>((resolve) => {
|
||||
expect(options?.signal).toBeDefined();
|
||||
releasePrompt = resolve;
|
||||
}),
|
||||
);
|
||||
|
||||
if (!session?.agent) {
|
||||
throw new Error("Expected session agent");
|
||||
}
|
||||
session.agent.session.prompt = promptMock as any;
|
||||
|
||||
const firstSubmit = submitResponse(sessionId, { "q-scope": "medium" }, TEST_ROOT_DIR);
|
||||
await vi.waitFor(() => {
|
||||
expect(promptMock).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
await expect(submitResponse(sessionId, { "q-scope": "medium" }, TEST_ROOT_DIR)).rejects.toThrow(
|
||||
GenerationInProgressError,
|
||||
);
|
||||
|
||||
expect(getSession(sessionId)?.history).toHaveLength(0);
|
||||
releasePrompt?.();
|
||||
const firstResponse = await firstSubmit;
|
||||
expect(firstResponse.type).toBe("question");
|
||||
expect(getSession(sessionId)?.history).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("processes response and returns next question", async () => {
|
||||
const mockIp = getUniqueIp();
|
||||
const { sessionId } = await createSession(mockIp, initialPlan, MOCK_TASK_STORE, TEST_ROOT_DIR);
|
||||
@@ -1341,6 +1381,71 @@ describe("planning module", () => {
|
||||
});
|
||||
|
||||
describe("generation controls", () => {
|
||||
it("older generation cleanup does not remove newer active entry", async () => {
|
||||
const { sessionId } = await createSession(getUniqueIp(), initialPlan, MOCK_TASK_STORE, TEST_ROOT_DIR);
|
||||
|
||||
let resolveFirst: (() => void) | undefined;
|
||||
const firstGeneration = __runGenerationWithTimeoutForTests(sessionId, async () =>
|
||||
new Promise<void>((resolve) => {
|
||||
resolveFirst = resolve;
|
||||
}),
|
||||
);
|
||||
|
||||
await vi.waitFor(() => {
|
||||
expect(__getActiveGenerationForTests(sessionId)).toBeDefined();
|
||||
});
|
||||
const firstRecord = __getActiveGenerationForTests(sessionId);
|
||||
|
||||
let resolveSecond: (() => void) | undefined;
|
||||
const secondGeneration = __runGenerationWithTimeoutForTests(sessionId, async () =>
|
||||
new Promise<void>((resolve) => {
|
||||
resolveSecond = resolve;
|
||||
}),
|
||||
);
|
||||
|
||||
await vi.waitFor(() => {
|
||||
const current = __getActiveGenerationForTests(sessionId);
|
||||
expect(current).toBeDefined();
|
||||
expect(current).not.toBe(firstRecord);
|
||||
});
|
||||
const secondRecord = __getActiveGenerationForTests(sessionId);
|
||||
|
||||
await expect(firstGeneration).rejects.toThrow("Generation aborted");
|
||||
expect(__getActiveGenerationForTests(sessionId)).toBe(secondRecord);
|
||||
|
||||
resolveSecond?.();
|
||||
await secondGeneration;
|
||||
expect(__getActiveGenerationForTests(sessionId)).toBeUndefined();
|
||||
resolveFirst?.();
|
||||
});
|
||||
|
||||
it("timeout path never leaves persisted session in generating", async () => {
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const store = new MockAiSessionStore();
|
||||
setAiSessionStore(store as any);
|
||||
|
||||
const hangingAgent = {
|
||||
session: {
|
||||
state: { messages: [] as Array<{ role: string; content: string }> },
|
||||
prompt: vi.fn(() => new Promise<void>(() => {})),
|
||||
dispose: vi.fn(),
|
||||
},
|
||||
};
|
||||
__setCreateFnAgent(async () => hangingAgent as any);
|
||||
|
||||
const sessionId = await createSessionWithAgent(getUniqueIp(), initialPlan, TEST_ROOT_DIR, MOCK_TASK_STORE);
|
||||
|
||||
await vi.advanceTimersByTimeAsync(GENERATION_TIMEOUT_MS + 10);
|
||||
await flushAsyncWork();
|
||||
|
||||
expect(getSession(sessionId)?.error).toContain("timed out");
|
||||
expect(store.rows.get(sessionId)?.status).toBe("error");
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it("returns false when stopping unknown session", () => {
|
||||
expect(stopGeneration("missing-session")).toBe(false);
|
||||
});
|
||||
@@ -1376,6 +1481,42 @@ describe("planning module", () => {
|
||||
resolvePrompt?.();
|
||||
});
|
||||
|
||||
it("does not append history when generation is aborted", async () => {
|
||||
let resolvePrompt: (() => void) | undefined;
|
||||
const hangingAgent = {
|
||||
session: {
|
||||
state: { messages: [] as Array<{ role: string; content: string }> },
|
||||
prompt: vi.fn(
|
||||
() =>
|
||||
new Promise<void>((resolve) => {
|
||||
resolvePrompt = resolve;
|
||||
}),
|
||||
),
|
||||
dispose: vi.fn(),
|
||||
},
|
||||
};
|
||||
__setCreateFnAgent(async () => hangingAgent as any);
|
||||
|
||||
const sessionId = await createSessionWithAgent(getUniqueIp(), initialPlan, TEST_ROOT_DIR, MOCK_TASK_STORE);
|
||||
await vi.waitFor(() => {
|
||||
expect(hangingAgent.session.prompt).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
const submitPromise = submitResponse(sessionId, { "q-scope": "medium" }, TEST_ROOT_DIR);
|
||||
await vi.waitFor(() => {
|
||||
expect(hangingAgent.session.prompt).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
expect(getSession(sessionId)?.history).toHaveLength(0);
|
||||
expect(stopGeneration(sessionId)).toBe(true);
|
||||
|
||||
const response = await submitPromise;
|
||||
expect(response.type).toBe("question");
|
||||
expect(getSession(sessionId)?.history).toHaveLength(0);
|
||||
|
||||
resolvePrompt?.();
|
||||
});
|
||||
|
||||
it("times out stalled generation and transitions session to error", async () => {
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
|
||||
@@ -1496,7 +1496,7 @@ function createAbortError(): Error {
|
||||
return error;
|
||||
}
|
||||
|
||||
async function runGenerationWithTimeout<T>(session: Session, operation: () => Promise<T>): Promise<T> {
|
||||
async function runGenerationWithTimeout<T>(session: Session, operation: (abortSignal: AbortSignal) => Promise<T>): Promise<T> {
|
||||
const existing = activeGenerations.get(session.id);
|
||||
if (existing) {
|
||||
clearTimeout(existing.timer);
|
||||
@@ -1510,8 +1510,9 @@ async function runGenerationWithTimeout<T>(session: Session, operation: () => Pr
|
||||
setSessionError(session, "AI generation timed out. You can retry or start a new session.");
|
||||
abortController.abort();
|
||||
}, GENERATION_TIMEOUT_MS);
|
||||
const generationRecord = { abortController, timer };
|
||||
|
||||
activeGenerations.set(session.id, { abortController, timer });
|
||||
activeGenerations.set(session.id, generationRecord);
|
||||
|
||||
const abortPromise = new Promise<never>((_, reject) => {
|
||||
abortController.signal.addEventListener(
|
||||
@@ -1522,7 +1523,7 @@ async function runGenerationWithTimeout<T>(session: Session, operation: () => Pr
|
||||
});
|
||||
|
||||
try {
|
||||
return await Promise.race([operation(), abortPromise]);
|
||||
return await Promise.race([operation(abortController.signal), abortPromise]);
|
||||
} catch (error) {
|
||||
if (error instanceof Error && error.name === "AbortError") {
|
||||
if (!timeoutTriggered && !session.error) {
|
||||
@@ -1532,7 +1533,9 @@ async function runGenerationWithTimeout<T>(session: Session, operation: () => Pr
|
||||
throw error;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
activeGenerations.delete(session.id);
|
||||
if (activeGenerations.get(session.id) === generationRecord) {
|
||||
activeGenerations.delete(session.id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1542,12 +1545,15 @@ async function continueAgentConversation(session: Session, message: string): Pro
|
||||
}
|
||||
|
||||
try {
|
||||
await runGenerationWithTimeout(session, async () => {
|
||||
await runGenerationWithTimeout(session, async (abortSignal) => {
|
||||
// Clear thinking output for this turn
|
||||
session.thinkingOutput = "";
|
||||
|
||||
// Send message to agent using .prompt() - it will stream thinking via onThinking callback
|
||||
await session.agent.session.prompt(message);
|
||||
// Send message to agent using .prompt() - it will stream thinking via onThinking callback.
|
||||
// Pass abort signal so timeout/user-stop can cancel the underlying prompt when supported.
|
||||
await (session.agent.session.prompt as (input: string, options?: { signal?: AbortSignal }) => Promise<void>)(message, {
|
||||
signal: abortSignal,
|
||||
});
|
||||
|
||||
// Get the response text from the agent's state
|
||||
interface AgentMessage {
|
||||
@@ -1611,10 +1617,11 @@ async function continueAgentConversation(session: Session, message: string): Pro
|
||||
);
|
||||
try {
|
||||
session.thinkingOutput = "";
|
||||
await session.agent.session.prompt(
|
||||
await (session.agent.session.prompt as (input: string, options?: { signal?: AbortSignal }) => Promise<void>)(
|
||||
"Your previous response could not be parsed as JSON. " +
|
||||
'Please respond with ONLY a valid JSON object: either {"type":"question","data":{...}} ' +
|
||||
'or {"type":"complete","data":{...}}. No markdown, no explanation, just the JSON.'
|
||||
'Please respond with ONLY a valid JSON object: either {"type":"question","data":{...}} ' +
|
||||
'or {"type":"complete","data":{...}}. No markdown, no explanation, just the JSON.',
|
||||
{ signal: abortSignal },
|
||||
);
|
||||
|
||||
// Get the new response text
|
||||
@@ -1908,6 +1915,20 @@ function formatRefineRequestForAgent(summary: PlanningSummary): string {
|
||||
].join("\n\n");
|
||||
}
|
||||
|
||||
function didSubmitSameAnswer(
|
||||
session: Session,
|
||||
responses: Record<string, unknown>,
|
||||
): boolean {
|
||||
if (!session.currentQuestion) {
|
||||
return false;
|
||||
}
|
||||
const lastEntry = session.history[session.history.length - 1];
|
||||
if (!lastEntry || lastEntry.question.id !== session.currentQuestion.id) {
|
||||
return false;
|
||||
}
|
||||
return JSON.stringify(lastEntry.response) === JSON.stringify(responses);
|
||||
}
|
||||
|
||||
export async function submitResponse(
|
||||
sessionId: string,
|
||||
responses: Record<string, unknown>,
|
||||
@@ -1926,6 +1947,13 @@ export async function submitResponse(
|
||||
if (store && !session.store) session.store = store;
|
||||
if (rootDir && !session.rootDir) session.rootDir = rootDir;
|
||||
|
||||
if (activeGenerations.has(session.id)) {
|
||||
if (didSubmitSameAnswer(session, responses)) {
|
||||
throw new GenerationInProgressError("Generation already in progress for this response");
|
||||
}
|
||||
throw new GenerationInProgressError("Generation already in progress");
|
||||
}
|
||||
|
||||
if (!session.currentQuestion) {
|
||||
if (!isRefineRequest(responses) || !session.summary) {
|
||||
throw new InvalidSessionStateError("No active question in session");
|
||||
@@ -1938,22 +1966,26 @@ export async function submitResponse(
|
||||
const refineMessage = formatRefineRequestForAgent(session.summary);
|
||||
await continueAgentConversation(session, refineMessage);
|
||||
} else {
|
||||
// Record the response
|
||||
session.history.push({
|
||||
question: session.currentQuestion,
|
||||
const currentQuestion = session.currentQuestion;
|
||||
const historyEntry = {
|
||||
question: currentQuestion,
|
||||
response: responses,
|
||||
thinkingOutput: session.lastGeneratedThinking || "",
|
||||
});
|
||||
};
|
||||
|
||||
session.error = undefined;
|
||||
persistSession(session, "generating");
|
||||
|
||||
if (!session.agent) {
|
||||
const replayHistory = session.history.slice(0, -1);
|
||||
await ensureSessionAgent(session, rootDir, replayHistory, promptOverrides, store);
|
||||
await ensureSessionAgent(session, rootDir, session.history, promptOverrides, store);
|
||||
}
|
||||
|
||||
const message = formatResponseForAgent(session.currentQuestion, responses);
|
||||
const message = formatResponseForAgent(currentQuestion, responses);
|
||||
await continueAgentConversation(session, message);
|
||||
|
||||
if (!session.error) {
|
||||
session.history.push(historyEntry);
|
||||
}
|
||||
}
|
||||
|
||||
// Return the current state (will be updated via SSE)
|
||||
@@ -2491,6 +2523,25 @@ export function __setPlanningNtfyHelpers(mock: PlanningNtfyHelpers | undefined):
|
||||
planningNtfyHelpers = mock;
|
||||
}
|
||||
|
||||
/** Test-only helper for validating generation tracking behavior. */
|
||||
export function __getActiveGenerationForTests(sessionId: string):
|
||||
| { abortController: AbortController; timer: NodeJS.Timeout }
|
||||
| undefined {
|
||||
return activeGenerations.get(sessionId);
|
||||
}
|
||||
|
||||
/** Test-only helper for exercising generation timeout orchestration directly. */
|
||||
export async function __runGenerationWithTimeoutForTests<T>(
|
||||
sessionId: string,
|
||||
operation: (abortSignal: AbortSignal) => Promise<T>,
|
||||
): Promise<T> {
|
||||
const session = getSession(sessionId);
|
||||
if (!session) {
|
||||
throw new SessionNotFoundError(`Planning session ${sessionId} not found or expired`);
|
||||
}
|
||||
return runGenerationWithTimeout(session, operation);
|
||||
}
|
||||
|
||||
// ── Custom Errors ───────────────────────────────────────────────────────────
|
||||
|
||||
export class RateLimitError extends Error {
|
||||
@@ -2513,3 +2564,10 @@ export class InvalidSessionStateError extends Error {
|
||||
this.name = "InvalidSessionStateError";
|
||||
}
|
||||
}
|
||||
|
||||
export class GenerationInProgressError extends Error {
|
||||
constructor(message: string) {
|
||||
super(message);
|
||||
this.name = "GenerationInProgressError";
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,7 +6,7 @@ import {
|
||||
type TaskPriority,
|
||||
type TaskStore,
|
||||
} from "@fusion/core";
|
||||
import { ApiError, badRequest, notFound, rateLimited } from "../api-error.js";
|
||||
import { ApiError, badRequest, conflict, notFound, rateLimited } from "../api-error.js";
|
||||
import { writeSSEEvent, type SessionBufferedEvent } from "../sse-buffer.js";
|
||||
import type { AiSessionStore } from "../ai-session-store.js";
|
||||
import type { ApiRoutesContext } from "./types.js";
|
||||
@@ -773,6 +773,8 @@ export function registerPlanningSubtaskRoutes(ctx: ApiRoutesContext, deps: Plann
|
||||
throw notFound(err instanceof Error ? err.message : String(err));
|
||||
} else if (err instanceof Error && err.name === "InvalidSessionStateError") {
|
||||
throw badRequest(err instanceof Error ? err.message : String(err));
|
||||
} else if (err instanceof Error && err.name === "GenerationInProgressError") {
|
||||
throw conflict(err instanceof Error ? err.message : String(err));
|
||||
} else {
|
||||
rethrowAsApiError(err, "Failed to process response");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user