From e6df8d78d496cf9640bd5a1e9d7e8f4c2cc2adf7 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Tue, 29 Sep 2026 16:49:25 +0200 Subject: [PATCH] fix: bound a resumed session fork's wait, keep its intent while in flight (#11413) * fix: bound a resumed session fork's wait, keep its intent while the fork is in flight Co-Authored-By: Claude Opus 5.5 (1M context) * fix: wait for a session fork another request is creating instead of dropping it When the fork is refused as already being created by a creation whose id was lost, the session waits for it to show up among the user's workspaces and adopts it, rather than aborting the send. An 'already exists' answer is adopted like a duplicate key. Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) --- .../tests/fork_background.rs | 16 ++++++ .../sessions/sessionState.svelte.ts | 56 +++++++++++++++---- .../components/sessions/sessionState.test.ts | 48 ++++++++++++++++ frontend/src/lib/utils/forkCreation.ts | 10 +++- 4 files changed, 115 insertions(+), 15 deletions(-) diff --git a/backend/windmill-api-integration-tests/tests/fork_background.rs b/backend/windmill-api-integration-tests/tests/fork_background.rs index 0e82fb0d70..24ef3468ae 100644 --- a/backend/windmill-api-integration-tests/tests/fork_background.rs +++ b/backend/windmill-api-integration-tests/tests/fork_background.rs @@ -103,6 +103,22 @@ async fn test_fork_created_in_background(db: Pool) -> anyhow::Result<( .await?; assert_eq!(resp.status(), 404); + // An attempt whose server stopped heartbeating before its fork committed reads as failed. + let abandoned = "5c3e1b1e-0000-4000-8000-000000000000"; + sqlx::query( + "INSERT INTO workspace_fork_creation + (fork_workspace_id, parent_workspace_id, created_by, creation_id, heartbeat_at) + VALUES ('wm-fork-gone', 'test-workspace', 'test@windmill.dev', $1::uuid, + now() - interval '2 minutes')", + ) + .bind(abandoned) + .execute(&db) + .await?; + assert_eq!( + wait_for_fork(&client, &base_url, abandoned).await["status"], + "failed" + ); + // Without the flag the fork is created before the answer, which clients asking for a // background fork recognise a server without it by. sqlx::query("UPDATE workspace SET deleted = true WHERE id = 'wm-fork-bg-err'") diff --git a/frontend/src/lib/components/sessions/sessionState.svelte.ts b/frontend/src/lib/components/sessions/sessionState.svelte.ts index e81df02120..f032385d31 100644 --- a/frontend/src/lib/components/sessions/sessionState.svelte.ts +++ b/frontend/src/lib/components/sessions/sessionState.svelte.ts @@ -1293,17 +1293,26 @@ export function setGeneratedSessionSummary( return true } -// Create a new fork workspace via the API, refresh the user-workspaces -// store, and return the new fork id. Used by both the first-send commit -// path (commitSessionWorkspace) and the move-session-to-a-new-fork path -// in the unavailable-session banner. Returns undefined on failure (a -// user-facing toast is already emitted). +// A resumed attempt's polls fail for good once it is gone, so they are not waited out as long as a +// fresh creation's, whose failures a rolling deploy clears; a resume cut short by such a deploy +// falls back to requesting the fork, which then waits for it below. +const RESUMED_FORK_MAX_FAILING_MS = 10_000 +// How long a fork another request is still creating is waited for. +const IN_FLIGHT_FORK_WAIT_MS = 3 * 60 * 1000 + +// Make sure the fork exists, refresh the user-workspaces store, and return +// the fork id. Used by both the first-send commit path +// (commitSessionWorkspace) and the move-session-to-a-new-fork path in the +// unavailable-session banner. Returns undefined on failure (a user-facing +// toast is already emitted). // -// Self-heal: if `fork.id` is already present in the user-workspaces -// store, the previous create succeeded (whose response we apparently -// lost). Adopt it silently instead of re-POSTing — the API would -// otherwise reject with workspace_pkey. Likewise, if the API returns a -// duplicate-key error we refresh the store and adopt the existing row. +// Self-heal: a fork already present in the user-workspaces store was +// created by an earlier request whose response was lost, and is adopted +// without requesting it again; so is one the API reports as existing. A +// fork whose `creation_id` is set waits for that creation (started before a +// reload), and one the API reports as still being created, by a creation +// whose id was lost, is waited for until it shows up among the user's +// workspaces. `onCreationStarted` receives the id of a new creation. export async function materializeFork( fork: PendingFork, onCreationStarted?: (creationId: string) => void @@ -1313,7 +1322,12 @@ export async function materializeFork( let resumed = false if (fork.creation_id) { // An attempt that failed or is gone is requested again below. - resumed = await waitForForkCreation(fork.parent_workspace_id, fork.creation_id).then( + resumed = await waitForForkCreation( + fork.parent_workspace_id, + fork.creation_id, + undefined, + RESUMED_FORK_MAX_FAILING_MS + ).then( () => true, () => false ) @@ -1330,7 +1344,11 @@ export async function materializeFork( return fork.id } catch (e: any) { const msg = String(e?.body ?? e?.message ?? e) - if (/workspace_pkey|duplicate key/i.test(msg)) { + if (msg.includes('is already being created') && (await forkJoinedWorkspaces(fork.id))) { + sendUserToast(`Created fork ${fork.name}`) + return fork.id + } + if (/workspace_pkey|duplicate key|already exists/i.test(msg)) { // Self-heal: the create likely already succeeded. Refresh + adopt the // existing row. Guard this refresh — a second network failure here must // NOT rethrow out of materializeFork (callers rely on the @@ -1348,6 +1366,20 @@ export async function materializeFork( } } +// Whether the fork shows up among the user's workspaces within IN_FLIGHT_FORK_WAIT_MS. +async function forkJoinedWorkspaces(forkId: string): Promise { + const deadline = Date.now() + IN_FLIGHT_FORK_WAIT_MS + while (Date.now() < deadline) { + await new Promise((resolve) => setTimeout(resolve, 1500)) + const listed = await WorkspaceService.listUserWorkspaces().catch(() => undefined) + if (listed?.workspaces?.some((w) => w.id === forkId)) { + usersWorkspaceStore.set(listed) + return true + } + } + return false +} + // Re-assign a committed session to a different workspace. Used to rescue // sessions whose original workspace was deleted / archived / had access // revoked — the chat history (stored in IndexedDB keyed by session id) is diff --git a/frontend/src/lib/components/sessions/sessionState.test.ts b/frontend/src/lib/components/sessions/sessionState.test.ts index 62d2b8d39f..710028413e 100644 --- a/frontend/src/lib/components/sessions/sessionState.test.ts +++ b/frontend/src/lib/components/sessions/sessionState.test.ts @@ -141,6 +141,54 @@ describe('commitSessionWorkspace — fork creation resumed after a reload', () = }) }) +describe('commitSessionWorkspace — fork still being created', () => { + it('gives up a gone creation quickly, then adopts the fork once it is created', async () => { + const id = 'test-commit-fork-in-flight' + const prevLicense = get(enterpriseLicense) + const prevWorkspaces = get(usersWorkspaceStore) + enterpriseLicense.set('test-license') + vi.useFakeTimers() + const status = vi.mocked(WorkspaceService.getForkCreationStatus) + status.mockClear() + status.mockRejectedValue(Object.assign(new Error('Not Found'), { status: 404 })) + vi.mocked(WorkspaceService.createWorkspaceFork).mockRejectedValue({ + body: "Bad request: workspace 'wm-fork-flight' is already being created" + }) + vi.mocked(WorkspaceService.listUserWorkspaces) + .mockResolvedValueOnce({ email: 't@t', workspaces: [] } as never) + .mockResolvedValue({ email: 't@t', workspaces: [ws('wm-fork-flight', 'parent_ws')] } as never) + sessionState.sessions.push({ + id, + name: 'fork-flight', + createdAt: 0, + pending_fork: { + parent_workspace_id: 'parent_ws', + id: 'wm-fork-flight', + name: 'flight', + creation_id: '5c3e1b1e-0000-4000-8000-000000000000' + } + } as Session) + try { + const committed = commitSessionWorkspace(id, 'parent_ws') + await vi.runAllTimersAsync() + expect(await committed).toBe('wm-fork-flight') + // The resumed creation is given up after its own short budget, not a fresh creation's. + expect(status.mock.calls.length).toBeLessThan(10) + } finally { + vi.useRealTimers() + status.mockResolvedValue({ status: 'completed' }) + vi.mocked(WorkspaceService.createWorkspaceFork).mockRejectedValue( + new Error('fork creation failed') + ) + vi.mocked(WorkspaceService.listUserWorkspaces).mockResolvedValue([] as never) + usersWorkspaceStore.set(prevWorkspaces) + const i = sessionState.sessions.findIndex((x) => x.id === id) + if (i >= 0) sessionState.sessions.splice(i, 1) + enterpriseLicense.set(prevLicense) + } + }) +}) + // Set the workspace list to a given number of non-'admins' workspaces so the // commit-path's CE workspace-cap check (mirror of backend _check_nb_of_workspaces) // can be exercised. Returns a restore fn. diff --git a/frontend/src/lib/utils/forkCreation.ts b/frontend/src/lib/utils/forkCreation.ts index d0fb4e5edf..aadd67e70f 100644 --- a/frontend/src/lib/utils/forkCreation.ts +++ b/frontend/src/lib/utils/forkCreation.ts @@ -28,11 +28,15 @@ export async function createWorkspaceForkAndWait( await waitForForkCreation(parentWorkspace, response, opts.onStep) } -/** Wait for a fork creation started in the background, by the id its request answered with. */ +/** + * Wait for a fork creation started in the background, by the id its request answered with. + * `maxFailingMs` bounds how long polls may keep failing before the wait gives up. + */ export async function waitForForkCreation( parentWorkspace: string, creationId: string, - onStep?: (step: string) => void + onStep?: (step: string) => void, + maxFailingMs = MAX_FAILING_POLLS_MS ): Promise { let failingSince: number | undefined while (true) { @@ -47,7 +51,7 @@ export async function waitForForkCreation( // A failed poll says nothing about the fork: it can be lost to the same proxy the fork // runs in the background to avoid, or reach a replica that has no status route yet. failingSince ??= Date.now() - if (Date.now() - failingSince >= MAX_FAILING_POLLS_MS) throw e + if (Date.now() - failingSince >= maxFailingMs) throw e } if (result?.status === 'running' && result.step) onStep?.(result.step) if (result?.status === 'completed') return