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) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-16 00:59:48 -07:00
parent dd1e0f22e1
commit 7527d2651f
4 changed files with 614 additions and 23 deletions

View File

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

View File

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

View File

@@ -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. mutate the newly selected session, while successful progress still resets the three-attempt budget.
*/ */
const planningAutoRetryAttemptRef = useRef(0); const planningAutoRetryAttemptRef = useRef(0);
const planningAutoRetryOwnerRef = useRef<{ sessionId: string; token: symbol } | null>(null); const planningAutoRetryOwnerRef = useRef<{ sessionId: string; token: symbol; ownsTurn?: () => boolean } | null>(null);
const startPlanningAutoRetryRef = useRef<(sessionId: string) => Promise<boolean>>(async () => false); const startPlanningAutoRetryRef = useRef<(sessionId: string, ownsTurn?: () => boolean) => Promise<boolean>>(async () => false);
const planningSessionLoadEpochRef = useRef(0); 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: 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. 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 tick = async () => {
const sessionId = currentSessionIdRef.current; const sessionId = currentSessionIdRef.current;
if (!sessionId) return; 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 { try {
const session = await fetchAiSession(sessionId); const session = await fetchAiSession(sessionId);
if (cancelled || !session) return; if (cancelled || !session) return;
if (currentSessionIdRef.current !== sessionId) 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) { 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. // Recover a legacy or partially persisted plan when its question event was missed.
// New sequential turns normally settle with both result and currentQuestion. // New sequential turns normally settle with both result and currentQuestion.
resetPlanningAutoRetryBudget(); resetPlanningAutoRetryBudget();
@@ -1281,6 +1298,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
setView({ type: "plan_review", session: { sessionId, currentQuestion: null, summary }, summary }); setView({ type: "plan_review", session: { sessionId, currentQuestion: null, summary }, summary });
setStreamingOutput(""); setStreamingOutput("");
} else if (session.status === "awaiting_input" && session.currentQuestion) { } else if (session.status === "awaiting_input" && session.currentQuestion) {
planningTurnEpochRef.current += 1;
/* /*
FNXC:PlanningTurnReconciliation 2026-07-20-10:36: FNXC:PlanningTurnReconciliation 2026-07-20-10:36:
Missed SSE recovery must hydrate the server's entire interview turn together. Keeping 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(""); setStreamingOutput("");
} else if (session.status === "complete" && session.result) { } else if (session.status === "complete" && session.result) {
planningTurnEpochRef.current += 1;
const resume = resolveCompletePlanningResume(session); const resume = resolveCompletePlanningResume(session);
if (resume.kind === "unrecoverable") { if (resume.kind === "unrecoverable") {
setView({ setView({
@@ -1322,9 +1341,23 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
setStreamingOutput(""); setStreamingOutput("");
} else if (session.status === "error") { } else if (session.status === "error") {
const errorMessage = session.error || t("planning.sessionFailed2", "Session failed"); 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 (handled) return;
if (cancelled || currentSessionIdRef.current !== sessionId) return; if (!pollRecoveryStillOwnsTurn()) return;
/* /*
FNXC:PlanningRetry 2026-07-13-00:05: FNXC:PlanningRetry 2026-07-13-00:05:
Mirror the SSE onError terminal-error transition here: when this poll is the one that 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. passive stream catch-up event that overwrites a newer awaiting-input question.
*/ */
if (isAnsweredQuestion && editingQuestionIdRef.current !== normalizedQuestion.id) return; if (isAnsweredQuestion && editingQuestionIdRef.current !== normalizedQuestion.id) return;
planningTurnEpochRef.current += 1;
setIsRetrying(false); setIsRetrying(false);
resetPlanningAutoRetryBudget(); resetPlanningAutoRetryBudget();
setIsRefiningSummary(false); setIsRefiningSummary(false);
@@ -1539,6 +1573,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
}, },
onSummary: (summary) => { onSummary: (summary) => {
if (isStaleEvent()) return; if (isStaleEvent()) return;
planningTurnEpochRef.current += 1;
const normalizedSummary = normalizePlanningSummary(summary); const normalizedSummary = normalizePlanningSummary(summary);
setIsRetrying(false); setIsRetrying(false);
resetPlanningAutoRetryBudget(); resetPlanningAutoRetryBudget();
@@ -1578,6 +1613,20 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
}, },
onError: (message) => { onError: (message) => {
if (isStaleEvent()) return; 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."); const errorMessage = message || t("planning.sessionFailed", "Session failed while contacting the AI.");
// A single transient stream error (e.g. tab was backgrounded long // A single transient stream error (e.g. tab was backgrounded long
@@ -1588,7 +1637,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
(async () => { (async () => {
try { try {
const session = await fetchAiSession(sessionId); const session = await fetchAiSession(sessionId);
if (isStaleEvent()) return; if (!recoveryStillOwnsTurn()) return;
if ( if (
session && session &&
(session.status === "generating" || session.status === "awaiting_input") (session.status === "generating" || session.status === "awaiting_input")
@@ -1599,7 +1648,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
} catch { } catch {
// fall through to error view below // fall through to error view below
} }
if (isStaleEvent()) return; if (!recoveryStillOwnsTurn()) return;
/* /*
FNXC:PlanningRetry 2026-07-21-10:00: 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 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. 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; return;
} }
if (!recoveryStillOwnsTurn()) return;
setIsRetrying(false); setIsRetrying(false);
setIsAutoRetrying(false); setIsAutoRetrying(false);
setIsRefiningSummary(false); setIsRefiningSummary(false);
@@ -1647,7 +1697,7 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
const startPlanningRetry = useCallback( const startPlanningRetry = useCallback(
async ( async (
retryTarget: { sessionId: string; currentQuestion: PlanningQuestion | null; summary: PlanningSummary | null }, retryTarget: { sessionId: string; currentQuestion: PlanningQuestion | null; summary: PlanningSummary | null },
options: { auto: boolean; retryToken?: symbol }, options: { auto: boolean; retryToken?: symbol; ownsTurn?: () => boolean },
) => { ) => {
setError(null); setError(null);
setIsRetrying(!options.auto); setIsRetrying(!options.auto);
@@ -1664,7 +1714,8 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
await retryPlanningSession(retryTarget.sessionId, projectId); await retryPlanningSession(retryTarget.sessionId, projectId);
} catch (err) { } catch (err) {
const retryStillOwnsSession = () => currentSessionIdRef.current === retryTarget.sessionId 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; if (!retryStillOwnsSession()) return;
let retryError: unknown = err; 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 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 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; 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) { if (options.retryToken && planningAutoRetryOwnerRef.current?.token === options.retryToken) {
planningAutoRetryOwnerRef.current = null; planningAutoRetryOwnerRef.current = null;
} }
queueMicrotask(() => { queueMicrotask(() => {
if (currentSessionIdRef.current === retryTarget.sessionId) { if (currentSessionIdRef.current === retryTarget.sessionId) {
void startPlanningAutoRetryRef.current(retryTarget.sessionId); void startPlanningAutoRetryRef.current(retryTarget.sessionId, options.ownsTurn);
} }
}); });
return; return;
@@ -1812,9 +1864,15 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
); );
const startPlanningAutoRetry = useCallback( const startPlanningAutoRetry = useCallback(
async (sessionId: string) => { async (sessionId: string, ownsTurn?: () => boolean) => {
if (planningAutoRetryOwnerRef.current?.sessionId === sessionId) { if (!(ownsTurn?.() ?? true)) return false;
return true; 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; if (viewRef.current.type === "error") return false;
// FNXC:PlanningRetry 2026-07-22-21:00: budget is per-session and survives remounts. // 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}`); const retryToken = Symbol(`planning-auto-retry:${sessionId}:${attempt}`);
planningAutoRetryAttemptsBySession.set(sessionId, attempt); planningAutoRetryAttemptsBySession.set(sessionId, attempt);
planningAutoRetryAttemptRef.current = attempt; planningAutoRetryAttemptRef.current = attempt;
planningAutoRetryOwnerRef.current = { sessionId, token: retryToken }; planningAutoRetryOwnerRef.current = { sessionId, token: retryToken, ownsTurn };
setAutoRetryAttempt(attempt); setAutoRetryAttempt(attempt);
setIsAutoRetrying(true); setIsAutoRetrying(true);
await startPlanningRetry( await startPlanningRetry(
{ sessionId, currentQuestion: null, summary: null }, { sessionId, currentQuestion: null, summary: null },
{ auto: true, retryToken }, { auto: true, retryToken, ownsTurn }
); );
return true; return true;
}, },
@@ -2809,6 +2867,8 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
} }
setError(null); setError(null);
const responseTurnEpoch = ++planningTurnEpochRef.current;
const responseLoadEpoch = planningSessionLoadEpochRef.current;
// 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;
@@ -2872,14 +2932,21 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
} catch (err) { } catch (err) {
const errorMessage = getErrorMessage(err) || t("planning.failedSubmitResponse", "Failed to submit response"); const errorMessage = getErrorMessage(err) || t("planning.failedSubmitResponse", "Failed to submit response");
const isDuplicateResponseConflict = isDuplicateResponseGenerationConflict(err); const isDuplicateResponseConflict = isDuplicateResponseGenerationConflict(err);
const reconciliationStillOwnsTurn = () => currentSessionIdRef.current === sessionId
&& planningSessionLoadEpochRef.current === responseLoadEpoch
&& planningTurnEpochRef.current === responseTurnEpoch;
/* /*
FNXC:PlanningTurnReconciliation 2026-07-20-10:36: FNXC:PlanningTurnReconciliation 2026-08-16-06:22:
A rejected HTTP response is ambiguous: the server may have accepted the answer before A rejected response may have been durably accepted, so reconciliation rehydrates before
the connection failed. Rehydrate durable state before restoring the form. If it was not rolling back optimism. Its snapshot loses to a newer session load, reset/stop, response
accepted, roll back the optimistic answer so history and the active question still agree. 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 { try {
const persisted = await fetchAiSession(sessionId); const persisted = await fetchAiSession(sessionId);
if (!reconciliationStillOwnsTurn()) return;
if (persisted?.status === "awaiting_input" && !persisted.currentQuestion && persisted.result) { if (persisted?.status === "awaiting_input" && !persisted.currentQuestion && persisted.result) {
const history = parseConversationHistory(persisted.conversationHistory); const history = parseConversationHistory(persisted.conversationHistory);
const summary = normalizePlanningSummary(JSON.parse(persisted.result) as PlanningSummary); const summary = normalizePlanningSummary(JSON.parse(persisted.result) as PlanningSummary);
@@ -2922,8 +2989,10 @@ export function PlanningModeModal({ isOpen, onClose, onTaskCreated, onTasksCreat
return; return;
} }
} catch { } catch {
if (!reconciliationStillOwnsTurn()) return;
// Fall back to the known pre-submit turn and remove its optimistic answer. // Fall back to the known pre-submit turn and remove its optimistic answer.
} }
if (!reconciliationStillOwnsTurn()) return;
conversationHistoryRef.current = historyBeforeSubmit; conversationHistoryRef.current = historyBeforeSubmit;
setConversationHistory(historyBeforeSubmit); setConversationHistory(historyBeforeSubmit);
setError(errorMessage); setError(errorMessage);

View File

@@ -57,6 +57,19 @@ async function clickProceedAfterHydration() {
fireEvent.click(screen.getByRole("button", { name: "Proceed with plan" })); 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", () => { describe("PlanningModeModal sequential flow", () => {
beforeEach(() => { beforeEach(() => {
vi.useRealTimers(); vi.useRealTimers();
@@ -260,6 +273,129 @@ describe("PlanningModeModal sequential flow", () => {
expect(screen.queryByText("session-a stream failed")).toBeNull(); 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<string, unknown>) => 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<void>)();
});
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<string, unknown>) => 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<void>;
act(() => {
pendingPoll = (poll as () => Promise<void>)();
});
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 () => { it("automatically retries a resumed error discovered by the loading poll", async () => {
const sessionId = "polled-error-session"; const sessionId = "polled-error-session";
const intervalSpy = vi.spyOn(globalThis, "setInterval"); const intervalSpy = vi.spyOn(globalThis, "setInterval");
@@ -1043,7 +1179,11 @@ describe("PlanningModeModal sequential flow", () => {
renderSession(); renderSession();
fireEvent.click(await screen.findByLabelText("Secure defaults")); 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") { if (status === "awaiting_input") {
expect(await screen.findByText("What should the durable session ask next?")).toBeInTheDocument(); 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(); 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<string, unknown>) => 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(<PlanningModeModal {...props} resumeSessionId="session-a" />);
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(<PlanningModeModal {...props} resumeSessionId="session-b" />);
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<string, unknown>) => 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<string, unknown>) => 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<string, unknown>) => 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(<PlanningModeModal {...props} resumeSessionId="session-a" />);
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(<PlanningModeModal {...props} resumeSessionId="session-b" />);
expect(await screen.findByText("What did session B ask?")).toBeInTheDocument();
rerender(<PlanningModeModal {...props} resumeSessionId="session-a" />);
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<string, unknown>) => 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 () => { it("retains an actionable response error after durable question reconciliation", async () => {
const submittedQuestion = { const submittedQuestion = {
id: "q-submitted", id: "q-submitted",
@@ -1081,8 +1578,7 @@ describe("PlanningModeModal sequential flow", () => {
}); });
renderSession(); renderSession();
fireEvent.click(await screen.findByLabelText("Secure defaults")); await selectResponseAfterHydration("Secure defaults");
fireEvent.click(screen.getByRole("button", { name: "Next" }));
expect(await screen.findByText("What should the durable session ask next?")).toBeInTheDocument(); expect(await screen.findByText("What should the durable session ask next?")).toBeInTheDocument();
expect(screen.getByText("Response submission timed out")).toBeInTheDocument(); expect(screen.getByText("Response submission timed out")).toBeInTheDocument();