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:
7
.changeset/planning-turn-sync.md
Normal file
7
.changeset/planning-turn-sync.md
Normal 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.
|
||||
@@ -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<string, unknown>,
|
||||
projectId?: string,
|
||||
): Promise<PlanningSession> {
|
||||
return api<PlanningSession>(withProjectId("/planning/respond", projectId), {
|
||||
): Promise<PlanningResponse> {
|
||||
return api<PlanningResponse>(withProjectId("/planning/respond", projectId), {
|
||||
method: "POST",
|
||||
body: JSON.stringify({ sessionId, responses }),
|
||||
});
|
||||
|
||||
@@ -334,9 +334,11 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
|
||||
const [error, setError] = useState<string | null>(null);
|
||||
const [, setResponseHistory] = useState<QuestionResponse[]>([]);
|
||||
const [conversationHistory, setConversationHistory] = useState<ConversationHistoryEntry[]>([]);
|
||||
const conversationHistoryRef = useRef<ConversationHistoryEntry[]>([]);
|
||||
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.
|
||||
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 [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<string | null>(null);
|
||||
const editingQuestionIdRef = useRef<string | null>(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<string, unknown>
|
||||
: { [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) {
|
||||
|
||||
@@ -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", () => {
|
||||
function renderActiveSessionForRename() {
|
||||
mockFetchAiSession.mockResolvedValueOnce({
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user