diff --git a/.changeset/planning-turn-sync.md b/.changeset/planning-turn-sync.md new file mode 100644 index 0000000000..fc091f22bc --- /dev/null +++ b/.changeset/planning-turn-sync.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Keep Planning Mode questions and the running plan in sync after each answer. +category: fix +dev: Mid-interview SSE no longer terminalizes on running summary; client ignores stale answered questions and reconciles submit failures against server state. diff --git a/packages/dashboard/app/api/legacy.ts b/packages/dashboard/app/api/legacy.ts index 026b86fab6..95ce5e8bf2 100644 --- a/packages/dashboard/app/api/legacy.ts +++ b/packages/dashboard/app/api/legacy.ts @@ -2500,6 +2500,8 @@ export interface PlanningSession { summary: PlanningSummary | null; } +/** The response endpoint may synchronously return a generated next question before SSE delivers it. */ +export type PlanningResponse = PlanningSession | { type: "question"; data: PlanningQuestion }; /** SSE event types for planning session streaming */ export type PlanningStreamEvent = @@ -2632,8 +2634,8 @@ export function respondToPlanning( sessionId: string, responses: Record, projectId?: string, -): Promise { - return api(withProjectId("/planning/respond", projectId), { +): Promise { + return api(withProjectId("/planning/respond", projectId), { method: "POST", body: JSON.stringify({ sessionId, responses }), }); diff --git a/packages/dashboard/app/components/PlanningModeModal.tsx b/packages/dashboard/app/components/PlanningModeModal.tsx index 38de00e88e..b3cc7960ac 100644 --- a/packages/dashboard/app/components/PlanningModeModal.tsx +++ b/packages/dashboard/app/components/PlanningModeModal.tsx @@ -334,9 +334,11 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat const [error, setError] = useState(null); const [, setResponseHistory] = useState([]); const [conversationHistory, setConversationHistory] = useState([]); + const conversationHistoryRef = useRef([]); const [editedSummary, setEditedSummary] = useState(null); // FNXC:PlanningMode 2026-07-19-15:35: FN-8400 keeps the in-progress plan independent of the center-pane view so it remains visible while the next question is generating. const [runningSummary, setRunningSummary] = useState(null); + const runningSummaryRef = useRef(null); const [branchMode, setBranchMode] = useState<"project-default" | "auto-new" | "existing" | "custom-new">("project-default"); const [branchName, setBranchName] = useState(""); const [baseBranch, setBaseBranch] = useState(""); @@ -361,6 +363,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat Pending state belongs to the selected history entry and never replaces the running plan pane. */ const [editingQuestionId, setEditingQuestionId] = useState(null); + const editingQuestionIdRef = useRef(null); const [isHistoryEditPending, setIsHistoryEditPending] = useState(false); const [isRenamingSession, setIsRenamingSession] = useState(false); const [sessionTitleDraft, setSessionTitleDraft] = useState(""); @@ -641,6 +644,18 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat streamingOutputRef.current = streamingOutput; }, [streamingOutput]); + useEffect(() => { + conversationHistoryRef.current = conversationHistory; + }, [conversationHistory]); + + useEffect(() => { + runningSummaryRef.current = runningSummary; + }, [runningSummary]); + + useEffect(() => { + editingQuestionIdRef.current = editingQuestionId; + }, [editingQuestionId]); + // Keep the streaming AI thinking pane pinned to the bottom as new tokens // arrive. If the user has scrolled up to read earlier output, we leave the // scroll position alone — only auto-follow when they're already near the @@ -698,11 +713,29 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat if (cancelled || !session) return; if (currentSessionIdRef.current !== sessionId) return; if (session.status === "awaiting_input" && session.currentQuestion) { + /* + FNXC:PlanningTurnReconciliation 2026-07-20-10:36: + Missed SSE recovery must hydrate the server's entire interview turn together. Keeping + the persisted Q&A history and running result while replacing only the center question + prevents an already-answered question or a blank plan pane from representing a + different turn than the server. + */ resetPlanningAutoRetryBudget(); - const question = JSON.parse(session.currentQuestion) as PlanningQuestion; + const question = normalizeQuestionOptions(JSON.parse(session.currentQuestion) as PlanningQuestion); + const history = parseConversationHistory(session.conversationHistory); + const summary = session.result + ? normalizePlanningSummary(JSON.parse(session.result) as PlanningSummary) + : null; + conversationHistoryRef.current = history; + runningSummaryRef.current = summary; + setConversationHistory(history); + setResponseHistory(history + .map((entry) => entry.response) + .filter((response): response is QuestionResponse => Boolean(response && typeof response === "object" && !Array.isArray(response)))); + setRunningSummary(summary); setView({ type: "question", - session: { sessionId, currentQuestion: question, summary: null }, + session: { sessionId, currentQuestion: question, summary }, }); setStreamingOutput(""); } else if (session.status === "complete" && session.result) { @@ -887,6 +920,16 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat onQuestion: (question) => { if (isStaleEvent()) return; const normalizedQuestion = normalizeQuestionOptions(question); + const isAnsweredQuestion = conversationHistoryRef.current.some( + (entry) => entry.question?.id === normalizedQuestion.id && entry.response !== undefined, + ); + /* + FNXC:PlanningTurnReconciliation 2026-07-20-10:36: + Buffered SSE reconnects may replay a question that the user already submitted. An + answered question may only return through the explicit rewind/edit branch, never as a + passive stream catch-up event that overwrites a newer awaiting-input question. + */ + if (isAnsweredQuestion && editingQuestionIdRef.current !== normalizedQuestion.id) return; setIsReconnecting(false); setIsRetrying(false); resetPlanningAutoRetryBudget(); @@ -912,7 +955,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat setView({ type: "question", - session: { sessionId, currentQuestion: normalizedQuestion, summary: runningSummary }, + session: { sessionId, currentQuestion: normalizedQuestion, summary: runningSummaryRef.current }, }); setStreamingOutput(""); }, @@ -943,6 +986,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat allowed to enter terminal SummaryView, preventing a first-answer SSE race from ending the interview. */ + runningSummaryRef.current = normalizedSummary; setRunningSummary(normalizedSummary); setView((previous) => previous.type === "question" ? { @@ -1934,7 +1978,9 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat // Capture before clearing state: the edit branch rewrites this exact history row while // the server preserves the other answers and generates the appended next question. const submittedEditingQuestionId = editingQuestionId; + const historyBeforeSubmit = conversationHistoryRef.current; setEditingQuestionId(null); + editingQuestionIdRef.current = null; // Keep the existing SSE connection alive - do NOT close it! // The connection established in handleStartPlanning will continue @@ -1945,23 +1991,22 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat setResponseHistory((prev) => submittedEditingQuestionId ? prev.map((response, index) => conversationHistory.filter((entry) => entry.question && entry.response)[index]?.question?.id === submittedEditingQuestionId ? responses : response) : [...prev, responses]); - setConversationHistory((prev) => { - // Capture any reasoning that accumulated since the last question - // (e.g. thinking streamed while the user was reading the question). - const currentThinking = streamingOutputRef.current.trim(); - let updated = prev; - if (currentThinking) { - const lastEntry = updated[updated.length - 1]; - if (lastEntry?.thinkingOutput !== currentThinking) { - updated = [...updated, { thinkingOutput: currentThinking }]; - } + // Capture any reasoning that accumulated since the last question + // (e.g. thinking streamed while the user was reading the question). + const currentThinking = streamingOutputRef.current.trim(); + let optimisticHistory = historyBeforeSubmit; + if (currentThinking) { + const lastEntry = optimisticHistory[optimisticHistory.length - 1]; + if (lastEntry?.thinkingOutput !== currentThinking) { + optimisticHistory = [...optimisticHistory, { thinkingOutput: currentThinking }]; } - const answer = { question: activeQuestion, response: responses }; - if (submittedEditingQuestionId) { - return updated.map((entry) => entry.question?.id === submittedEditingQuestionId ? answer : entry); - } - return [...updated, answer]; - }); + } + const answer = { question: activeQuestion, response: responses }; + optimisticHistory = submittedEditingQuestionId + ? optimisticHistory.map((entry) => entry.question?.id === submittedEditingQuestionId ? answer : entry) + : [...optimisticHistory, answer]; + conversationHistoryRef.current = optimisticHistory; + setConversationHistory(optimisticHistory); resetPlanningAutoRetryBudget(); setView({ type: "loading" }); setStreamingOutput(""); // Clear old thinking output when entering loading state @@ -1969,11 +2014,63 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat try { // Submit response - AI will broadcast events via the already-connected stream - await respondToPlanning(sessionId, responses, projectId); - // Events (question/summary) will arrive via the existing SSE stream + const response = await respondToPlanning(sessionId, responses, projectId); + // The stream normally drives this transition. Adopt an already-returned next question + // when its event was lost so loading cannot later be replaced by a stale replay. + if ("type" in response && response.type === "question" && response.data.id !== activeQuestion.id) { + const nextQuestion = normalizeQuestionOptions(response.data); + setView({ + type: "question", + session: { sessionId, currentQuestion: nextQuestion, summary: runningSummaryRef.current }, + }); + } } catch (err) { - setError(getErrorMessage(err) || t("planning.failedSubmitResponse", "Failed to submit response")); - setView({ type: "question", session }); + const errorMessage = getErrorMessage(err) || t("planning.failedSubmitResponse", "Failed to submit response"); + /* + FNXC:PlanningTurnReconciliation 2026-07-20-10:36: + A rejected HTTP response is ambiguous: the server may have accepted the answer before + the connection failed. Rehydrate durable state before restoring the form. If it was not + accepted, roll back the optimistic answer so history and the active question still agree. + */ + try { + const persisted = await fetchAiSession(sessionId); + if (persisted?.status === "awaiting_input" && persisted.currentQuestion) { + const history = parseConversationHistory(persisted.conversationHistory); + const summary = persisted.result + ? normalizePlanningSummary(JSON.parse(persisted.result) as PlanningSummary) + : null; + const currentQuestion = normalizeQuestionOptions(JSON.parse(persisted.currentQuestion) as PlanningQuestion); + conversationHistoryRef.current = history; + runningSummaryRef.current = summary; + setConversationHistory(history); + setResponseHistory(history + .map((entry) => entry.response) + .filter((response): response is QuestionResponse => Boolean(response && typeof response === "object" && !Array.isArray(response)))); + setRunningSummary(summary); + setError(errorMessage); + setView({ type: "question", session: { sessionId, currentQuestion, summary } }); + return; + } + if (persisted?.status === "generating") { + const history = parseConversationHistory(persisted.conversationHistory); + const summary = persisted.result + ? normalizePlanningSummary(JSON.parse(persisted.result) as PlanningSummary) + : null; + conversationHistoryRef.current = history; + runningSummaryRef.current = summary; + setConversationHistory(history); + setRunningSummary(summary); + setError(errorMessage); + setView({ type: "loading" }); + return; + } + } catch { + // Fall back to the known pre-submit turn and remove its optimistic answer. + } + conversationHistoryRef.current = historyBeforeSubmit; + setConversationHistory(historyBeforeSubmit); + setError(errorMessage); + setView({ type: "question", session: { ...session, summary: runningSummaryRef.current } }); } }, [conversationHistory, editingQuestionId, projectId, resetPlanningAutoRetryBudget, view] @@ -2158,14 +2255,18 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat try { const rewound = await rewindPlanningSession(view.session.sessionId, projectId, questionId); setEditingQuestionId(questionId); - setConversationHistory(rewound.history.map((item) => ({ + editingQuestionIdRef.current = questionId; + const rewoundHistory = rewound.history.map((item) => ({ question: item.question, response: item.response && typeof item.response === "object" && !Array.isArray(item.response) ? item.response as Record : { [item.question.id]: item.response }, thinkingOutput: item.thinkingOutput, - }))); - const nextSummary = rewound.summary ? normalizePlanningSummary(rewound.summary) : runningSummary; + })); + conversationHistoryRef.current = rewoundHistory; + setConversationHistory(rewoundHistory); + const nextSummary = rewound.summary ? normalizePlanningSummary(rewound.summary) : runningSummaryRef.current; + runningSummaryRef.current = nextSummary; setRunningSummary(nextSummary); setView({ type: "question", session: { ...view.session, currentQuestion: rewound.currentQuestion, summary: nextSummary } }); } catch (err) { diff --git a/packages/dashboard/app/components/__tests__/PlanningModeModal.planning-flow.test.tsx b/packages/dashboard/app/components/__tests__/PlanningModeModal.planning-flow.test.tsx index 9d98a5ba9c..b6b94d9882 100644 --- a/packages/dashboard/app/components/__tests__/PlanningModeModal.planning-flow.test.tsx +++ b/packages/dashboard/app/components/__tests__/PlanningModeModal.planning-flow.test.tsx @@ -3576,6 +3576,123 @@ describe("PlanningModeModal", () => { }); }); + /* + FNXC:PlanningTurnReconciliation 2026-07-20-10:36: + These regressions reproduce the operator-visible desync: an answered Q1 must never displace + Q2, and recovery must hydrate question, answered history, and running plan as one server turn. + */ + describe("interview turn reconciliation", () => { + const secondQuestion: PlanningQuestion = { + id: "q-turn-reconciliation-next", + type: "text", + question: "Which constraint matters most next?", + }; + const secondSummary = { + ...mockSummary, + title: "Updated synchronized plan", + description: "The plan reflects Q1 before asking Q2.", + }; + + it("keeps Q2 and the latest plan when a stale answered Q1 stream event replays on tablet", async () => { + mockViewport("tablet"); + let streamHandlers: any; + mockConnectPlanningStream.mockImplementationOnce((_sessionId: string, _projectId: string | undefined, handlers: any) => { + streamHandlers = handlers; + queuePlanningStreamEvent(() => handlers.onQuestion?.(mockQuestion)); + return { close: vi.fn(), isConnected: vi.fn().mockReturnValue(true) }; + }); + mockRespondToPlanning.mockImplementationOnce(async () => { + queuePlanningStreamEvent(() => { + streamHandlers.onSummary?.(secondSummary); + streamHandlers.onQuestion?.(secondQuestion); + }); + return { sessionId: "session-123", currentQuestion: null, summary: null }; + }); + + render(); + fireEvent.change(screen.getByPlaceholderText(/e.g., Build a user authentication/), { target: { value: "Synchronize every turn" } }); + fireEvent.click(screen.getByText("Start Planning")); + fireEvent.click(await screen.findByText("Medium")); + fireEvent.click(screen.getByRole("button", { name: "Next question" })); + + expect(await screen.findByText(secondQuestion.question)).toBeDefined(); + expect(screen.getByRole("button", { name: "Question" })).toHaveAttribute("aria-pressed", "true"); + fireEvent.click(screen.getByRole("button", { name: "Running plan" })); + expect(within(screen.getByRole("complementary", { name: "Running plan" })).getByText(secondSummary.title)).toBeDefined(); + fireEvent.click(screen.getByRole("button", { name: "Answered questions" })); + expect(within(screen.getByRole("complementary", { name: "Answered questions" })).getByText(mockQuestion.question)).toBeDefined(); + fireEvent.click(screen.getByRole("button", { name: "Question" })); + + await act(async () => { + streamHandlers.onQuestion?.(mockQuestion); + }); + + expect(screen.getByText(secondQuestion.question)).toBeDefined(); + expect(screen.queryByText("Planning Complete!")).toBeNull(); + }); + + it("rolls back an optimistic answer when submit fails before server acceptance", async () => { + mockRespondToPlanning.mockRejectedValueOnce(new Error("submit timed out")); + mockFetchAiSession.mockResolvedValueOnce({ + id: "session-123", + type: "planning", + status: "awaiting_input", + title: "Still awaiting Q1", + inputPayload: JSON.stringify({ initialPlan: "Recover submit" }), + conversationHistory: "[]", + currentQuestion: JSON.stringify(mockQuestion), + result: JSON.stringify(mockSummary), + thinkingOutput: "", + projectId: null, + }); + + render(); + await screen.findByText(mockQuestion.question); + fireEvent.click(screen.getByText("Medium")); + fireEvent.click(screen.getByRole("button", { name: "Next question" })); + + expect(await screen.findByText("submit timed out")).toBeDefined(); + expect(screen.getByText(mockQuestion.question)).toBeDefined(); + expect(screen.queryByTestId("conversation-history")).toBeNull(); + expect(within(screen.getByRole("complementary", { name: "Running plan" })).getByText(mockSummary.title)).toBeDefined(); + }); + + it("hydrates Q2, Q1 history, and running plan after loading poll misses SSE", async () => { + const persistedHistory = [{ question: mockQuestion, response: { [mockQuestion.id]: "medium" } }]; + try { + mockRespondToPlanning.mockResolvedValueOnce({ sessionId: "session-123", currentQuestion: null, summary: null }); + mockFetchAiSession.mockResolvedValueOnce({ + id: "session-123", + type: "planning", + status: "awaiting_input", + title: "Recovered turn", + inputPayload: JSON.stringify({ initialPlan: "Poll recovery" }), + conversationHistory: JSON.stringify(persistedHistory), + currentQuestion: JSON.stringify(secondQuestion), + result: JSON.stringify(secondSummary), + thinkingOutput: "", + projectId: null, + }); + + render(); + await screen.findByText(mockQuestion.question); + vi.useFakeTimers(); + fireEvent.click(screen.getByText("Medium")); + fireEvent.click(screen.getByRole("button", { name: "Next question" })); + + await act(async () => { + await vi.advanceTimersByTimeAsync(8000); + }); + + expect(screen.getByText(secondQuestion.question)).toBeDefined(); + expect(within(screen.getByRole("complementary", { name: "Answered questions" })).getByText(mockQuestion.question)).toBeDefined(); + expect(within(screen.getByRole("complementary", { name: "Running plan" })).getByText(secondSummary.title)).toBeDefined(); + } finally { + vi.useRealTimers(); + } + }); + }); + describe("session rename", () => { function renderActiveSessionForRename() { mockFetchAiSession.mockResolvedValueOnce({ diff --git a/packages/dashboard/src/__tests__/routes-planning.test.ts b/packages/dashboard/src/__tests__/routes-planning.test.ts index 77595078be..5f0b1473ba 100644 --- a/packages/dashboard/src/__tests__/routes-planning.test.ts +++ b/packages/dashboard/src/__tests__/routes-planning.test.ts @@ -1459,6 +1459,74 @@ describe("Planning Mode Routes", () => { expect(streamRes.body).toContain("What is your preference?"); }); + /* + FNXC:PlanningStreamTurnIdentity 2026-07-20-10:36: + A persisted summary is the active interview's running plan, not a completion sentinel. + Reconnect must send it before the awaiting-input question and keep the subscription alive. + */ + it("keeps a mid-interview stream open when a running summary and question are persisted", async () => { + const startRes = await REQUEST( + buildApp(), + "POST", + "/api/planning/start", + JSON.stringify({ initialPlan: "Reconnect a running interview" }), + { "Content-Type": "application/json" }, + ); + const sessionId = startRes.body.sessionId as string; + const { getSession } = await import("../planning.js"); + const session = await getSession(sessionId); + expect(session).toBeDefined(); + + const runningSummary = { + title: "Running plan", + description: "Keep the interview turn aligned.", + suggestedSize: "M", + keyDeliverables: ["A synchronized interview"], + }; + const nextQuestion = { + id: "q-mid-interview", + type: "text", + question: "What should happen next?", + }; + // @ts-expect-error - test setup mutates the in-memory active session. + session!.summary = runningSummary; + // @ts-expect-error - test setup mutates the in-memory active session. + session!.currentQuestion = nextQuestion; + + const streamPromise = REQUEST(buildApp(), "GET", `/api/planning/${sessionId}/stream`); + setTimeout(() => planningStreamManager.broadcast(sessionId, { type: "complete" }), 10); + const streamRes = await streamPromise; + + expect(streamRes.body).toContain("event: summary"); + expect(streamRes.body).toContain(runningSummary.title); + expect(streamRes.body).toContain("event: question"); + expect(streamRes.body).toContain(nextQuestion.question); + expect(streamRes.body.indexOf("event: summary")).toBeLessThan(streamRes.body.indexOf("event: question")); + }); + + it("terminalizes a validated persisted summary session", async () => { + const startRes = await REQUEST( + buildApp(), + "POST", + "/api/planning/start", + JSON.stringify({ initialPlan: "Validated reconnect" }), + { "Content-Type": "application/json" }, + ); + const sessionId = startRes.body.sessionId as string; + const { getSession } = await import("../planning.js"); + const session = await getSession(sessionId); + expect(session).toBeDefined(); + // @ts-expect-error - test setup mutates the in-memory terminal session. + session!.summary = { title: "Validated plan", description: "Ready", suggestedSize: "S", keyDeliverables: [] }; + // @ts-expect-error - test setup mutates the in-memory terminal session. + session!.validated = true; + + const streamRes = await REQUEST(buildApp(), "GET", `/api/planning/${sessionId}/stream`); + + expect(streamRes.body).toContain("event: summary"); + expect(streamRes.body).toContain("event: complete"); + }); + it("emits catch-up question event for awaiting_input sessions", async () => { // This test verifies the fix for the mismatch where a session was advertised as // needing input but the resume path initially entered loading state. diff --git a/packages/dashboard/src/routes/register-planning-subtask-routes.ts b/packages/dashboard/src/routes/register-planning-subtask-routes.ts index a451969a83..3bb97992c4 100644 --- a/packages/dashboard/src/routes/register-planning-subtask-routes.ts +++ b/packages/dashboard/src/routes/register-planning-subtask-routes.ts @@ -1551,6 +1551,13 @@ export function registerPlanningSubtaskRoutes(ctx: ApiRoutesContext, deps: Plann } } + /* + FNXC:PlanningStreamTurnIdentity 2026-07-20-10:36: + A running summary is persisted after every interview turn, so it is catch-up state rather + than completion evidence. Reconnect must refresh that plan and continue into the current + awaiting-input question; only Validate writes `session.validated`, which authorizes a + terminal complete event and closes the stream. + */ if (session.summary) { const existing = planningStreamManager.getBufferedEvents(sessionId, 0); const lastSummaryEvent = [...existing].reverse().find((event) => event.event === "summary"); @@ -1567,16 +1574,18 @@ export function registerPlanningSubtaskRoutes(ctx: ApiRoutesContext, deps: Plann } } - const lastCompleteEvent = [...existing].reverse().find((event) => event.event === "complete"); - const completeEventId = lastCompleteEvent?.id - ?? planningStreamManager.broadcast(sessionId, { type: "complete" }); + if (session.validated) { + const lastCompleteEvent = [...existing].reverse().find((event) => event.event === "complete"); + const completeEventId = lastCompleteEvent?.id + ?? planningStreamManager.broadcast(sessionId, { type: "complete" }); - if (lastEventId === undefined || completeEventId > lastEventId) { - writeSSEEvent(res, "complete", JSON.stringify({}), completeEventId); + if (lastEventId === undefined || completeEventId > lastEventId) { + writeSSEEvent(res, "complete", JSON.stringify({}), completeEventId); + } + + res.end(); + return; } - - res.end(); - return; } // First-connect catch-up should replay buffered thinking chunks so the