From 7527d2651f8619e2aac8ac8958adf21d12ce05ff Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Sun, 16 Aug 2026 00:59:48 -0700 Subject: [PATCH] FN-9116: fence planning reconciliation by turn ownership Prevent stale Planning Mode snapshots and recovery callbacks from replacing a newer interview turn. - track session-load and turn epochs across response reconciliation, polling, streaming, and automatic retry - fence loading-poll fetches at launch so accepted SSE questions and newer responses retain ownership - add desktop and mobile race-ordering coverage and document the resolved suite-only flake - publish a patch changeset for the operator-visible recovery fix Files changed: .changeset/fn-9116-planning-reconciliation.md | 7 + .../suite-only-flakes-observed-register.md | 19 + .../dashboard/app/components/PlanningModeModal.tsx | 109 ++++- .../PlanningModeModal.planning-flow.test.tsx | 502 ++++++++++++++++++++- 4 files changed, 614 insertions(+), 23 deletions(-) Fusion-Task-Id: FN-9116 Fusion-Task-Lineage: b861176e-9104-4caf-8b6c-e645972aa5f6 Co-authored-by: Fusion (runfusion.ai) --- .changeset/fn-9116-planning-reconciliation.md | 7 + .../suite-only-flakes-observed-register.md | 19 + .../app/components/PlanningModeModal.tsx | 109 +++- .../PlanningModeModal.planning-flow.test.tsx | 502 +++++++++++++++++- 4 files changed, 614 insertions(+), 23 deletions(-) create mode 100644 .changeset/fn-9116-planning-reconciliation.md diff --git a/.changeset/fn-9116-planning-reconciliation.md b/.changeset/fn-9116-planning-reconciliation.md new file mode 100644 index 0000000000..af4012416f --- /dev/null +++ b/.changeset/fn-9116-planning-reconciliation.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Keep Planning Mode on the current session after a stale response refresh. +category: fix +dev: Fence duplicate-response, accepted stream-error, and loading-poll recovery by session, load, and turn ownership. diff --git a/docs/solutions/test-failures/suite-only-flakes-observed-register.md b/docs/solutions/test-failures/suite-only-flakes-observed-register.md index ececbdc314..b607c39563 100644 --- a/docs/solutions/test-failures/suite-only-flakes-observed-register.md +++ b/docs/solutions/test-failures/suite-only-flakes-observed-register.md @@ -151,3 +151,22 @@ The timeout occurred after all test assertions and is unrelated to FN-8979's can | earlier busy-machine standard lane | **1 failed** (`'desktop'` row); targeted rerun 58/58 passed | This file now carries THREE distinct register/ledger histories (entries 4 and 5 above plus this one) and one prior FN-8936 stabilization. Under the AGENTS.md repeated-quarantine rule this is a subsystem product-race smell: the duplicate-response generation reconciliation path (FN-8756 banner suppression / duplicate-generation dedup) should be investigated as a product race rather than stabilized a fourth time. Filed as a Fusion task; a second clean sighting of this exact test is an ordinary on-sight quarantine. + +**Resolved 2026-08-16 (FN-9116): Product race.** `handleSubmitResponse` caught a duplicate response-generation rejection, awaited `fetchAiSession(sessionId)`, then wrote its old session snapshot after a newer writer could already own the UI. The fix captures the response load and turn epochs before the response await, so an A → B → A reload cannot let the old A response adopt the new A load epoch. Every reconciliation/fallback write drops when a newer load, response, stream event, or recovery transition owns the view. + +Crucially, an accepted SSE `onError` is a turn boundary only after stale-event rejection. Its recovery captures that turn token across fetch and auto-retry awaits; a later response cannot be overwritten by an old reconnect, retry failure, or permanent error, and reconciliation from the errored turn cannot overwrite the recovery. The loading-poll error path now also claims its recovery turn *before* auto-retry: a successful retry returns early, so claiming afterward had left a held reconciliation authorized to overwrite recovery loading state. + +FN-9116 adds deterministic ordering coverage for desktop and mobile rows across durable-question, result-only plan-review, generating snapshots, A → B → A reload/rejection, `onError`-before-reconciliation, `onError` recovery losing ownership to a later response, and loading-poll recovery landing before a held reconciliation. The non-duplicate actionable-error assertions remain intact and passing. Response actions now settle hydration and query the live control before dispatch, removing the detached hydration-node test seam without changing product semantics. + +- **Resolved tree/SHA:** `d5f29bbdbc` (FN-9116 worktree; final documentation commit follows). + +| verification | result | +|---|---| +| targeted planning-flow file ×3 | **passed** (76 tests each run) | +| shared-helper sibling suites ×1 | **passed** | +| `app:backfill-3` run 1 | **passed** (5,693 tests) | +| `app:backfill-3` run 2 | **passed** (5,693 tests) | +| `app:backfill-3` run 3 | **passed** (5,693 tests) | +| `pnpm lint`, `pnpm verify:fast`, `pnpm build` | **passed** | + +The flake is structurally removed rather than stabilized: every hydration/recovery writer now has an ownership boundary before it can overwrite a newer turn. This is a published behavior fix, so FN-9116 includes a patch changeset. diff --git a/packages/dashboard/app/components/PlanningModeModal.tsx b/packages/dashboard/app/components/PlanningModeModal.tsx index 3284d5ae5c..3df637341e 100644 --- a/packages/dashboard/app/components/PlanningModeModal.tsx +++ b/packages/dashboard/app/components/PlanningModeModal.tsx @@ -614,10 +614,16 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat mutate the newly selected session, while successful progress still resets the three-attempt budget. */ const planningAutoRetryAttemptRef = useRef(0); - const planningAutoRetryOwnerRef = useRef<{ sessionId: string; token: symbol } | null>(null); - const startPlanningAutoRetryRef = useRef<(sessionId: string) => Promise>(async () => false); + const planningAutoRetryOwnerRef = useRef<{ sessionId: string; token: symbol; ownsTurn?: () => boolean } | null>(null); + const startPlanningAutoRetryRef = useRef<(sessionId: string, ownsTurn?: () => boolean) => Promise>(async () => false); const planningSessionLoadEpochRef = useRef(0); /* + FNXC:PlanningTurnReconciliation 2026-08-16-06:22: + Same-session streamed durable turns and response submissions can supersede a pending + reconciliation without changing session identity. This epoch lets their newer writer win. + */ + const planningTurnEpochRef = useRef(0); + /* FNXC:PlanningMode 2026-07-02-07:56: Refine Further is a single-flight completed-summary turn. Guard synchronously with a ref so duplicate click, touch, or keyboard activations cannot submit a second refine request or close the active stream with a generation-in-progress error before React renders the disabled state. */ @@ -1261,11 +1267,22 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat const tick = async () => { const sessionId = currentSessionIdRef.current; if (!sessionId) return; + /* + FNXC:PlanningTurnReconciliation 2026-08-16-07:57: + A loading poll may resolve after an accepted SSE question, response submission, or newer + session load has claimed the same session. Fence the fetch at launch so its stale snapshot + cannot advance the turn epoch and overwrite the newer question, summary, or recovery state. + */ + const pollTurnEpoch = planningTurnEpochRef.current; + const pollLoadEpoch = planningSessionLoadEpochRef.current; try { const session = await fetchAiSession(sessionId); if (cancelled || !session) return; if (currentSessionIdRef.current !== sessionId) return; + if (planningSessionLoadEpochRef.current !== pollLoadEpoch) return; + if (planningTurnEpochRef.current !== pollTurnEpoch) return; if (session.status === "awaiting_input" && !session.currentQuestion && session.result) { + planningTurnEpochRef.current += 1; // Recover a legacy or partially persisted plan when its question event was missed. // New sequential turns normally settle with both result and currentQuestion. resetPlanningAutoRetryBudget(); @@ -1281,6 +1298,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat setView({ type: "plan_review", session: { sessionId, currentQuestion: null, summary }, summary }); setStreamingOutput(""); } else if (session.status === "awaiting_input" && session.currentQuestion) { + planningTurnEpochRef.current += 1; /* FNXC:PlanningTurnReconciliation 2026-07-20-10:36: Missed SSE recovery must hydrate the server's entire interview turn together. Keeping @@ -1308,6 +1326,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat }); setStreamingOutput(""); } else if (session.status === "complete" && session.result) { + planningTurnEpochRef.current += 1; const resume = resolveCompletePlanningResume(session); if (resume.kind === "unrecoverable") { setView({ @@ -1322,9 +1341,23 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat setStreamingOutput(""); } else if (session.status === "error") { const errorMessage = session.error || t("planning.sessionFailed2", "Session failed"); - const handled = await startPlanningAutoRetryRef.current(sessionId); + /* + FNXC:PlanningTurnReconciliation 2026-08-16-07:19: + A loading-poll error starts a recovery turn before auto-retry awaits. A successful retry + returns early, so claiming afterward left a pending duplicate-response reconciliation + authorized to overwrite recovery's loading state. The poll's turn owns both retry and + terminal-error writes; a newer question, response, stream recovery, or session switch + invalidates this predicate before either path can mutate the view. + */ + const recoveryTurnEpoch = ++planningTurnEpochRef.current; + const recoveryLoadEpoch = planningSessionLoadEpochRef.current; + const pollRecoveryStillOwnsTurn = () => !cancelled + && currentSessionIdRef.current === sessionId + && planningSessionLoadEpochRef.current === recoveryLoadEpoch + && planningTurnEpochRef.current === recoveryTurnEpoch; + const handled = await startPlanningAutoRetryRef.current(sessionId, pollRecoveryStillOwnsTurn); if (handled) return; - if (cancelled || currentSessionIdRef.current !== sessionId) return; + if (!pollRecoveryStillOwnsTurn()) return; /* FNXC:PlanningRetry 2026-07-13-00:05: Mirror the SSE onError terminal-error transition here: when this poll is the one that @@ -1509,6 +1542,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat passive stream catch-up event that overwrites a newer awaiting-input question. */ if (isAnsweredQuestion && editingQuestionIdRef.current !== normalizedQuestion.id) return; + planningTurnEpochRef.current += 1; setIsRetrying(false); resetPlanningAutoRetryBudget(); setIsRefiningSummary(false); @@ -1539,6 +1573,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat }, onSummary: (summary) => { if (isStaleEvent()) return; + planningTurnEpochRef.current += 1; const normalizedSummary = normalizePlanningSummary(summary); setIsRetrying(false); resetPlanningAutoRetryBudget(); @@ -1578,6 +1613,20 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat }, onError: (message) => { if (isStaleEvent()) return; + /* + FNXC:PlanningTurnReconciliation 2026-08-16-06:45: + An accepted stream error starts recovery for the current turn, whether it reconnects, + auto-retries, or renders a permanent error. Its turn token remains authoritative across + recovery awaits, so a later response cannot be replaced by an older recovery either. + Advance only after stale-event rejection so a delayed duplicate-response reconciliation + from the errored turn cannot overwrite recovery, while an obsolete stream cannot + invalidate the live turn. + */ + const recoveryTurnEpoch = ++planningTurnEpochRef.current; + const recoveryLoadEpoch = planningSessionLoadEpochRef.current; + const recoveryStillOwnsTurn = () => currentSessionIdRef.current === sessionId + && planningSessionLoadEpochRef.current === recoveryLoadEpoch + && planningTurnEpochRef.current === recoveryTurnEpoch; const errorMessage = message || t("planning.sessionFailed", "Session failed while contacting the AI."); // A single transient stream error (e.g. tab was backgrounded long @@ -1588,7 +1637,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat (async () => { try { const session = await fetchAiSession(sessionId); - if (isStaleEvent()) return; + if (!recoveryStillOwnsTurn()) return; if ( session && (session.status === "generating" || session.status === "awaiting_input") @@ -1599,7 +1648,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat } catch { // fall through to error view below } - if (isStaleEvent()) return; + if (!recoveryStillOwnsTurn()) return; /* FNXC:PlanningRetry 2026-07-21-10:00: @@ -1607,9 +1656,10 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat Returning to Planning must use the same bounded, single-flight retry path as a live turn so tab suspension or navigation never turns a resumable session into an error UI. */ - if (await startPlanningAutoRetryRef.current(sessionId)) { + if (await startPlanningAutoRetryRef.current(sessionId, recoveryStillOwnsTurn)) { return; } + if (!recoveryStillOwnsTurn()) return; setIsRetrying(false); setIsAutoRetrying(false); setIsRefiningSummary(false); @@ -1647,7 +1697,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat const startPlanningRetry = useCallback( async ( retryTarget: { sessionId: string; currentQuestion: PlanningQuestion | null; summary: PlanningSummary | null }, - options: { auto: boolean; retryToken?: symbol }, + options: { auto: boolean; retryToken?: symbol; ownsTurn?: () => boolean }, ) => { setError(null); setIsRetrying(!options.auto); @@ -1664,7 +1714,8 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat await retryPlanningSession(retryTarget.sessionId, projectId); } catch (err) { const retryStillOwnsSession = () => currentSessionIdRef.current === retryTarget.sessionId - && (!options.auto || planningAutoRetryOwnerRef.current?.token === options.retryToken); + && (!options.auto || planningAutoRetryOwnerRef.current?.token === options.retryToken) + && (options.ownsTurn?.() ?? true); if (!retryStillOwnsSession()) return; let retryError: unknown = err; @@ -1779,14 +1830,15 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat A rejected automatic attempt has settled before its bounded successor is queued, so release only its matching token first. The successor can then acquire ownership, while duplicate SSE/poll reports still coalesce only during a genuinely pending invocation; - a stale callback cannot clear a newer session's owner. + a stale callback cannot clear a newer session's owner. Preserve a recovery's turn + predicate into the queued successor because another turn can land before its microtask. */ if (options.retryToken && planningAutoRetryOwnerRef.current?.token === options.retryToken) { planningAutoRetryOwnerRef.current = null; } queueMicrotask(() => { if (currentSessionIdRef.current === retryTarget.sessionId) { - void startPlanningAutoRetryRef.current(retryTarget.sessionId); + void startPlanningAutoRetryRef.current(retryTarget.sessionId, options.ownsTurn); } }); return; @@ -1812,9 +1864,15 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat ); const startPlanningAutoRetry = useCallback( - async (sessionId: string) => { - if (planningAutoRetryOwnerRef.current?.sessionId === sessionId) { - return true; + async (sessionId: string, ownsTurn?: () => boolean) => { + if (!(ownsTurn?.() ?? true)) return false; + const existingOwner = planningAutoRetryOwnerRef.current; + if (existingOwner?.sessionId === sessionId) { + if (!existingOwner.ownsTurn || existingOwner.ownsTurn()) return true; + // FNXC:PlanningTurnReconciliation 2026-08-16-07:19: An invalidated recovery + // attempt cannot satisfy a newer error's recovery. Release its token so the current + // turn starts the bounded retry that its caller is waiting to observe. + planningAutoRetryOwnerRef.current = null; } if (viewRef.current.type === "error") return false; // FNXC:PlanningRetry 2026-07-22-21:00: budget is per-session and survives remounts. @@ -1828,12 +1886,12 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat const retryToken = Symbol(`planning-auto-retry:${sessionId}:${attempt}`); planningAutoRetryAttemptsBySession.set(sessionId, attempt); planningAutoRetryAttemptRef.current = attempt; - planningAutoRetryOwnerRef.current = { sessionId, token: retryToken }; + planningAutoRetryOwnerRef.current = { sessionId, token: retryToken, ownsTurn }; setAutoRetryAttempt(attempt); setIsAutoRetrying(true); await startPlanningRetry( { sessionId, currentQuestion: null, summary: null }, - { auto: true, retryToken }, + { auto: true, retryToken, ownsTurn } ); return true; }, @@ -2809,6 +2867,8 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat } setError(null); + const responseTurnEpoch = ++planningTurnEpochRef.current; + const responseLoadEpoch = planningSessionLoadEpochRef.current; // 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; @@ -2872,14 +2932,21 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat } catch (err) { const errorMessage = getErrorMessage(err) || t("planning.failedSubmitResponse", "Failed to submit response"); const isDuplicateResponseConflict = isDuplicateResponseGenerationConflict(err); + const reconciliationStillOwnsTurn = () => currentSessionIdRef.current === sessionId + && planningSessionLoadEpochRef.current === responseLoadEpoch + && planningTurnEpochRef.current === responseTurnEpoch; /* - 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. + FNXC:PlanningTurnReconciliation 2026-08-16-06:22: + A rejected response may have been durably accepted, so reconciliation rehydrates before + rolling back optimism. Its snapshot loses to a newer session load, reset/stop, response + submission, or same-session streamed question/summary. Capture both ownership epochs + before the response await: an A → B → A reload must not let the old A response adopt the + new load epoch. A delayed refresh must then write nothing, or an old durable snapshot can + replace the question, history, summary, and error state that the newer turn already owns. */ try { const persisted = await fetchAiSession(sessionId); + if (!reconciliationStillOwnsTurn()) return; if (persisted?.status === "awaiting_input" && !persisted.currentQuestion && persisted.result) { const history = parseConversationHistory(persisted.conversationHistory); const summary = normalizePlanningSummary(JSON.parse(persisted.result) as PlanningSummary); @@ -2922,8 +2989,10 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat return; } } catch { + if (!reconciliationStillOwnsTurn()) return; // Fall back to the known pre-submit turn and remove its optimistic answer. } + if (!reconciliationStillOwnsTurn()) return; conversationHistoryRef.current = historyBeforeSubmit; setConversationHistory(historyBeforeSubmit); setError(errorMessage); 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 c6a0cb097b..5312836471 100644 --- a/packages/dashboard/app/components/__tests__/PlanningModeModal.planning-flow.test.tsx +++ b/packages/dashboard/app/components/__tests__/PlanningModeModal.planning-flow.test.tsx @@ -57,6 +57,19 @@ async function clickProceedAfterHydration() { fireEvent.click(screen.getByRole("button", { name: "Proceed with plan" })); } +/* +FNXC:PlanningTurnReconciliation 2026-08-16-07:19: +Resume hydration can replace the selected answer control before the user advances the interview. +Settle and re-query the live control so ordering tests exercise the response/reconciliation race, +not a detached DOM node that a user could never activate. +*/ +async function selectResponseAfterHydration(label: string) { + await screen.findByLabelText(label); + await act(async () => {}); + fireEvent.click(screen.getByLabelText(label)); + fireEvent.click(screen.getByRole("button", { name: "Next" })); +} + describe("PlanningModeModal sequential flow", () => { beforeEach(() => { vi.useRealTimers(); @@ -260,6 +273,129 @@ describe("PlanningModeModal sequential flow", () => { expect(screen.queryByText("session-a stream failed")).toBeNull(); }); + it.each(["desktop", "mobile"] as const)("does not let delayed duplicate reconciliation overwrite loading-poll recovery on %s", async (viewport) => { + mockViewportMode.mockReturnValue(viewport); + const sessionId = `poll-recovery-${viewport}`; + const submittedQuestion = { + id: "q-submitted", + type: "single_select", + question: "Which poll recovery owns this turn?", + options: [{ id: "secure", label: "Secure defaults" }], + }; + const intervalSpy = vi.spyOn(globalThis, "setInterval"); + let resolveReconciliation!: (session: Record) => void; + let fetchCount = 0; + mockFetchAiSession.mockImplementation(() => { + fetchCount += 1; + if (fetchCount === 1) { + return Promise.resolve({ + ...base, + id: sessionId, + status: "awaiting_input", + currentQuestion: JSON.stringify(submittedQuestion), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + } + if (fetchCount === 2) { + return new Promise((resolve) => { resolveReconciliation = resolve; }); + } + return Promise.resolve({ + ...base, + id: sessionId, + status: "error", + error: "Poll recovery owns this turn", + currentQuestion: null, + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + }); + mockRespondToPlanning.mockRejectedValue(new Error("Generation already in progress for this response")); + + renderSession(sessionId); + await selectResponseAfterHydration("Secure defaults"); + await waitFor(() => expect(fetchCount).toBe(2)); + + const poll = intervalSpy.mock.calls.find(([, delay]) => delay === 8000)?.[0]; + expect(poll).toBeTypeOf("function"); + await act(async () => { + await (poll as () => Promise)(); + }); + expect(mockRetryPlanningSession).toHaveBeenCalledWith(sessionId, "project-1"); + + await act(async () => { + resolveReconciliation({ + ...base, + id: sessionId, + status: "awaiting_input", + currentQuestion: JSON.stringify({ id: "q-stale", type: "text", question: "What did stale reconciliation ask?" }), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + await Promise.resolve(); + }); + + expect(screen.queryByText("What did stale reconciliation ask?")).toBeNull(); + expect(screen.queryByText("Generation already in progress for this response")).toBeNull(); + }); + + it.each(["desktop", "mobile"] as const)("keeps a newer streamed question when an older loading poll resolves on %s", async (viewport) => { + mockViewportMode.mockReturnValue(viewport); + const sessionId = `stale-loading-poll-${viewport}`; + const intervalSpy = vi.spyOn(globalThis, "setInterval"); + const streamedQuestion = { + id: "q-streamed", + type: "text", + question: "What did the current streamed turn ask?", + }; + let resolvePoll!: (session: Record) => void; + let fetchCount = 0; + mockFetchAiSession.mockImplementation(() => { + fetchCount += 1; + if (fetchCount === 1) { + return Promise.resolve({ + ...base, + id: sessionId, + status: "generating", + currentQuestion: null, + result: JSON.stringify(summaryWithRefinements), + inputPayload: JSON.stringify({ generationPurpose: "plan_update" }), + }); + } + return new Promise((resolve) => { resolvePoll = resolve; }); + }); + + renderSession(sessionId); + await waitFor(() => expect(mockConnectPlanningStream).toHaveBeenCalledTimes(1)); + await waitFor(() => expect(intervalSpy.mock.calls.some(([, delay]) => delay === 8000)).toBe(true)); + const poll = intervalSpy.mock.calls.find(([, delay]) => delay === 8000)?.[0]; + expect(poll).toBeTypeOf("function"); + let pendingPoll!: Promise; + act(() => { + pendingPoll = (poll as () => Promise)(); + }); + await waitFor(() => expect(fetchCount).toBe(2)); + + const handlers = mockConnectPlanningStream.mock.calls[0]?.[2]; + act(() => handlers?.onQuestion?.(streamedQuestion)); + expect(await screen.findByText("What did the current streamed turn ask?")).toBeInTheDocument(); + + await act(async () => { + resolvePoll({ + ...base, + id: sessionId, + status: "awaiting_input", + currentQuestion: JSON.stringify({ id: "q-stale-poll", type: "text", question: "What did the stale loading poll ask?" }), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + await pendingPoll; + }); + + expect(screen.getByText("What did the current streamed turn ask?")).toBeInTheDocument(); + expect(screen.queryByText("What did the stale loading poll ask?")).toBeNull(); + }); + it("automatically retries a resumed error discovered by the loading poll", async () => { const sessionId = "polled-error-session"; const intervalSpy = vi.spyOn(globalThis, "setInterval"); @@ -1043,7 +1179,11 @@ describe("PlanningModeModal sequential flow", () => { renderSession(); fireEvent.click(await screen.findByLabelText("Secure defaults")); - fireEvent.click(screen.getByRole("button", { name: "Next" })); + await act(async () => { + fireEvent.click(screen.getByRole("button", { name: "Next" })); + await Promise.resolve(); + }); + await waitFor(() => expect(mockRespondToPlanning).toHaveBeenCalledTimes(1)); if (status === "awaiting_input") { expect(await screen.findByText("What should the durable session ask next?")).toBeInTheDocument(); @@ -1054,6 +1194,363 @@ describe("PlanningModeModal sequential flow", () => { expect(document.querySelector(".planning-error")).toBeNull(); }); + it.each([ + { viewport: "desktop", persistedStatus: "awaiting_input", label: "a durable question" }, + { viewport: "mobile", persistedStatus: "awaiting_input", label: "a durable question" }, + { viewport: "desktop", persistedStatus: "plan_review", label: "a durable plan review" }, + { viewport: "mobile", persistedStatus: "plan_review", label: "a durable plan review" }, + { viewport: "desktop", persistedStatus: "generating", label: "generation progress" }, + { viewport: "mobile", persistedStatus: "generating", label: "generation progress" }, + ] as const)("keeps the newer session when delayed duplicate reconciliation returns $label on $viewport", async ({ viewport, persistedStatus }) => { + mockViewportMode.mockReturnValue(viewport); + const submittedQuestion = { + id: "q-submitted", + type: "single_select", + question: "Which outcome matters most?", + options: [{ id: "secure", label: "Secure defaults" }], + }; + const staleQuestion = { + id: "q-stale", + type: "text", + question: "What did the stale session ask?", + }; + const currentQuestion = { + id: "q-current", + type: "text", + question: "What should the current session ask?", + }; + let resolveStaleSession!: (session: Record) => void; + let sessionAReads = 0; + mockFetchAiSession.mockImplementation((sessionId: string) => { + if (sessionId === "session-a") { + sessionAReads += 1; + if (sessionAReads === 1) { + return Promise.resolve({ + ...base, + id: sessionId, + status: "awaiting_input", + currentQuestion: JSON.stringify(submittedQuestion), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + } + return new Promise((resolve) => { + resolveStaleSession = resolve; + }); + } + return Promise.resolve({ + ...base, + id: sessionId, + status: "awaiting_input", + currentQuestion: JSON.stringify(currentQuestion), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + }); + mockRespondToPlanning.mockRejectedValue(new Error("Generation already in progress for this response")); + + const props = { isOpen: true, onClose: vi.fn(), onTaskCreated: vi.fn(), onTasksCreated: vi.fn(), tasks: mockTasks, projectId: "project-1" }; + const { rerender } = render(); + fireEvent.click(await screen.findByLabelText("Secure defaults")); + await act(async () => { + fireEvent.click(screen.getByRole("button", { name: "Next" })); + await Promise.resolve(); + }); + await waitFor(() => expect(sessionAReads).toBe(2)); + + rerender(); + expect(await screen.findByText("What should the current session ask?")).toBeInTheDocument(); + + const staleSession = { + ...base, + id: "session-a", + status: persistedStatus === "generating" ? "generating" : "awaiting_input", + currentQuestion: persistedStatus === "awaiting_input" ? JSON.stringify(staleQuestion) : null, + result: JSON.stringify(summaryWithRefinements), + inputPayload: JSON.stringify({ generationPurpose: "plan_update" }), + }; + if (persistedStatus === "plan_review") { + staleSession.currentQuestion = null; + } + await act(async () => { + resolveStaleSession(staleSession); + await Promise.resolve(); + }); + + expect(screen.getByText("What should the current session ask?")).toBeInTheDocument(); + expect(screen.queryByText("What did the stale session ask?")).toBeNull(); + if (persistedStatus === "generating") { + expect(screen.queryByText("Generating plan…")).toBeNull(); + } + expect(screen.queryByText("Generation already in progress for this response")).toBeNull(); + }); + + it.each(["desktop", "mobile"] as const)("keeps a newer streamed question when delayed duplicate reconciliation resolves on %s", async (viewport) => { + mockViewportMode.mockReturnValue(viewport); + const submittedQuestion = { + id: "q-submitted", + type: "single_select", + question: "Which outcome matters most?", + options: [{ id: "secure", label: "Secure defaults" }], + }; + const streamedQuestion = { + id: "q-streamed", + type: "text", + question: "What did the newer streamed turn ask?", + }; + let resolveStaleSession!: (session: Record) => void; + let sessionReads = 0; + mockFetchAiSession.mockImplementation(() => { + sessionReads += 1; + if (sessionReads === 1) { + return Promise.resolve({ + ...base, + status: "awaiting_input", + currentQuestion: JSON.stringify(submittedQuestion), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + } + return new Promise((resolve) => { + resolveStaleSession = resolve; + }); + }); + mockRespondToPlanning.mockRejectedValue(new Error("Generation already in progress for this response")); + + renderSession(); + fireEvent.click(await screen.findByLabelText("Secure defaults")); + fireEvent.click(screen.getByRole("button", { name: "Next" })); + await waitFor(() => expect(sessionReads).toBe(2)); + await waitFor(() => expect(mockConnectPlanningStream).toHaveBeenCalledWith("session-1", "project-1", expect.any(Object))); + + const handlers = mockConnectPlanningStream.mock.calls.at(-1)?.[2]; + act(() => handlers?.onQuestion?.(streamedQuestion)); + expect(await screen.findByText("What did the newer streamed turn ask?")).toBeInTheDocument(); + + await act(async () => { + resolveStaleSession({ + ...base, + status: "generating", + currentQuestion: null, + result: JSON.stringify(summaryWithRefinements), + inputPayload: JSON.stringify({ generationPurpose: "plan_update" }), + }); + await Promise.resolve(); + }); + + expect(screen.getByText("What did the newer streamed turn ask?")).toBeInTheDocument(); + expect(screen.queryByText("Generating plan…")).toBeNull(); + expect(screen.queryByText("Generation already in progress for this response")).toBeNull(); + }); + + it.each([ + { viewport: "desktop", recovery: "reconnect" }, + { viewport: "mobile", recovery: "reconnect" }, + { viewport: "desktop", recovery: "permanent-error" }, + { viewport: "mobile", recovery: "permanent-error" }, + ] as const)("keeps $recovery stream-error recovery when stale duplicate reconciliation resolves on $viewport", async ({ viewport, recovery }) => { + mockViewportMode.mockReturnValue(viewport); + const sessionId = `duplicate-stream-${viewport}-${recovery}`; + const submittedQuestion = { + id: "q-submitted", + type: "single_select", + question: "Which recovery should win?", + options: [{ id: "secure", label: "Secure defaults" }], + }; + let resolveReconciliation!: (session: Record) => void; + let fetchCount = 0; + mockFetchAiSession.mockImplementation(() => { + fetchCount += 1; + if (fetchCount === 1) { + return Promise.resolve({ + ...base, + id: sessionId, + status: "awaiting_input", + currentQuestion: JSON.stringify(submittedQuestion), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + } + if (fetchCount === 2) { + return new Promise((resolve) => { + resolveReconciliation = resolve; + }); + } + return Promise.resolve(recovery === "reconnect" + ? { + ...base, + id: sessionId, + status: "generating", + currentQuestion: null, + result: JSON.stringify(summaryWithRefinements), + inputPayload: JSON.stringify({ generationPurpose: "plan_update" }), + } + : null); + }); + mockRespondToPlanning.mockRejectedValue(new Error("Generation already in progress for this response")); + if (recovery === "permanent-error") { + mockRetryPlanningSession.mockRejectedValue(new Error("Recovery could not continue")); + } + + renderSession(sessionId); + fireEvent.click(await screen.findByLabelText("Secure defaults")); + fireEvent.click(screen.getByRole("button", { name: "Next" })); + await waitFor(() => expect(fetchCount).toBe(2)); + + await act(async () => { + mockConnectPlanningStream.mock.calls[0]?.[2]?.onError?.("The planning stream failed"); + await Promise.resolve(); + }); + if (recovery === "reconnect") { + await waitFor(() => expect(mockConnectPlanningStream).toHaveBeenCalledTimes(2)); + } else { + expect(await screen.findByText("Recovery could not continue")).toBeInTheDocument(); + } + + await act(async () => { + resolveReconciliation({ + ...base, + id: sessionId, + status: "awaiting_input", + currentQuestion: JSON.stringify({ id: "q-stale", type: "text", question: "What did stale reconciliation ask?" }), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + await Promise.resolve(); + }); + + expect(screen.queryByText("What did stale reconciliation ask?")).toBeNull(); + if (recovery === "reconnect") { + expect(mockConnectPlanningStream).toHaveBeenCalledTimes(2); + expect(screen.queryByText("The planning stream failed")).toBeNull(); + } else { + expect(screen.getByText("Recovery could not continue")).toBeInTheDocument(); + } + }); + + it.each(["desktop", "mobile"] as const)("does not let a rejected response adopt a reloaded session epoch on %s", async (viewport) => { + mockViewportMode.mockReturnValue(viewport); + const submittedQuestion = { + id: "q-submitted", + type: "single_select", + question: "Which session owns this answer?", + options: [{ id: "secure", label: "Secure defaults" }], + }; + let rejectResponse!: (error: Error) => void; + let resolveStaleReconciliation!: (session: Record) => void; + let sessionACalls = 0; + mockFetchAiSession.mockImplementation((id: string) => { + if (id === "session-b") { + return Promise.resolve({ + ...base, + id, + status: "awaiting_input", + currentQuestion: JSON.stringify({ id: "q-b", type: "text", question: "What did session B ask?" }), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + } + sessionACalls += 1; + if (sessionACalls === 3) { + return new Promise((resolve) => { resolveStaleReconciliation = resolve; }); + } + return Promise.resolve({ + ...base, + id: "session-a", + status: "awaiting_input", + currentQuestion: JSON.stringify(sessionACalls === 1 + ? submittedQuestion + : { id: "q-reloaded", type: "text", question: "What did the reloaded session A ask?" }), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + }); + mockRespondToPlanning.mockReturnValue(new Promise((_resolve, reject) => { rejectResponse = reject; })); + const props = { isOpen: true, onClose: vi.fn(), onTaskCreated: vi.fn(), onTasksCreated: vi.fn(), tasks: mockTasks, projectId: "project-1" }; + const { rerender } = render(); + + fireEvent.click(await screen.findByLabelText("Secure defaults")); + await act(async () => { + fireEvent.click(screen.getByRole("button", { name: "Next" })); + await Promise.resolve(); + }); + await waitFor(() => expect(mockRespondToPlanning).toHaveBeenCalledTimes(1)); + rerender(); + expect(await screen.findByText("What did session B ask?")).toBeInTheDocument(); + rerender(); + expect(await screen.findByText("What did the reloaded session A ask?")).toBeInTheDocument(); + + await act(async () => { + rejectResponse(new Error("Generation already in progress for this response")); + await Promise.resolve(); + }); + await waitFor(() => expect(sessionACalls).toBe(3)); + await act(async () => { + resolveStaleReconciliation({ + ...base, + id: "session-a", + status: "awaiting_input", + currentQuestion: JSON.stringify({ id: "q-stale", type: "text", question: "What did stale reconciliation ask?" }), + result: JSON.stringify(summaryWithRefinements), + inputPayload: "{}", + }); + await Promise.resolve(); + }); + + expect(screen.getByText("What did the reloaded session A ask?")).toBeInTheDocument(); + expect(screen.queryByText("What did stale reconciliation ask?")).toBeNull(); + }); + + it.each(["desktop", "mobile"] as const)("does not let an older stream-error recovery reconnect over a later response on %s", async (viewport) => { + mockViewportMode.mockReturnValue(viewport); + const submittedQuestion = { + id: "q-submitted", + type: "single_select", + question: "Which response wins after stream recovery?", + options: [{ id: "secure", label: "Secure defaults" }], + }; + let resolveRecovery!: (session: Record) => void; + let fetchCount = 0; + mockFetchAiSession.mockImplementation(() => { + fetchCount += 1; + if (fetchCount === 1) { + return Promise.resolve({ + ...base, + status: "generating", + currentQuestion: null, + result: JSON.stringify(summaryWithRefinements), + inputPayload: JSON.stringify({ generationPurpose: "plan_update" }), + }); + } + return new Promise((resolve) => { resolveRecovery = resolve; }); + }); + mockRespondToPlanning.mockReturnValue(new Promise(() => undefined)); + + renderSession(); + await waitFor(() => expect(mockConnectPlanningStream).toHaveBeenCalledTimes(1)); + const handlers = mockConnectPlanningStream.mock.calls[0]?.[2]; + act(() => handlers?.onQuestion?.(submittedQuestion)); + expect(await screen.findByLabelText("Secure defaults")).toBeInTheDocument(); + handlers?.onError?.("The old stream failed"); + await waitFor(() => expect(fetchCount).toBe(2)); + fireEvent.click(screen.getByLabelText("Secure defaults")); + fireEvent.click(screen.getByRole("button", { name: "Next" })); + const connectionCountAfterResponse = mockConnectPlanningStream.mock.calls.length; + await act(async () => { + resolveRecovery({ + ...base, + status: "generating", + currentQuestion: null, + result: JSON.stringify(summaryWithRefinements), + inputPayload: JSON.stringify({ generationPurpose: "plan_update" }), + }); + await Promise.resolve(); + }); + + expect(mockConnectPlanningStream).toHaveBeenCalledTimes(connectionCountAfterResponse); + expect(screen.queryByText("The old stream failed")).toBeNull(); + }); + it("retains an actionable response error after durable question reconciliation", async () => { const submittedQuestion = { id: "q-submitted", @@ -1081,8 +1578,7 @@ describe("PlanningModeModal sequential flow", () => { }); renderSession(); - fireEvent.click(await screen.findByLabelText("Secure defaults")); - fireEvent.click(screen.getByRole("button", { name: "Next" })); + await selectResponseAfterHydration("Secure defaults"); expect(await screen.findByText("What should the durable session ask next?")).toBeInTheDocument(); expect(screen.getByText("Response submission timed out")).toBeInTheDocument();