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