FN-8433: synchronize Planning Mode interview turns

Keep Planning Mode questions, answer history, and the running plan aligned across streaming and recovery paths.

- Reconcile server and client interview state after answers, retries, and missed SSE events.
- Keep persisted mid-interview summaries streaming until explicit validation.
- Add client and route regressions for stale question replay and reconnect behavior.

Files changed:
 .changeset/planning-turn-sync.md                   |   7 +
 packages/dashboard/app/api/legacy.ts               |   6 +-
 .../dashboard/app/components/PlanningModeModal.tsx | 153 +++++++++++++++++----
 .../PlanningModeModal.planning-flow.test.tsx       | 117 ++++++++++++++++
 .../src/__tests__/routes-planning.test.ts          |  68 +++++++++
 .../src/routes/register-planning-subtask-routes.ts |  25 ++--
 6 files changed, 340 insertions(+), 36 deletions(-)

Fusion-Task-Id: FN-8433

Fusion-Task-Lineage: d59a5b75-6ea0-4f97-83ab-ede770b4ae06

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-07-20 11:31:14 -07:00
parent 3962222863
commit 02e297aab5
6 changed files with 340 additions and 36 deletions

View File

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

View File

@@ -2500,6 +2500,8 @@ export interface PlanningSession {
summary: PlanningSummary | null; 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 */ /** SSE event types for planning session streaming */
export type PlanningStreamEvent = export type PlanningStreamEvent =
@@ -2632,8 +2634,8 @@ export function respondToPlanning(
sessionId: string, sessionId: string,
responses: Record<string, unknown>, responses: Record<string, unknown>,
projectId?: string, projectId?: string,
): Promise<PlanningSession> { ): Promise<PlanningResponse> {
return api<PlanningSession>(withProjectId("/planning/respond", projectId), { return api<PlanningResponse>(withProjectId("/planning/respond", projectId), {
method: "POST", method: "POST",
body: JSON.stringify({ sessionId, responses }), body: JSON.stringify({ sessionId, responses }),
}); });

View File

@@ -334,9 +334,11 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
const [error, setError] = useState<string | null>(null); const [error, setError] = useState<string | null>(null);
const [, setResponseHistory] = useState<QuestionResponse[]>([]); const [, setResponseHistory] = useState<QuestionResponse[]>([]);
const [conversationHistory, setConversationHistory] = useState<ConversationHistoryEntry[]>([]); const [conversationHistory, setConversationHistory] = useState<ConversationHistoryEntry[]>([]);
const conversationHistoryRef = useRef<ConversationHistoryEntry[]>([]);
const [editedSummary, setEditedSummary] = useState<PlanningSummary | null>(null); const [editedSummary, setEditedSummary] = useState<PlanningSummary | null>(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. // 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<PlanningSummary | null>(null); const [runningSummary, setRunningSummary] = useState<PlanningSummary | null>(null);
const runningSummaryRef = useRef<PlanningSummary | null>(null);
const [branchMode, setBranchMode] = useState<"project-default" | "auto-new" | "existing" | "custom-new">("project-default"); const [branchMode, setBranchMode] = useState<"project-default" | "auto-new" | "existing" | "custom-new">("project-default");
const [branchName, setBranchName] = useState(""); const [branchName, setBranchName] = useState("");
const [baseBranch, setBaseBranch] = 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. Pending state belongs to the selected history entry and never replaces the running plan pane.
*/ */
const [editingQuestionId, setEditingQuestionId] = useState<string | null>(null); const [editingQuestionId, setEditingQuestionId] = useState<string | null>(null);
const editingQuestionIdRef = useRef<string | null>(null);
const [isHistoryEditPending, setIsHistoryEditPending] = useState(false); const [isHistoryEditPending, setIsHistoryEditPending] = useState(false);
const [isRenamingSession, setIsRenamingSession] = useState(false); const [isRenamingSession, setIsRenamingSession] = useState(false);
const [sessionTitleDraft, setSessionTitleDraft] = useState(""); const [sessionTitleDraft, setSessionTitleDraft] = useState("");
@@ -641,6 +644,18 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
streamingOutputRef.current = streamingOutput; streamingOutputRef.current = streamingOutput;
}, [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 // 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 // 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 // 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 (cancelled || !session) return;
if (currentSessionIdRef.current !== sessionId) return; if (currentSessionIdRef.current !== sessionId) return;
if (session.status === "awaiting_input" && session.currentQuestion) { 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(); 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({ setView({
type: "question", type: "question",
session: { sessionId, currentQuestion: question, summary: null }, session: { sessionId, currentQuestion: question, summary },
}); });
setStreamingOutput(""); setStreamingOutput("");
} else if (session.status === "complete" && session.result) { } else if (session.status === "complete" && session.result) {
@@ -887,6 +920,16 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
onQuestion: (question) => { onQuestion: (question) => {
if (isStaleEvent()) return; if (isStaleEvent()) return;
const normalizedQuestion = normalizeQuestionOptions(question); 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); setIsReconnecting(false);
setIsRetrying(false); setIsRetrying(false);
resetPlanningAutoRetryBudget(); resetPlanningAutoRetryBudget();
@@ -912,7 +955,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
setView({ setView({
type: "question", type: "question",
session: { sessionId, currentQuestion: normalizedQuestion, summary: runningSummary }, session: { sessionId, currentQuestion: normalizedQuestion, summary: runningSummaryRef.current },
}); });
setStreamingOutput(""); 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 allowed to enter terminal SummaryView, preventing a first-answer SSE race from ending
the interview. the interview.
*/ */
runningSummaryRef.current = normalizedSummary;
setRunningSummary(normalizedSummary); setRunningSummary(normalizedSummary);
setView((previous) => previous.type === "question" 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 // 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. // the server preserves the other answers and generates the appended next question.
const submittedEditingQuestionId = editingQuestionId; const submittedEditingQuestionId = editingQuestionId;
const historyBeforeSubmit = conversationHistoryRef.current;
setEditingQuestionId(null); setEditingQuestionId(null);
editingQuestionIdRef.current = null;
// Keep the existing SSE connection alive - do NOT close it! // Keep the existing SSE connection alive - do NOT close it!
// The connection established in handleStartPlanning will continue // The connection established in handleStartPlanning will continue
@@ -1945,23 +1991,22 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
setResponseHistory((prev) => submittedEditingQuestionId setResponseHistory((prev) => submittedEditingQuestionId
? prev.map((response, index) => conversationHistory.filter((entry) => entry.question && entry.response)[index]?.question?.id === submittedEditingQuestionId ? responses : response) ? prev.map((response, index) => conversationHistory.filter((entry) => entry.question && entry.response)[index]?.question?.id === submittedEditingQuestionId ? responses : response)
: [...prev, responses]); : [...prev, responses]);
setConversationHistory((prev) => { // Capture any reasoning that accumulated since the last question
// Capture any reasoning that accumulated since the last question // (e.g. thinking streamed while the user was reading the question).
// (e.g. thinking streamed while the user was reading the question). const currentThinking = streamingOutputRef.current.trim();
const currentThinking = streamingOutputRef.current.trim(); let optimisticHistory = historyBeforeSubmit;
let updated = prev; if (currentThinking) {
if (currentThinking) { const lastEntry = optimisticHistory[optimisticHistory.length - 1];
const lastEntry = updated[updated.length - 1]; if (lastEntry?.thinkingOutput !== currentThinking) {
if (lastEntry?.thinkingOutput !== currentThinking) { optimisticHistory = [...optimisticHistory, { thinkingOutput: currentThinking }];
updated = [...updated, { thinkingOutput: currentThinking }];
}
} }
const answer = { question: activeQuestion, response: responses }; }
if (submittedEditingQuestionId) { const answer = { question: activeQuestion, response: responses };
return updated.map((entry) => entry.question?.id === submittedEditingQuestionId ? answer : entry); optimisticHistory = submittedEditingQuestionId
} ? optimisticHistory.map((entry) => entry.question?.id === submittedEditingQuestionId ? answer : entry)
return [...updated, answer]; : [...optimisticHistory, answer];
}); conversationHistoryRef.current = optimisticHistory;
setConversationHistory(optimisticHistory);
resetPlanningAutoRetryBudget(); resetPlanningAutoRetryBudget();
setView({ type: "loading" }); setView({ type: "loading" });
setStreamingOutput(""); // Clear old thinking output when entering loading state setStreamingOutput(""); // Clear old thinking output when entering loading state
@@ -1969,11 +2014,63 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
try { try {
// Submit response - AI will broadcast events via the already-connected stream // Submit response - AI will broadcast events via the already-connected stream
await respondToPlanning(sessionId, responses, projectId); const response = await respondToPlanning(sessionId, responses, projectId);
// Events (question/summary) will arrive via the existing SSE stream // 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) { } catch (err) {
setError(getErrorMessage(err) || t("planning.failedSubmitResponse", "Failed to submit response")); const errorMessage = getErrorMessage(err) || t("planning.failedSubmitResponse", "Failed to submit response");
setView({ type: "question", session }); /*
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] [conversationHistory, editingQuestionId, projectId, resetPlanningAutoRetryBudget, view]
@@ -2158,14 +2255,18 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
try { try {
const rewound = await rewindPlanningSession(view.session.sessionId, projectId, questionId); const rewound = await rewindPlanningSession(view.session.sessionId, projectId, questionId);
setEditingQuestionId(questionId); setEditingQuestionId(questionId);
setConversationHistory(rewound.history.map((item) => ({ editingQuestionIdRef.current = questionId;
const rewoundHistory = rewound.history.map((item) => ({
question: item.question, question: item.question,
response: item.response && typeof item.response === "object" && !Array.isArray(item.response) response: item.response && typeof item.response === "object" && !Array.isArray(item.response)
? item.response as Record<string, unknown> ? item.response as Record<string, unknown>
: { [item.question.id]: item.response }, : { [item.question.id]: item.response },
thinkingOutput: item.thinkingOutput, 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); setRunningSummary(nextSummary);
setView({ type: "question", session: { ...view.session, currentQuestion: rewound.currentQuestion, summary: nextSummary } }); setView({ type: "question", session: { ...view.session, currentQuestion: rewound.currentQuestion, summary: nextSummary } });
} catch (err) { } catch (err) {

View File

@@ -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(<PlanningModeModal isOpen={true} onClose={mockOnClose} onTaskCreated={mockOnTaskCreated} onTasksCreated={vi.fn()} tasks={mockTasks} />);
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(<PlanningModeModal isOpen={true} onClose={mockOnClose} onTaskCreated={mockOnTaskCreated} onTasksCreated={vi.fn()} tasks={mockTasks} initialPlan="Recover submit" />);
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(<PlanningModeModal isOpen={true} onClose={mockOnClose} onTaskCreated={mockOnTaskCreated} onTasksCreated={vi.fn()} tasks={mockTasks} initialPlan="Poll recovery" />);
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", () => { describe("session rename", () => {
function renderActiveSessionForRename() { function renderActiveSessionForRename() {
mockFetchAiSession.mockResolvedValueOnce({ mockFetchAiSession.mockResolvedValueOnce({

View File

@@ -1459,6 +1459,74 @@ describe("Planning Mode Routes", () => {
expect(streamRes.body).toContain("What is your preference?"); 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 () => { 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 // This test verifies the fix for the mismatch where a session was advertised as
// needing input but the resume path initially entered loading state. // needing input but the resume path initially entered loading state.

View File

@@ -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) { if (session.summary) {
const existing = planningStreamManager.getBufferedEvents(sessionId, 0); const existing = planningStreamManager.getBufferedEvents(sessionId, 0);
const lastSummaryEvent = [...existing].reverse().find((event) => event.event === "summary"); 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"); if (session.validated) {
const completeEventId = lastCompleteEvent?.id const lastCompleteEvent = [...existing].reverse().find((event) => event.event === "complete");
?? planningStreamManager.broadcast(sessionId, { type: "complete" }); const completeEventId = lastCompleteEvent?.id
?? planningStreamManager.broadcast(sessionId, { type: "complete" });
if (lastEventId === undefined || completeEventId > lastEventId) { if (lastEventId === undefined || completeEventId > lastEventId) {
writeSSEEvent(res, "complete", JSON.stringify({}), completeEventId); writeSSEEvent(res, "complete", JSON.stringify({}), completeEventId);
}
res.end();
return;
} }
res.end();
return;
} }
// First-connect catch-up should replay buffered thinking chunks so the // First-connect catch-up should replay buffered thinking chunks so the