From 8938aaecce1a846a90c472dfe9bfb7999de70d98 Mon Sep 17 00:00:00 2001 From: Yousef Abdelhadi Date: Sat, 26 Sep 2026 04:46:30 -0400 Subject: [PATCH 1/4] fix(gen2): harden workspace lifecycle recovery --- .../[workspaceId]/instance/route.ts | 2 +- apps/web/lib/gen2/instance-lifecycle.test.ts | 135 ++++++++++- apps/web/lib/gen2/instance.ts | 170 ++++++++++--- apps/web/lib/runtime/azure-host.ts | 4 +- apps/web/lib/runtime/orchestrator-health.ts | 18 +- docs/gen2-workspace.md | 7 +- .../orchestrator/src/backend/firecracker.rs | 225 +++++++++++++++--- 7 files changed, 474 insertions(+), 87 deletions(-) diff --git a/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.ts b/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.ts index 52feb907c..a7c4fa1a9 100644 --- a/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.ts +++ b/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.ts @@ -2,7 +2,7 @@ import { withUser } from "@/lib/http/api-route"; import { ensureGen2Instance, stopGen2Instance } from "@/lib/gen2/instance"; import { getGen2WorkspaceDetail } from "@/lib/gen2/workspaces"; -/** Provisioning a Firecracker guest can outlast a normal request. */ +/** Host wake, Firecracker creation, and persistence fit Vercel's 300s limit. */ export const maxDuration = 300; type Params = { workspaceId: string }; diff --git a/apps/web/lib/gen2/instance-lifecycle.test.ts b/apps/web/lib/gen2/instance-lifecycle.test.ts index 6d450b238..37a7a85e0 100644 --- a/apps/web/lib/gen2/instance-lifecycle.test.ts +++ b/apps/web/lib/gen2/instance-lifecycle.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; const mocks = vi.hoisted(() => { const updates: Array> = []; @@ -22,7 +22,7 @@ const mocks = vi.hoisted(() => { // `.returning()`. `claimed` says whether the compare-and-set matched a row; // when it does not, the write must not land -- that is the whole point of // the claim, and modelling it is how the race test means anything. - const state = { claimed: true }; + const state = { claimed: true, leaseCurrent: true }; let pending: Record | null = null; function apply() { @@ -37,6 +37,9 @@ const mocks = vi.hoisted(() => { if ("lastError" in values) { member.lastError = values.lastError as string | null; } + if ("updatedAt" in values) { + member.updatedAt = values.updatedAt as Date; + } } const updateQuery = { @@ -46,12 +49,16 @@ const mocks = vi.hoisted(() => { }), where: vi.fn(() => updateQuery), returning: vi.fn(async () => { - if (!state.claimed) { + const status = pending?.status; + if ( + (status === "provisioning" && !state.claimed) || + (status !== "provisioning" && !state.leaseCurrent) + ) { pending = null; return []; } apply(); - return [{ id: member.id }]; + return [{ id: member.id, updatedAt: member.updatedAt }]; }), then: (resolve: (value: undefined) => unknown) => { apply(); @@ -96,13 +103,19 @@ import { Gen2LifecycleError } from "./errors"; import { ensureGen2Instance, stopGen2Instance } from "./instance"; describe("gen2 instance lifecycle", () => { + afterEach(() => { + vi.useRealTimers(); + }); + beforeEach(() => { mocks.updates.length = 0; mocks.member.status = "pending"; mocks.member.sandboxId = null; mocks.member.lastError = null; + mocks.member.updatedAt = new Date("2026-09-20T20:00:00.000Z"); mocks.member.role = "owner"; mocks.state.claimed = true; + mocks.state.leaseCurrent = true; for (const method of ["from", "innerJoin", "where"] as const) { mocks.selectQuery[method].mockReturnValue(mocks.selectQuery); } @@ -130,7 +143,7 @@ describe("gen2 instance lifecycle", () => { ]); expect(workspace.status).toBe("ready"); expect(workspace.sandboxId).toBe("sandbox-1"); - expect(mocks.ensureHostReady).toHaveBeenCalledOnce(); + expect(mocks.ensureHostReady).toHaveBeenCalledWith(150_000); }); it("reattaches when the Firecracker machine is already on the host", async () => { @@ -196,6 +209,43 @@ describe("gen2 instance lifecycle", () => { }); }); + it("recovers a stale provisioning claim", async () => { + mocks.member.status = "provisioning"; + mocks.member.updatedAt = new Date(Date.now() - 10 * 60_000); + + const workspace = await ensureGen2Instance(mocks.member.id, "user-1", { + provision: mocks.provision, + destroy: mocks.destroy, + }); + + expect(mocks.provision).toHaveBeenCalledOnce(); + expect(mocks.updates.map((update) => update.status)).toEqual([ + "provisioning", + "ready", + ]); + expect(workspace.status).toBe("ready"); + }); + + it("records a retryable failure when host wake fails during stale recovery", async () => { + mocks.member.status = "provisioning"; + mocks.member.updatedAt = new Date(Date.now() - 10 * 60_000); + mocks.ensureHostReady.mockRejectedValue(new TypeError("fetch failed")); + + await expect( + ensureGen2Instance(mocks.member.id, "user-1", { + provision: mocks.provision, + destroy: mocks.destroy, + }), + ).rejects.toMatchObject({ status: 503 }); + + expect(mocks.provision).not.toHaveBeenCalled(); + expect(mocks.member.status).toBe("failed"); + expect(mocks.updates.at(-1)).toMatchObject({ + status: "failed", + lastError: expect.stringMatching(/Firecracker host could not be reached/), + }); + }); + it("does nothing when the machine is already running", async () => { mocks.member.status = "ready"; mocks.member.sandboxId = "sandbox-1"; @@ -272,14 +322,83 @@ describe("gen2 instance lifecycle", () => { expect(mocks.provision).toHaveBeenCalledOnce(); }); - it("steps aside when another member already claimed provisioning", async () => { + it("waits for another member's provisioning attempt to finish", async () => { + vi.useFakeTimers(); + mocks.member.status = "provisioning"; + mocks.member.updatedAt = new Date(); mocks.state.claimed = false; - const workspace = await ensureGen2Instance(mocks.member.id, "user-2", { + let reads = 0; + mocks.selectQuery.limit.mockImplementation(async () => { + reads += 1; + if (reads === 3) { + mocks.member.status = "ready"; + mocks.member.sandboxId = "sandbox-1"; + } + return [mocks.member]; + }); + + const opening = ensureGen2Instance(mocks.member.id, "user-2", { provision: mocks.provision, destroy: mocks.destroy, }); + + await vi.advanceTimersByTimeAsync(1_000); + const workspace = await opening; + expect(mocks.provision).not.toHaveBeenCalled(); - expect(workspace.status).toBe("pending"); + expect(reads).toBeGreaterThanOrEqual(3); + expect(workspace.status).toBe("ready"); + }); + + it("does not let an old provisioner overwrite a newer startup lease", async () => { + const newerLease = new Date(Date.now() + 1_000); + mocks.provision.mockImplementation(async () => { + mocks.member.status = "provisioning"; + mocks.member.updatedAt = newerLease; + mocks.state.leaseCurrent = false; + return { id: "sandbox-1" }; + }); + + await expect( + ensureGen2Instance(mocks.member.id, "user-1", { + provision: mocks.provision, + destroy: mocks.destroy, + }), + ).rejects.toMatchObject({ + message: /newer startup attempt took over/, + status: 503, + }); + + expect(mocks.member.status).toBe("provisioning"); + expect(mocks.member.updatedAt).toBe(newerLease); + expect(mocks.updates.map((update) => update.status)).toEqual([ + "provisioning", + ]); + }); + + it("reopens the machine after the owner stops it", async () => { + mocks.member.status = "ready"; + mocks.member.sandboxId = "sandbox-1"; + const runtime = { + provision: mocks.provision, + destroy: mocks.destroy, + }; + + await stopGen2Instance(mocks.member.id, "user-1", runtime); + const workspace = await ensureGen2Instance( + mocks.member.id, + "user-1", + runtime, + ); + + expect(mocks.destroy).toHaveBeenCalledOnce(); + expect(mocks.provision).toHaveBeenCalledOnce(); + expect(mocks.updates.map((update) => update.status)).toEqual([ + "stopped", + "provisioning", + "ready", + ]); + expect(workspace.status).toBe("ready"); }); it("destroys the sandbox when the owner stops it", async () => { diff --git a/apps/web/lib/gen2/instance.ts b/apps/web/lib/gen2/instance.ts index 8fe5e339a..9c15c368e 100644 --- a/apps/web/lib/gen2/instance.ts +++ b/apps/web/lib/gen2/instance.ts @@ -1,6 +1,6 @@ import "server-only"; -import { and, eq, inArray } from "drizzle-orm"; +import { and, eq, inArray, lt, or } from "drizzle-orm"; import { schema } from "@codev/db"; import type { Gen2WorkspaceStatus } from "@codev/contracts"; @@ -62,8 +62,17 @@ const HOST_UNREACHABLE_MESSAGE = /** Firecracker create can outlast the default 70s orchestrator timeout. */ const GEN2_PROVISION_TIMEOUT_MS = 120_000; +/** Leave room for guest creation within Vercel's 300-second function limit. */ +const GEN2_HOST_READY_TIMEOUT_MS = 150_000; /** Keep expiry inside the orchestrator's exclusive four-hour window. */ const GEN2_EXPIRES_SLACK_MS = 60_000; +/** + * The host wake and guest creation budgets total 270 seconds. Leave another + * thirty seconds for source preparation and persisting the final state. + */ +const GEN2_PROVISIONING_STALE_AFTER_MS = + GEN2_HOST_READY_TIMEOUT_MS + GEN2_PROVISION_TIMEOUT_MS + 30_000; +const GEN2_PROVISIONING_POLL_MS = 1_000; export function describeGen2RuntimeFailure(error: unknown): string { if (isGen2HostUnreachable(error)) { @@ -162,14 +171,30 @@ async function writeGen2Instance( sandboxId?: string | null; lastError: string | null; }, + expectedProvisioningAt?: Date, ) { - await getDatabase() + const update = getDatabase() .update(schema.gen2Workspaces) .set({ ...values, updatedAt: new Date(), }) - .where(eq(schema.gen2Workspaces.id, workspaceId)); + .where( + expectedProvisioningAt + ? and( + eq(schema.gen2Workspaces.id, workspaceId), + eq(schema.gen2Workspaces.status, "provisioning"), + eq(schema.gen2Workspaces.updatedAt, expectedProvisioningAt), + ) + : eq(schema.gen2Workspaces.id, workspaceId), + ); + if (expectedProvisioningAt) { + return ( + (await update.returning({ id: schema.gen2Workspaces.id })).length > 0 + ); + } + await update; + return true; } /** @@ -177,25 +202,72 @@ async function writeGen2Instance( * * Two members opening the same workspace at once would otherwise both see * `pending` and both provision. The compare-and-set means exactly one wins; - * the loser gets `null` and simply waits for the winner's VM. + * the loser gets `null` and waits for the winner to finish starting the VM. */ async function claimProvisioning(workspaceId: string) { - const claimed = await getDatabase() + const now = new Date(); + const staleBefore = new Date( + now.getTime() - GEN2_PROVISIONING_STALE_AFTER_MS, + ); + const [claim] = await getDatabase() .update(schema.gen2Workspaces) - .set({ status: "provisioning", lastError: null, updatedAt: new Date() }) + .set({ status: "provisioning", lastError: null, updatedAt: now }) .where( and( eq(schema.gen2Workspaces.id, workspaceId), - inArray(schema.gen2Workspaces.status, [ - "pending", - "stopped", - "failed", - "ready", - ]), + or( + inArray(schema.gen2Workspaces.status, [ + "pending", + "stopped", + "failed", + "ready", + ]), + and( + eq(schema.gen2Workspaces.status, "provisioning"), + lt(schema.gen2Workspaces.updatedAt, staleBefore), + ), + ), ), ) - .returning({ id: schema.gen2Workspaces.id }); - return claimed.length > 0; + .returning({ updatedAt: schema.gen2Workspaces.updatedAt }); + return claim?.updatedAt ?? null; +} + +async function waitForProvisioning(workspaceId: string, userId: string) { + const deadline = Date.now() + GEN2_PROVISIONING_STALE_AFTER_MS; + while (Date.now() < deadline) { + const workspace = await requireGen2Member(workspaceId, userId); + if (workspace.status === "ready") return workspace; + if (workspace.status === "failed") { + throw new Gen2LifecycleError( + workspace.lastError ?? "The Firecracker instance could not start.", + 502, + ); + } + if (workspace.status !== "provisioning") { + throw new Gen2LifecycleError( + workspace.lastError ?? + "Workspace startup ended before the machine was ready.", + 503, + ); + } + if ( + Date.now() - Date.parse(workspace.updatedAt) >= + GEN2_PROVISIONING_STALE_AFTER_MS + ) { + throw new Gen2LifecycleError( + "The previous startup attempt stopped responding. Try again to resume this workspace.", + 503, + ); + } + await new Promise((resolve) => + setTimeout(resolve, GEN2_PROVISIONING_POLL_MS), + ); + } + throw new Gen2LifecycleError( + "The Firecracker instance is still starting. Try again in a moment.", + 503, + ); } /** @@ -221,7 +293,7 @@ export async function ensureGen2Instance( if (membership.status === "ready") { if (!currentRuntime.current) return membership; try { - await ensureHostReady(); + await ensureHostReady(GEN2_HOST_READY_TIMEOUT_MS); hostReady = true; if (await currentRuntime.current(workspaceId)) return membership; } catch (error) { @@ -233,25 +305,35 @@ export async function ensureGen2Instance( } } - if (!(await claimProvisioning(workspaceId))) { - // Someone else is already bringing it up; report the live status rather - // than racing them for the same guest. - return requireGen2Member(workspaceId, userId); + const provisioningAt = await claimProvisioning(workspaceId); + if (!provisioningAt) { + // Another member owns the startup lease. Join that operation instead of + // returning a successful response that leaves this browser stuck at + // "Starting" without an active runtime. + return waitForProvisioning(workspaceId, userId); } - const previousStatus = membership.status; + // A process can be terminated by its platform before its catch block runs. + // A stale provisioning lease is recoverable, so if host wake fails during + // recovery leave a retryable failure instead of refreshing the stuck lease. + const previousStatus = + membership.status === "provisioning" ? "failed" : membership.status; if (!hostReady) { try { - await ensureHostReady(); + await ensureHostReady(GEN2_HOST_READY_TIMEOUT_MS); } catch (error) { const message = describeGen2RuntimeFailure(error); logEvent("error", "gen2.instance.host_unready", { detail: error instanceof Error ? error.message : "unknown", }); - await writeGen2Instance(workspaceId, { - status: previousStatus, - lastError: message, - }); + await writeGen2Instance( + workspaceId, + { + status: previousStatus, + lastError: message, + }, + provisioningAt, + ); throw new Gen2LifecycleError(message, 503); } } @@ -278,20 +360,40 @@ export async function ensureGen2Instance( Date.now() + GEN2_SANDBOX_LIFECYCLE.timeoutMs - GEN2_EXPIRES_SLACK_MS, ), )); - await writeGen2Instance(workspaceId, { - status: "ready", - sandboxId: sandbox.id, - lastError: null, - }); + const committed = await writeGen2Instance( + workspaceId, + { + status: "ready", + sandboxId: sandbox.id, + lastError: null, + }, + provisioningAt, + ); + if (!committed) { + throw new Gen2LifecycleError( + "A newer startup attempt took over. Try opening the workspace again.", + 503, + ); + } } catch (error) { const message = describeGen2RuntimeFailure(error); logEvent("error", "gen2.instance.start_failed", { detail: error instanceof Error ? error.message : "unknown", }); - await writeGen2Instance(workspaceId, { - status: "failed", - lastError: message, - }); + const markedFailed = await writeGen2Instance( + workspaceId, + { + status: "failed", + lastError: message, + }, + provisioningAt, + ); + if (!markedFailed) { + throw new Gen2LifecycleError( + "A newer startup attempt took over. Try opening the workspace again.", + 503, + ); + } throw new Gen2LifecycleError(message, 502); } diff --git a/apps/web/lib/runtime/azure-host.ts b/apps/web/lib/runtime/azure-host.ts index cf2f2a3a6..08192819a 100644 --- a/apps/web/lib/runtime/azure-host.ts +++ b/apps/web/lib/runtime/azure-host.ts @@ -167,8 +167,8 @@ export async function getHostState(): Promise { * right trade for a caller with no poll of its own, and the wrong one for the * workspace open path, which passes a small number and reports `host-starting` * instead. A returning member hits `stopping` routinely: the idle timer - * deallocates at ten minutes, so coming back a moment later lands squarely on - * it. + * deallocates after a one-minute quiet window, so coming back a moment later + * lands squarely on it. */ const DEFAULT_STOPPING_ATTEMPTS = 30; diff --git a/apps/web/lib/runtime/orchestrator-health.ts b/apps/web/lib/runtime/orchestrator-health.ts index 3a135c8bb..48c503f1c 100644 --- a/apps/web/lib/runtime/orchestrator-health.ts +++ b/apps/web/lib/runtime/orchestrator-health.ts @@ -49,16 +49,16 @@ async function parseHealth(response: Response) { /** * Wait until the Firecracker host is running *and* its orchestrator answers. * - * The host stops itself after ten minutes idle, so the first call after any - * quiet period lands on a stopped instance. Starting it takes roughly ten - * seconds before the orchestrator is even up, and longer before it serves -- - * far longer than a single provision attempt is willing to wait. Callers that - * skip this see "Firecracker host unavailable" on the first click and success - * on the second, which is the whole of that bug. + * The Azure Firecracker host deallocates after one quiet minute, so the first + * call after a quiet period can land on a stopped instance. Starting it takes + * roughly ten seconds before the orchestrator is even up, and longer before + * it serves -- far longer than a single provision attempt is willing to wait. + * Callers that skip this see "Firecracker host unavailable" on the first + * click and success on the second, which is the whole of that bug. * - * `requestHostWake` absorbs transient EC2 failures itself and reports the host - * as starting, so a capacity refusal or a mid-restart instance costs another - * turn of this loop rather than failing the action outright. + * `requestHostWake` handles transient Azure start failures and reports the + * host as starting, so a capacity refusal or a mid-restart instance costs + * another turn of this loop rather than failing the action outright. */ export async function ensureHostReady(timeoutMs = HOST_START_TIMEOUT_MS) { // The local stand-in has no host to wake. diff --git a/docs/gen2-workspace.md b/docs/gen2-workspace.md index 253db6532..cbe46c2f9 100644 --- a/docs/gen2-workspace.md +++ b/docs/gen2-workspace.md @@ -45,9 +45,10 @@ Opening a workspace is the intent to use it, so `ensureGen2Instance` runs on open and is safe for any member to call -- the person who follows a share link should not have to wait for the owner to press something. The composer is never disabled either: type into a cold workspace and the machine is brought -up as part of sending. The orchestrator hibernates an idle guest after 15 -minutes, preserving its workspace disk while releasing its sandbox slot. A -host with no active sandboxes deallocates after one quiet minute. Opening the +up as part of sending. The orchestrator hibernates an idle guest after fifteen +minutes, preserving its workspace disk while releasing its sandbox slot. Once +no guest or recently used Orca IDE session keeps the host active, it deallocates +after a one-minute quiet window checked every thirty seconds. Opening the workspace resumes it from the saved disk. ## Interface diff --git a/services/orchestrator/src/backend/firecracker.rs b/services/orchestrator/src/backend/firecracker.rs index 5f34e1908..be8830877 100644 --- a/services/orchestrator/src/backend/firecracker.rs +++ b/services/orchestrator/src/backend/firecracker.rs @@ -365,6 +365,54 @@ impl FirecrackerApiClient { } } +async fn load_snapshot_while_firecracker_runs( + api: &FirecrackerApiClient, + child: &Mutex, +) -> std::result::Result<(), SnapshotRestoreFailure> { + tokio::select! { + result = api.load_snapshot() => result.map_err(SnapshotRestoreFailure::Api), + exit = async { + let mut child = child.lock().await; + child.wait().await + } => { + match exit { + Ok(status) => Err(SnapshotRestoreFailure::FirecrackerExited(status)), + Err(error) => Err(SnapshotRestoreFailure::Api(RuntimeError::internal(error))), + } + } + } +} + +enum SnapshotRestoreFailure { + Api(RuntimeError), + FirecrackerExited(std::process::ExitStatus), +} + +async fn restore_snapshot_while_firecracker_runs( + api: &FirecrackerApiClient, + child: &Mutex, + timeout_duration: Duration, +) -> Result<()> { + timeout(timeout_duration, async { + loop { + match load_snapshot_while_firecracker_runs(api, child).await { + Ok(()) => return Ok(()), + Err(SnapshotRestoreFailure::Api(RuntimeError::Unavailable(_))) => { + sleep(Duration::from_millis(100)).await + } + Err(SnapshotRestoreFailure::Api(error)) => return Err(error), + Err(SnapshotRestoreFailure::FirecrackerExited(status)) => { + return Err(RuntimeError::Unavailable(format!( + "Firecracker exited before snapshot restore: {status}" + ))); + } + } + } + }) + .await + .map_err(|_| RuntimeError::Timeout("Firecracker snapshot restore timed out".into()))? +} + async fn read_until_sequence( stream: &mut UnixStream, sequence: &[u8], @@ -1179,6 +1227,37 @@ impl FirecrackerBackend { } async fn prepare_and_start(&self, request: &CreateRequest) -> Result { + let mut restore_from_saved_disks = false; + loop { + let mut full_snapshot_restore_failed = false; + match self + .prepare_and_start_attempt( + request, + restore_from_saved_disks, + &mut full_snapshot_restore_failed, + ) + .await + { + Ok(machine) => return Ok(machine), + Err(error) if full_snapshot_restore_failed && !restore_from_saved_disks => { + warn!( + workspace_id = %request.workspace_id, + %error, + "full Firecracker snapshot restore failed; retrying from saved workspace disks" + ); + restore_from_saved_disks = true; + } + Err(error) => return Err(error), + } + } + } + + async fn prepare_and_start_attempt( + &self, + request: &CreateRequest, + restore_from_saved_disks: bool, + full_snapshot_restore_failed: &mut bool, + ) -> Result { let snapshot_state = if request.resume_from_snapshot && request.persistent_disk_lun.is_none() { self.snapshot_metadata(&request.workspace_id).await? @@ -1190,10 +1269,16 @@ impl FirecrackerBackend { .as_ref() .map(|(directory, _)| directory.clone()) .unwrap_or_else(|| self.snapshot_dir(&request.workspace_id)); - let restore_snapshot = - snapshot_metadata.is_some_and(|metadata| metadata.kind == SnapshotKind::FullMachine); - let disk_checkpoint = - snapshot_metadata.filter(|metadata| metadata.kind == SnapshotKind::WorkspaceDisks); + let restore_snapshot = !restore_from_saved_disks + && snapshot_metadata.is_some_and(|metadata| metadata.kind == SnapshotKind::FullMachine); + let disk_checkpoint = snapshot_metadata.filter(|metadata| { + restore_from_saved_disks || metadata.kind == SnapshotKind::WorkspaceDisks + }); + if restore_from_saved_disks && disk_checkpoint.is_none() { + return Err(RuntimeError::Unavailable( + "saved Firecracker workspace disks are unavailable for recovery".into(), + )); + } let slot = { let machines = self.machines.read().await; if restore_snapshot { @@ -1390,23 +1475,15 @@ impl FirecrackerBackend { if restore_snapshot { let restore_started_at = Instant::now(); - let restore = timeout(Duration::from_secs(45), async { - loop { - match FirecrackerApiClient::new(machine.api_socket.clone()) - .load_snapshot() - .await - { - Ok(()) => break Ok(()), - Err(RuntimeError::Unavailable(_)) => { - sleep(Duration::from_millis(100)).await - } - Err(error) => break Err(error), - } - } - }) - .await; - match restore { - Ok(Ok(())) => { + let api = FirecrackerApiClient::new(machine.api_socket.clone()); + match restore_snapshot_while_firecracker_runs( + &api, + &machine.child, + Duration::from_secs(45), + ) + .await + { + Ok(()) => { let restore_ms = restore_started_at.elapsed().as_millis() as u64; if restore_ms > 500 { warn!( @@ -1423,16 +1500,11 @@ impl FirecrackerBackend { ); } } - Ok(Err(error)) => { + Err(error) => { self.cleanup_failed_machine(&machine).await; + *full_snapshot_restore_failed = true; return Err(error); } - Err(_) => { - self.cleanup_failed_machine(&machine).await; - return Err(RuntimeError::Timeout( - "Firecracker snapshot restore timed out".into(), - )); - } } } @@ -2015,8 +2087,21 @@ pub fn parse_duration(value: &str) -> Option { #[cfg(test)] mod tests { use super::{ - GUEST_CID_BASE, MicroVmSnapshotMetadata, SnapshotKind, first_available_slot, guest_ip, - has_reaper_blocking_requests, host_ip, parse_duration, tap_name, + FirecrackerApiClient, FirecrackerBackend, FirecrackerConfig, GUEST_CID_BASE, + MicroVmSnapshotMetadata, SnapshotKind, first_available_slot, guest_ip, + has_reaper_blocking_requests, host_ip, parse_duration, + restore_snapshot_while_firecracker_runs, tap_name, + }; + use crate::model::RuntimeError; + use std::{ + collections::HashMap, + path::PathBuf, + sync::atomic::{AtomicBool, Ordering}, + time::{Duration, Instant}, + }; + use tokio::{ + process::Command, + sync::{Mutex, RwLock as AsyncRwLock}, }; #[test] @@ -2077,4 +2162,84 @@ mod tests { assert_eq!(metadata.kind, SnapshotKind::FullMachine); assert_eq!(metadata.slot, 3); } + + #[tokio::test] + async fn exits_during_snapshot_restore_are_reported_immediately() { + let child = Command::new("sh") + .args(["-c", "exit 17"]) + .spawn() + .expect("spawn an exiting Firecracker stand-in"); + let child = Mutex::new(child); + let api = FirecrackerApiClient::new( + tempfile::tempdir() + .expect("temporary directory") + .path() + .join("missing-api.socket"), + ); + let started_at = Instant::now(); + + let error = restore_snapshot_while_firecracker_runs(&api, &child, Duration::from_secs(45)) + .await + .expect_err("an exited Firecracker process cannot restore a snapshot"); + + assert!(matches!( + error, + RuntimeError::Unavailable(message) + if message.contains("Firecracker exited before snapshot restore") + && message.contains("17") + )); + assert!(started_at.elapsed() < Duration::from_secs(2)); + } + + #[tokio::test] + async fn shutdown_reservation_blocks_create_on_the_firecracker_backend() { + let backend = FirecrackerBackend { + config: FirecrackerConfig { + runtime_dir: PathBuf::new(), + kernel_image: PathBuf::new(), + rootfs_image: PathBuf::new(), + firecracker_bin: PathBuf::new(), + jailer_bin: PathBuf::new(), + jailer_dir: PathBuf::new(), + max_sandboxes: 1, + vcpu_count: 1, + memory_mib: 256, + workspace_disk_gib: 1, + idle_timeout: Duration::from_secs(60), + guest_network: false, + }, + machines: AsyncRwLock::new(HashMap::new()), + provision: Mutex::new(()), + host_shutdown_pending: AtomicBool::new(false), + }; + + assert!(backend.begin_host_shutdown_if_idle().await); + let request = serde_json::from_value(serde_json::json!({ + "workspaceId": "lifecycle-test", + "repositoryUrl": null, + "repositorySnapshot": { + "files": [], + "totalBytes": 0 + }, + "baseSha": "0000000000000000000000000000000000000000", + "expiresAt": "2026-09-26T08:00:00Z", + "resumeFromSnapshot": false, + "lifecycle": { + "timeoutMs": 14400000, + "lifecycle": { "onTimeout": "pause", "autoResume": true } + } + })) + .expect("valid create request"); + let error = backend + .create(request) + .await + .expect_err("create must not race host deallocation"); + + assert!(matches!( + error, + RuntimeError::Unavailable(message) if message.contains("host is shutting down") + )); + backend.cancel_host_shutdown(); + assert!(!backend.host_shutdown_pending.load(Ordering::Acquire)); + } } From 7d3cf858f4a7fe79da3f96f95081be7bf6822d1a Mon Sep 17 00:00:00 2001 From: Yousef Abdelhadi Date: Sat, 26 Sep 2026 06:04:42 -0400 Subject: [PATCH 2/4] fix(gen2): retry startup and serialize stop --- .../[workspaceId]/instance/route.test.ts | 120 +++++++++++++ .../[workspaceId]/instance/route.ts | 12 +- .../components/gen2/workspace-room.test.tsx | 4 +- apps/web/components/gen2/workspace-room.tsx | 48 +++--- apps/web/lib/gen2/instance-lifecycle.test.ts | 54 ++++-- apps/web/lib/gen2/instance.test.ts | 4 +- apps/web/lib/gen2/instance.ts | 122 +++++++------ apps/web/lib/gen2/startup-client.test.ts | 140 +++++++++++++++ apps/web/lib/gen2/startup-client.ts | 160 ++++++++++++++++++ apps/web/lib/gen2/workspaces-delete.test.ts | 33 ++-- apps/web/lib/gen2/workspaces.test.ts | 8 +- apps/web/lib/gen2/workspaces.ts | 29 +--- .../lib/runtime/orchestrator-health.test.ts | 49 ++++++ apps/web/lib/runtime/orchestrator-health.ts | 4 +- 14 files changed, 644 insertions(+), 143 deletions(-) create mode 100644 apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.test.ts create mode 100644 apps/web/lib/gen2/startup-client.test.ts create mode 100644 apps/web/lib/gen2/startup-client.ts create mode 100644 apps/web/lib/runtime/orchestrator-health.test.ts diff --git a/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.test.ts b/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.test.ts new file mode 100644 index 000000000..ae84d81be --- /dev/null +++ b/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.test.ts @@ -0,0 +1,120 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + ensure: vi.fn(), + detail: vi.fn(), +})); + +vi.mock("@/lib/http/api-route", () => ({ + withUser: ( + handler: (input: { + user: { id: string }; + params: { workspaceId: string }; + }) => Promise, + ) => { + return async ( + _request: Request, + context: { params: Promise<{ workspaceId: string }> }, + ) => + handler({ + user: { id: "user-1" }, + params: await context.params, + }); + }, +})); +vi.mock("@/lib/gen2/instance", () => ({ + ensureGen2Instance: mocks.ensure, + stopGen2Instance: vi.fn(), +})); +vi.mock("@/lib/gen2/workspaces", () => ({ + getGen2WorkspaceDetail: mocks.detail, +})); + +import { maxDuration, POST } from "./route"; + +const workspaceId = "e010bd2c-a3c1-438f-acef-166287a3b1cb"; +const workspace = { + id: workspaceId, + status: "ready", + sandboxId: "sandbox-1", +}; + +describe("gen2 instance route", () => { + beforeEach(() => { + mocks.ensure.mockResolvedValue(workspace); + mocks.detail.mockResolvedValue(workspace); + }); + + afterEach(() => vi.resetAllMocks()); + + it("allows host wake and Firecracker startup their bounded route budget", () => { + expect(maxDuration).toBe(300); + }); + + it("returns the ready workspace after startup", async () => { + const response = await POST( + new Request( + `https://codev.test/api/gen2/workspaces/${workspaceId}/instance`, + { + method: "POST", + }, + ), + { params: Promise.resolve({ workspaceId }) }, + ); + + expect(response.status).toBe(200); + await expect(response.json()).resolves.toEqual({ workspace }); + expect(mocks.ensure).toHaveBeenCalledWith(workspaceId, "user-1"); + expect(mocks.detail).toHaveBeenCalledWith(workspaceId, "user-1"); + }); + + it("returns 202 when another member still owns provisioning", async () => { + mocks.ensure.mockResolvedValueOnce({ + ...workspace, + status: "provisioning", + }); + mocks.detail.mockResolvedValueOnce({ + ...workspace, + status: "provisioning", + sandboxId: null, + }); + + const response = await POST( + new Request( + `https://codev.test/api/gen2/workspaces/${workspaceId}/instance`, + { + method: "POST", + }, + ), + { params: Promise.resolve({ workspaceId }) }, + ); + + expect(response.status).toBe(202); + await expect(response.json()).resolves.toEqual({ + workspace: { ...workspace, status: "provisioning", sandboxId: null }, + }); + }); + + it("uses the returned workspace state when a concurrent stop wins", async () => { + mocks.detail.mockResolvedValueOnce({ + ...workspace, + status: "stopped", + sandboxId: null, + }); + + const response = await POST( + new Request( + `https://codev.test/api/gen2/workspaces/${workspaceId}/instance`, + { + method: "POST", + }, + ), + { params: Promise.resolve({ workspaceId }) }, + ); + + expect(response.status).toBe(202); + await expect(response.json()).resolves.toEqual({ + workspace: { ...workspace, status: "stopped", sandboxId: null }, + }); + }); +}); diff --git a/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.ts b/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.ts index a7c4fa1a9..1da5e6c1e 100644 --- a/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.ts +++ b/apps/web/app/api/gen2/workspaces/[workspaceId]/instance/route.ts @@ -2,7 +2,7 @@ import { withUser } from "@/lib/http/api-route"; import { ensureGen2Instance, stopGen2Instance } from "@/lib/gen2/instance"; import { getGen2WorkspaceDetail } from "@/lib/gen2/workspaces"; -/** Host wake, Firecracker creation, and persistence fit Vercel's 300s limit. */ +/** One bounded host-wake attempt plus guest creation fits Vercel's 300s limit. */ export const maxDuration = 300; type Params = { workspaceId: string }; @@ -14,9 +14,13 @@ type Params = { workspaceId: string }; export const POST = withUser( async ({ user, params: { workspaceId } }) => { await ensureGen2Instance(workspaceId, user.id); - return Response.json({ - workspace: await getGen2WorkspaceDetail(workspaceId, user.id), - }); + const workspace = await getGen2WorkspaceDetail(workspaceId, user.id); + return Response.json( + { workspace }, + { + status: workspace.status === "ready" && workspace.sandboxId ? 200 : 202, + }, + ); }, { errorStatus: 502 }, ); diff --git a/apps/web/components/gen2/workspace-room.test.tsx b/apps/web/components/gen2/workspace-room.test.tsx index 1a8bf198c..09c09a847 100644 --- a/apps/web/components/gen2/workspace-room.test.tsx +++ b/apps/web/components/gen2/workspace-room.test.tsx @@ -47,7 +47,9 @@ function stubFetch(instance: () => Response) { const ready = () => new Response( - JSON.stringify({ workspace: { ...workspace, status: "ready" } }), + JSON.stringify({ + workspace: { ...workspace, status: "ready", sandboxId: "sandbox-1" }, + }), { status: 200 }, ); diff --git a/apps/web/components/gen2/workspace-room.tsx b/apps/web/components/gen2/workspace-room.tsx index 35bc8864a..5f7677914 100644 --- a/apps/web/components/gen2/workspace-room.tsx +++ b/apps/web/components/gen2/workspace-room.tsx @@ -6,6 +6,7 @@ import { ArrowLeft, Check, Link2 } from "lucide-react"; import type { Gen2WorkspaceDetail } from "@codev/contracts"; import { Gen2ChatPanel } from "./chat-panel"; +import { ensureGen2WorkspaceReady } from "@/lib/gen2/startup-client"; import { Gen2Workbench, type Gen2WorkbenchHandle } from "./workbench"; const STATUS_LABEL: Record = { @@ -29,6 +30,7 @@ export function Gen2WorkspaceRoom({ const [refreshToken, setRefreshToken] = useState(0); const [runtimeReady, setRuntimeReady] = useState(false); const workbenchRef = useRef(null); + const startupInFlightRef = useRef | null>(null); const ready = current.status === "ready" && runtimeReady; const displayStatus = current.status === "ready" && !runtimeReady @@ -47,35 +49,29 @@ export function Gen2WorkspaceRoom({ * decides whether anything needs doing, and a second member opening the * same workspace joins the boot already in progress. */ - const ensureRunning = useCallback(async () => { + const ensureRunning = useCallback(() => { + if (startupInFlightRef.current) return startupInFlightRef.current; + setRuntimeReady(false); - try { - const response = await fetch( - `/api/gen2/workspaces/${current.id}/instance`, - { method: "POST" }, - ); - const payload = (await response.json().catch(() => ({}))) as { - workspace?: Gen2WorkspaceDetail; - error?: string; - }; - if (payload.workspace) setCurrent(payload.workspace); - if (!response.ok) { - setCurrent((value) => ({ - ...value, - lastError: payload.error ?? "The machine could not start.", - })); - return false; + setCurrent((value) => ({ ...value, lastError: null })); + const attempt = (async () => { + const result = await ensureGen2WorkspaceReady(current.id); + if (result.workspace) { + setCurrent(result.workspace); + setRuntimeReady(true); + refresh(); + return true; } - setRuntimeReady(true); - refresh(); - return true; - } catch { - setCurrent((value) => ({ - ...value, - lastError: "The machine could not be reached. Try again.", - })); + setCurrent((value) => ({ ...value, lastError: result.error })); return false; - } + })(); + const tracked = attempt.finally(() => { + if (startupInFlightRef.current === tracked) { + startupInFlightRef.current = null; + } + }); + startupInFlightRef.current = tracked; + return tracked; }, [current.id, refresh]); const bootedRef = useRef(false); diff --git a/apps/web/lib/gen2/instance-lifecycle.test.ts b/apps/web/lib/gen2/instance-lifecycle.test.ts index 37a7a85e0..e14d2c53a 100644 --- a/apps/web/lib/gen2/instance-lifecycle.test.ts +++ b/apps/web/lib/gen2/instance-lifecycle.test.ts @@ -143,7 +143,7 @@ describe("gen2 instance lifecycle", () => { ]); expect(workspace.status).toBe("ready"); expect(workspace.sandboxId).toBe("sandbox-1"); - expect(mocks.ensureHostReady).toHaveBeenCalledWith(150_000); + expect(mocks.ensureHostReady).toHaveBeenCalledWith(60_000); }); it("reattaches when the Firecracker machine is already on the host", async () => { @@ -322,32 +322,18 @@ describe("gen2 instance lifecycle", () => { expect(mocks.provision).toHaveBeenCalledOnce(); }); - it("waits for another member's provisioning attempt to finish", async () => { - vi.useFakeTimers(); + it("returns the current workspace while another member is provisioning", async () => { mocks.member.status = "provisioning"; mocks.member.updatedAt = new Date(); mocks.state.claimed = false; - let reads = 0; - mocks.selectQuery.limit.mockImplementation(async () => { - reads += 1; - if (reads === 3) { - mocks.member.status = "ready"; - mocks.member.sandboxId = "sandbox-1"; - } - return [mocks.member]; - }); - const opening = ensureGen2Instance(mocks.member.id, "user-2", { + const workspace = await ensureGen2Instance(mocks.member.id, "user-2", { provision: mocks.provision, destroy: mocks.destroy, }); - await vi.advanceTimersByTimeAsync(1_000); - const workspace = await opening; - expect(mocks.provision).not.toHaveBeenCalled(); - expect(reads).toBeGreaterThanOrEqual(3); - expect(workspace.status).toBe("ready"); + expect(workspace.status).toBe("provisioning"); }); it("does not let an old provisioner overwrite a newer startup lease", async () => { @@ -394,6 +380,7 @@ describe("gen2 instance lifecycle", () => { expect(mocks.destroy).toHaveBeenCalledOnce(); expect(mocks.provision).toHaveBeenCalledOnce(); expect(mocks.updates.map((update) => update.status)).toEqual([ + "provisioning", "stopped", "provisioning", "ready", @@ -412,4 +399,35 @@ describe("gen2 instance lifecycle", () => { expect(workspace.status).toBe("stopped"); expect(workspace.sandboxId).toBeNull(); }); + + it("marks an unreachable host stopped without losing its snapshot", async () => { + mocks.member.status = "ready"; + mocks.member.sandboxId = "sandbox-1"; + mocks.destroy.mockRejectedValue(new TypeError("fetch failed")); + + const workspace = await stopGen2Instance(mocks.member.id, "user-1", { + provision: mocks.provision, + destroy: mocks.destroy, + }); + + expect(workspace.status).toBe("stopped"); + expect(workspace.sandboxId).toBeNull(); + expect(mocks.updates.map((update) => update.status)).toEqual([ + "provisioning", + "stopped", + ]); + }); + + it("does not stop a startup that another request is provisioning", async () => { + mocks.member.status = "provisioning"; + + await expect( + stopGen2Instance(mocks.member.id, "user-1", { + provision: mocks.provision, + destroy: mocks.destroy, + }), + ).rejects.toMatchObject({ status: 409 }); + + expect(mocks.destroy).not.toHaveBeenCalled(); + }); }); diff --git a/apps/web/lib/gen2/instance.test.ts b/apps/web/lib/gen2/instance.test.ts index 7fef10c49..ae4a80624 100644 --- a/apps/web/lib/gen2/instance.test.ts +++ b/apps/web/lib/gen2/instance.test.ts @@ -27,9 +27,9 @@ describe("gen2 Firecracker instance policy", () => { }); }); - it("only stops a live or starting instance", () => { + it("only stops a ready instance", () => { expect(canStopInstance("ready")).toBe(true); - expect(canStopInstance("provisioning")).toBe(true); + expect(canStopInstance("provisioning")).toBe(false); expect(canStopInstance("stopped")).toBe(false); expect(canStopInstance("pending")).toBe(false); }); diff --git a/apps/web/lib/gen2/instance.ts b/apps/web/lib/gen2/instance.ts index 9c15c368e..672601624 100644 --- a/apps/web/lib/gen2/instance.ts +++ b/apps/web/lib/gen2/instance.ts @@ -54,7 +54,7 @@ export function buildBlankSandboxSource() { } export function canStopInstance(status: Gen2WorkspaceStatus) { - return status === "ready" || status === "provisioning"; + return status === "ready"; } const HOST_UNREACHABLE_MESSAGE = @@ -62,17 +62,20 @@ const HOST_UNREACHABLE_MESSAGE = /** Firecracker create can outlast the default 70s orchestrator timeout. */ const GEN2_PROVISION_TIMEOUT_MS = 120_000; -/** Leave room for guest creation within Vercel's 300-second function limit. */ -const GEN2_HOST_READY_TIMEOUT_MS = 150_000; +/** + * A host wake is one bounded attempt. The workspace page retries this request + * while Azure finishes booting, keeping each Vercel invocation comfortably + * below its 300-second limit. + */ +const GEN2_HOST_READY_TIMEOUT_MS = 60_000; /** Keep expiry inside the orchestrator's exclusive four-hour window. */ const GEN2_EXPIRES_SLACK_MS = 60_000; /** - * The host wake and guest creation budgets total 270 seconds. Leave another - * thirty seconds for source preparation and persisting the final state. + * One host wake attempt and guest creation fit within 180 seconds. Leave + * another thirty seconds for source preparation and persisting the final state. */ const GEN2_PROVISIONING_STALE_AFTER_MS = GEN2_HOST_READY_TIMEOUT_MS + GEN2_PROVISION_TIMEOUT_MS + 30_000; -const GEN2_PROVISIONING_POLL_MS = 1_000; export function describeGen2RuntimeFailure(error: unknown): string { if (isGen2HostUnreachable(error)) { @@ -233,43 +236,6 @@ async function claimProvisioning(workspaceId: string) { return claim?.updatedAt ?? null; } -async function waitForProvisioning(workspaceId: string, userId: string) { - const deadline = Date.now() + GEN2_PROVISIONING_STALE_AFTER_MS; - while (Date.now() < deadline) { - const workspace = await requireGen2Member(workspaceId, userId); - if (workspace.status === "ready") return workspace; - if (workspace.status === "failed") { - throw new Gen2LifecycleError( - workspace.lastError ?? "The Firecracker instance could not start.", - 502, - ); - } - if (workspace.status !== "provisioning") { - throw new Gen2LifecycleError( - workspace.lastError ?? - "Workspace startup ended before the machine was ready.", - 503, - ); - } - if ( - Date.now() - Date.parse(workspace.updatedAt) >= - GEN2_PROVISIONING_STALE_AFTER_MS - ) { - throw new Gen2LifecycleError( - "The previous startup attempt stopped responding. Try again to resume this workspace.", - 503, - ); - } - await new Promise((resolve) => - setTimeout(resolve, GEN2_PROVISIONING_POLL_MS), - ); - } - throw new Gen2LifecycleError( - "The Firecracker instance is still starting. Try again in a moment.", - 503, - ); -} - /** * Makes sure this workspace has a machine, and returns once it does. * @@ -307,10 +273,9 @@ export async function ensureGen2Instance( const provisioningAt = await claimProvisioning(workspaceId); if (!provisioningAt) { - // Another member owns the startup lease. Join that operation instead of - // returning a successful response that leaves this browser stuck at - // "Starting" without an active runtime. - return waitForProvisioning(workspaceId, userId); + // Another member owns startup. Return its durable workspace state so this + // request stays short; the opening browser polls until it is ready. + return requireGen2Member(workspaceId, userId); } // A process can be terminated by its platform before its catch block runs. @@ -409,16 +374,67 @@ export async function stopGen2Instance( if (membership.role !== "owner") { throw new Gen2AccessError("Only the owner can stop this instance.", 403); } + if (membership.status === "provisioning") { + throw new Gen2LifecycleError( + "Wait for the instance to finish starting before stopping it.", + ); + } if (!canStopInstance(membership.status)) { throw new Gen2LifecycleError("This instance is not running."); } - await runtime.destroy(workspaceId); - await writeGen2Instance(workspaceId, { - status: "stopped", - sandboxId: null, - lastError: null, - }); + // Reserve the workspace before contacting Firecracker so a simultaneous + // open joins this stop instead of returning a guest that is being removed. + const stoppingAt = new Date(); + const [stopLease] = await getDatabase() + .update(schema.gen2Workspaces) + .set({ status: "provisioning", lastError: null, updatedAt: stoppingAt }) + .where( + and( + eq(schema.gen2Workspaces.id, workspaceId), + eq(schema.gen2Workspaces.status, "ready"), + ), + ) + .returning({ updatedAt: schema.gen2Workspaces.updatedAt }); + if (!stopLease) { + throw new Gen2LifecycleError( + "The instance changed state. Reload the workspace and try again.", + 409, + ); + } + + try { + await runtime.destroy(workspaceId); + } catch (error) { + if (!isGen2HostUnreachable(error)) { + const message = describeGen2RuntimeFailure(error); + await writeGen2Instance( + workspaceId, + { status: "ready", lastError: message }, + stoppingAt, + ); + throw new Gen2LifecycleError(message, 502); + } + + // A deallocated host cannot have a running guest to stop. Preserve the + // workspace snapshot; the next open wakes the host and restores it. + logEvent("warn", "gen2.instance.stop_host_unreachable", { + workspaceId, + detail: error instanceof Error ? error.message : "unknown", + }); + } + + const committed = await writeGen2Instance( + workspaceId, + { status: "stopped", sandboxId: null, lastError: null }, + stoppingAt, + ); + if (!committed) { + throw new Gen2LifecycleError( + "A newer lifecycle operation took over. Reload the workspace and try again.", + 503, + ); + } return requireGen2Member(workspaceId, userId); } diff --git a/apps/web/lib/gen2/startup-client.test.ts b/apps/web/lib/gen2/startup-client.test.ts new file mode 100644 index 000000000..fa29d930c --- /dev/null +++ b/apps/web/lib/gen2/startup-client.test.ts @@ -0,0 +1,140 @@ +// @vitest-environment node + +import { describe, expect, it, vi } from "vitest"; +import type { Gen2WorkspaceDetail } from "@codev/contracts"; + +import { ensureGen2WorkspaceReady } from "./startup-client"; + +function workspace( + status: Gen2WorkspaceDetail["status"], + sandboxId: string | null = null, +) { + return { + id: "workspace-1", + name: "Studio", + status, + sandboxId, + lastError: null, + role: "owner", + repository: null, + repositoryPrivate: false, + defaultBranch: null, + createdAt: "2026-09-20T20:00:00.000Z", + updatedAt: "2026-09-20T20:00:00.000Z", + members: [], + } as Gen2WorkspaceDetail; +} + +function response(status: number, body: unknown) { + return new Response(JSON.stringify(body), { + status, + headers: { "content-type": "application/json" }, + }); +} + +function fetchSequence(responses: Response[]) { + const calls: Array<{ url: string; method: string }> = []; + const fetcher = (async (input: RequestInfo | URL, init?: RequestInit) => { + calls.push({ + url: String(input), + method: init?.method ?? "GET", + }); + const next = responses.shift(); + if (!next) throw new Error("Unexpected request"); + return next; + }) as typeof fetch; + return { calls, fetcher }; +} + +describe("Gen 2 workspace startup polling", () => { + it("retries bounded host-wake failures until the server returns ready", async () => { + const { calls, fetcher } = fetchSequence([ + response(503, { error: "The Firecracker host is still starting." }), + response(200, { workspace: workspace("ready", "sandbox-1") }), + ]); + const pause = vi.fn(async () => {}); + + const result = await ensureGen2WorkspaceReady("workspace-1", { + fetcher, + pause, + random: () => 0.5, + }); + + expect(result.workspace?.sandboxId).toBe("sandbox-1"); + expect(calls).toEqual([ + { url: "/api/gen2/workspaces/workspace-1/instance", method: "POST" }, + { url: "/api/gen2/workspaces/workspace-1/instance", method: "POST" }, + ]); + expect(pause).toHaveBeenCalledOnce(); + }); + + it("joins another member's startup by polling persisted workspace state", async () => { + const { calls, fetcher } = fetchSequence([ + response(202, { workspace: workspace("provisioning") }), + response(200, { workspace: workspace("provisioning") }), + response(200, { workspace: workspace("ready", "sandbox-1") }), + ]); + const pause = vi.fn(async () => {}); + + const result = await ensureGen2WorkspaceReady("workspace-1", { + fetcher, + pause, + random: () => 0.5, + }); + + expect(result.workspace?.status).toBe("ready"); + expect(calls.map((call) => call.method)).toEqual(["POST", "GET", "GET"]); + expect(pause).toHaveBeenCalledTimes(2); + }); + + it("restarts startup when the previous host-wake attempt returns to pending", async () => { + const { calls, fetcher } = fetchSequence([ + response(202, { workspace: workspace("provisioning") }), + response(200, { workspace: workspace("pending") }), + response(200, { workspace: workspace("ready", "sandbox-1") }), + ]); + + const result = await ensureGen2WorkspaceReady("workspace-1", { + fetcher, + pause: async () => {}, + random: () => 0.5, + }); + + expect(result.workspace?.status).toBe("ready"); + expect(calls.map((call) => call.method)).toEqual(["POST", "GET", "POST"]); + }); + + it("surfaces non-retryable startup failures for the existing retry action", async () => { + const { fetcher } = fetchSequence([ + response(502, { error: "Firecracker guest creation failed." }), + ]); + + await expect( + ensureGen2WorkspaceReady("workspace-1", { fetcher }), + ).resolves.toEqual({ error: "Firecracker guest creation failed." }); + }); + + it("stops retrying after the bounded client wait", async () => { + let clock = 0; + const { calls, fetcher } = fetchSequence([ + response(503, { error: "The Firecracker host is still starting." }), + response(503, { error: "The Firecracker host is still starting." }), + response(503, { error: "The Firecracker host is still starting." }), + ]); + + const result = await ensureGen2WorkspaceReady("workspace-1", { + fetcher, + maxWaitMs: 5_000, + now: () => clock, + pause: async (milliseconds) => { + clock += milliseconds; + }, + random: () => 0.5, + }); + + expect(result).toEqual({ + error: "The Firecracker host is still starting. Try again in a moment.", + }); + expect(calls).toHaveLength(2); + }); +}); diff --git a/apps/web/lib/gen2/startup-client.ts b/apps/web/lib/gen2/startup-client.ts new file mode 100644 index 000000000..75aee80d4 --- /dev/null +++ b/apps/web/lib/gen2/startup-client.ts @@ -0,0 +1,160 @@ +import type { Gen2WorkspaceDetail } from "@codev/contracts"; + +const STARTUP_MAX_WAIT_MS = 12 * 60_000; +const STARTUP_RECHECK_MS = 60_000; +const MAX_BACKOFF_MS = 15_000; + +type WorkspaceResponse = { + workspace?: Gen2WorkspaceDetail; + error?: string; +}; + +type ReadyGen2WorkspaceDetail = Gen2WorkspaceDetail & { + status: "ready"; + sandboxId: string; +}; + +type StartupResult = + | { workspace: Gen2WorkspaceDetail; error?: never } + | { workspace?: never; error: string }; + +type StartupDependencies = { + fetcher?: typeof fetch; + pause?: (milliseconds: number) => Promise; + now?: () => number; + random?: () => number; + maxWaitMs?: number; + recheckMs?: number; +}; + +const pause = (milliseconds: number) => + new Promise((resolve) => setTimeout(resolve, milliseconds)); + +function retryDelay(attempt: number, random: () => number) { + const base = Math.min(MAX_BACKOFF_MS, 1_000 * 2 ** Math.min(attempt, 4)); + return Math.round(base * (0.8 + random() * 0.4)); +} + +async function readResponse(response: Response): Promise { + return (await response.json().catch(() => ({}))) as WorkspaceResponse; +} + +function readyWorkspace( + workspace: Gen2WorkspaceDetail | undefined, +): workspace is ReadyGen2WorkspaceDetail { + return workspace?.status === "ready" && Boolean(workspace.sandboxId); +} + +/** + * Keep the workspace opening while Azure wakes its host. Each POST is bounded + * by the server; a 503 means the host is still waking and can be retried. When + * another member owns provisioning, poll its persisted workspace state and + * periodically POST again so a stale startup claim can be recovered. + */ +export async function ensureGen2WorkspaceReady( + workspaceId: string, + dependencies: StartupDependencies = {}, +): Promise { + const fetcher = dependencies.fetcher ?? fetch; + const wait = dependencies.pause ?? pause; + const now = dependencies.now ?? Date.now; + const random = dependencies.random ?? Math.random; + const deadline = now() + (dependencies.maxWaitMs ?? STARTUP_MAX_WAIT_MS); + const recheckMs = dependencies.recheckMs ?? STARTUP_RECHECK_MS; + const url = `/api/gen2/workspaces/${workspaceId}`; + let attempt = 0; + let lastPostAt = Number.NEGATIVE_INFINITY; + let postNext = true; + + while (now() < deadline) { + if (postNext) { + lastPostAt = now(); + postNext = false; + try { + const response = await fetcher(`${url}/instance`, { method: "POST" }); + const payload = await readResponse(response); + if (readyWorkspace(payload.workspace)) { + return { workspace: payload.workspace! }; + } + if (payload.workspace) { + if (payload.workspace.status === "failed") { + return { + error: + payload.workspace.lastError ?? + "The Firecracker instance could not start.", + }; + } + if (payload.workspace.status === "deleting") { + return { error: "This workspace is being deleted." }; + } + } + if (!response.ok && response.status !== 503) { + return { + error: payload.error ?? "The machine could not start.", + }; + } + if (response.status === 503) { + attempt += 1; + postNext = true; + } + } catch { + // A response can be lost after the server claims startup. The next + // status read can join it without holding another long request open. + attempt += 1; + } + if (postNext) { + if (now() < deadline) { + await wait(retryDelay(attempt, random)); + } + continue; + } + } else { + await wait(retryDelay(attempt, random)); + if (now() >= deadline) break; + if (now() - lastPostAt >= recheckMs) { + postNext = true; + continue; + } + + try { + const response = await fetcher(url, { method: "GET" }); + const payload = await readResponse(response); + if (!response.ok) { + if (response.status === 404 || response.status === 409) { + return { + error: payload.error ?? "This workspace is no longer available.", + }; + } + attempt += 1; + continue; + } + const workspace = payload.workspace; + if (workspace && readyWorkspace(workspace)) return { workspace }; + if (!workspace) { + return { error: "The workspace status could not be loaded." }; + } + if (workspace.status === "failed") { + return { + error: + workspace.lastError ?? + "The Firecracker instance could not start.", + }; + } + if (workspace.status === "deleting") { + return { error: "This workspace is being deleted." }; + } + if (workspace.status === "pending" || workspace.status === "stopped") { + postNext = true; + } else { + attempt += 1; + } + } catch { + attempt += 1; + } + } + } + + return { + error: "The Firecracker host is still starting. Try again in a moment.", + }; +} diff --git a/apps/web/lib/gen2/workspaces-delete.test.ts b/apps/web/lib/gen2/workspaces-delete.test.ts index ebb680f49..1acc5a203 100644 --- a/apps/web/lib/gen2/workspaces-delete.test.ts +++ b/apps/web/lib/gen2/workspaces-delete.test.ts @@ -96,25 +96,18 @@ describe("Gen 2 workspace deletion", () => { mocks.databaseUpdateQuery.where.mockClear(); }); - it("wakes an unreachable host to purge workspace data before deleting", async () => { - mocks.destroySandbox - .mockRejectedValueOnce(new TypeError("fetch failed")) - .mockResolvedValueOnce(undefined); - + it("readies the host before purging workspace data", async () => { await expect( deleteGen2Workspace("11111111-1111-4111-8111-111111111111", "user-1"), ).resolves.toBeUndefined(); expect(mocks.ensureHostReady).toHaveBeenCalledOnce(); - expect(mocks.destroySandbox).toHaveBeenCalledTimes(2); + expect(mocks.ensureHostReady).toHaveBeenCalledWith(120_000); + expect(mocks.destroySandbox).toHaveBeenCalledOnce(); expect(mocks.discardSandboxSnapshot).toHaveBeenCalledOnce(); expect(mocks.database.delete).toHaveBeenCalledOnce(); - expect(mocks.logEvent).toHaveBeenCalledWith( - "warn", - "gen2.workspace.delete_host_unreachable", - expect.objectContaining({ - workspaceId: "11111111-1111-4111-8111-111111111111", - }), + expect(mocks.ensureHostReady.mock.invocationCallOrder[0]).toBeLessThan( + mocks.destroySandbox.mock.invocationCallOrder[0]!, ); }); @@ -133,4 +126,20 @@ describe("Gen 2 workspace deletion", () => { expect(mocks.database.delete).not.toHaveBeenCalled(); expect(mocks.database.update).toHaveBeenCalledOnce(); }); + + it("keeps deletion retryable when the host does not wake in time", async () => { + mocks.ensureHostReady.mockRejectedValue( + new Error("The Firecracker host is still starting."), + ); + + await expect( + deleteGen2Workspace("11111111-1111-4111-8111-111111111111", "user-1"), + ).rejects.toMatchObject({ status: 502 }); + + expect(mocks.ensureHostReady).toHaveBeenCalledWith(120_000); + expect(mocks.destroySandbox).not.toHaveBeenCalled(); + expect(mocks.discardSandboxSnapshot).not.toHaveBeenCalled(); + expect(mocks.database.delete).not.toHaveBeenCalled(); + expect(mocks.database.update).toHaveBeenCalledOnce(); + }); }); diff --git a/apps/web/lib/gen2/workspaces.test.ts b/apps/web/lib/gen2/workspaces.test.ts index cc5b525ff..e7024deef 100644 --- a/apps/web/lib/gen2/workspaces.test.ts +++ b/apps/web/lib/gen2/workspaces.test.ts @@ -226,6 +226,7 @@ describe("gen2 workspaces", () => { expect(mocks.events).toEqual([ "lock-delete-row", "update:gen2_workspaces", + "wake-host", "destroy", "discard-snapshot", "delete-row", @@ -242,14 +243,9 @@ describe("gen2 workspaces", () => { expect(mocks.ensureHostReady).not.toHaveBeenCalled(); }); - it("wakes an unreachable host and purges the snapshot before deletion", async () => { + it("wakes a stopped host before purging the snapshot", async () => { mocks.currentStatus = "stopped"; mocks.memberStatus = "stopped"; - mocks.destroySandbox - .mockRejectedValueOnce(new TypeError("fetch failed")) - .mockImplementationOnce(async () => { - mocks.events.push("destroy"); - }); await deleteGen2Workspace("11111111-1111-4111-8111-111111111111", "user-1"); diff --git a/apps/web/lib/gen2/workspaces.ts b/apps/web/lib/gen2/workspaces.ts index 18e702a02..4c5e46419 100644 --- a/apps/web/lib/gen2/workspaces.ts +++ b/apps/web/lib/gen2/workspaces.ts @@ -20,11 +20,7 @@ import { discardSandboxSnapshot, } from "../runtime/orchestrator-sandbox"; import { GEN2_MAX_OWNED_WORKSPACES } from "./constants"; -import { - Gen2AccessError, - Gen2LifecycleError, - isGen2HostUnreachable, -} from "./errors"; +import { Gen2AccessError, Gen2LifecycleError } from "./errors"; const DEFAULT_WORKSPACE_NAME = "Workspace"; const INVITE_TTL_MS = 7 * 24 * 60 * 60 * 1_000; @@ -255,22 +251,15 @@ export async function deleteGen2Workspace(workspaceId: string, userId: string) { try { // A never-started workspace has no guest or snapshot, so don't wake a host - // just to delete its database row. For all other states, purge the runtime - // data before freeing the owner's workspace slot. + // just to delete its database row. For all other states, make sure the + // host is responsive before teardown so an unreachable-host timeout does + // not consume most of Vercel's function budget before the wake attempt. if (currentStatus !== "pending") { - try { - await destroySandbox(workspaceId); - } catch (error) { - if (!isGen2HostUnreachable(error)) throw error; - logEvent("warn", "gen2.workspace.delete_host_unreachable", { - workspaceId, - detail: error instanceof Error ? error.message : "unknown", - }); - // The host may have deallocated while retaining the workspace disk. - // Wake it so deletion can remove the snapshot rather than orphaning it. - await ensureHostReady(); - await destroySandbox(workspaceId); - } + // Leave room in Vercel's 300-second request budget for guest teardown, + // snapshot removal, and the final database delete. If the Azure VM is + // still booting at this bound, the deleting record remains retryable. + await ensureHostReady(120_000); + await destroySandbox(workspaceId); await discardSandboxSnapshot(workspaceId); } await database diff --git a/apps/web/lib/runtime/orchestrator-health.test.ts b/apps/web/lib/runtime/orchestrator-health.test.ts new file mode 100644 index 000000000..6af1ac072 --- /dev/null +++ b/apps/web/lib/runtime/orchestrator-health.test.ts @@ -0,0 +1,49 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + fakeGuestEnabled: vi.fn(() => false), + requestHostWake: vi.fn(), + orchestratorRequest: vi.fn(), +})); + +vi.mock("./fake-guest", () => ({ + fakeGuestEnabled: mocks.fakeGuestEnabled, +})); + +vi.mock("./host", () => ({ + requestHostWake: mocks.requestHostWake, +})); + +vi.mock("./orchestrator-request", () => ({ + OrchestratorError: class OrchestratorError extends Error { + constructor( + message: string, + readonly status: number, + ) { + super(message); + } + }, + orchestratorRequest: mocks.orchestratorRequest, + orchestratorRequestAt: vi.fn(), +})); + +import { ensureHostReady } from "./orchestrator-health"; + +describe("ensureHostReady", () => { + beforeEach(() => { + mocks.fakeGuestEnabled.mockReturnValue(false); + mocks.requestHostWake.mockReset().mockResolvedValue("running"); + mocks.orchestratorRequest + .mockReset() + .mockResolvedValue( + Response.json({ status: "ok", service: "codev-orchestrator" }), + ); + }); + + it("leaves stopping-host waits to its bounded health poll loop", async () => { + await expect(ensureHostReady(1_000)).resolves.toBeUndefined(); + + expect(mocks.requestHostWake).toHaveBeenCalledOnce(); + expect(mocks.requestHostWake).toHaveBeenCalledWith(1); + }); +}); diff --git a/apps/web/lib/runtime/orchestrator-health.ts b/apps/web/lib/runtime/orchestrator-health.ts index 48c503f1c..818b25391 100644 --- a/apps/web/lib/runtime/orchestrator-health.ts +++ b/apps/web/lib/runtime/orchestrator-health.ts @@ -65,7 +65,9 @@ export async function ensureHostReady(timeoutMs = HOST_START_TIMEOUT_MS) { if (fakeGuestEnabled()) return; const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { - const state = await requestHostWake().catch(() => "starting" as const); + // This function owns its own poll loop; keep an individual Azure status + // check from spending the whole loop budget waiting out a stopping VM. + const state = await requestHostWake(1).catch(() => "starting" as const); if (state === "running") { try { await waitForOrchestrator(Math.min(45_000, deadline - Date.now())); From eb3f0fd801c48ddfc051a229c298c0e2c0e66f87 Mon Sep 17 00:00:00 2001 From: Yousef Abdelhadi Date: Sun, 27 Sep 2026 03:02:40 -0400 Subject: [PATCH 3/4] fix(gen2): preserve terminal output and input ordering --- .../components/gen2/terminal-pane.test.tsx | 77 ++++++++++++++- apps/web/components/gen2/terminal-pane.tsx | 14 ++- .../lib/gen2/workbench.integration.test.ts | 19 ++++ apps/web/lib/runtime/fake-guest.ts | 4 +- services/orchestrator/src/guest.rs | 97 ++++++++++++++++++- services/orchestrator/src/model.rs | 4 + 6 files changed, 207 insertions(+), 8 deletions(-) diff --git a/apps/web/components/gen2/terminal-pane.test.tsx b/apps/web/components/gen2/terminal-pane.test.tsx index 18e4b06e4..b890a6294 100644 --- a/apps/web/components/gen2/terminal-pane.test.tsx +++ b/apps/web/components/gen2/terminal-pane.test.tsx @@ -1,11 +1,24 @@ import { fireEvent, render, screen, waitFor } from "@testing-library/react"; import { beforeEach, describe, expect, it, vi } from "vitest"; +const terminalMocks = vi.hoisted(() => ({ + instances: [] as Array<{ onDataHandler?: (data: string) => void }>, +})); + vi.mock("@xterm/xterm", () => ({ Terminal: class { rows = 24; cols = 80; - onData = vi.fn(); + onDataHandler?: (data: string) => void; + + constructor() { + terminalMocks.instances.push(this); + } + + onData(handler: (data: string) => void) { + this.onDataHandler = handler; + } + loadAddon() {} open() {} dispose() {} @@ -25,6 +38,7 @@ const workspaceId = "11111111-1111-4111-8111-111111111111"; describe("Gen2TerminalPane", () => { beforeEach(() => { + terminalMocks.instances.length = 0; class ResizeObserverStub { observe() {} disconnect() {} @@ -71,4 +85,65 @@ describe("Gen2TerminalPane", () => { fireEvent.click(resume); await waitFor(() => expect(onResumeWorkspace).toHaveBeenCalledOnce()); }); + + it("serializes rapid terminal input so characters arrive in order", async () => { + let releaseFirstInput!: () => void; + const firstInput = new Promise((resolve) => { + releaseFirstInput = resolve; + }); + const inputCalls: string[] = []; + + vi.stubGlobal( + "fetch", + vi.fn(async (_input: RequestInfo | URL, init?: RequestInit) => { + const request = JSON.parse(String(init?.body ?? "{}")) as { + action?: string; + data?: string; + }; + if (request.action === "start") { + return new Response(JSON.stringify({ sessionId: "terminal-1" }), { + status: 201, + }); + } + if (request.action === "input") { + inputCalls.push(request.data ?? ""); + if (inputCalls.length === 1) await firstInput; + return new Response(null, { status: 204 }); + } + if (request.action === "poll") { + return new Promise(() => {}); + } + return new Response(null, { status: 204 }); + }), + ); + + const { unmount } = render( + true)} + />, + ); + + fireEvent.click(screen.getByRole("button", { name: "Start terminal" })); + await waitFor(() => + expect(terminalMocks.instances[0]?.onDataHandler).toBeDefined(), + ); + const sendInput = terminalMocks.instances[0]!.onDataHandler!; + sendInput("p"); + sendInput("w"); + sendInput("d"); + + try { + await waitFor(() => expect(inputCalls).toEqual(["p"])); + expect(inputCalls).toEqual(["p"]); + } finally { + releaseFirstInput(); + } + + await waitFor(() => expect(inputCalls).toEqual(["p", "w", "d"])); + unmount(); + }); }); diff --git a/apps/web/components/gen2/terminal-pane.tsx b/apps/web/components/gen2/terminal-pane.tsx index 48613c93a..1a5ec740c 100644 --- a/apps/web/components/gen2/terminal-pane.tsx +++ b/apps/web/components/gen2/terminal-pane.tsx @@ -100,12 +100,20 @@ export function Gen2TerminalPane({ setStatus("idle"); return; } - sessionRef.current = payload.sessionId; + const sessionId = payload.sessionId; + sessionRef.current = sessionId; afterRef.current = 0; setStatus("live"); + let inputQueue = Promise.resolve(); term.onData((data) => { - void post({ action: "input", sessionId: payload.sessionId, data }) - .then((inputResponse) => { + inputQueue = inputQueue + .then(async () => { + if (sessionRef.current !== sessionId) return; + const inputResponse = await post({ + action: "input", + sessionId, + data, + }); if ([404, 502, 503].includes(inputResponse.status)) { markWorkspacePaused(); } diff --git a/apps/web/lib/gen2/workbench.integration.test.ts b/apps/web/lib/gen2/workbench.integration.test.ts index c9f588888..f938e8344 100644 --- a/apps/web/lib/gen2/workbench.integration.test.ts +++ b/apps/web/lib/gen2/workbench.integration.test.ts @@ -211,6 +211,25 @@ describe("gen2 workbench against the guest", () => { expect(second.chunks).toEqual([]); }); + it("keeps terminal output at the next-sequence cursor", async () => { + const sessionId = await startGen2Terminal(workspaceId, userId, { + rows: 24, + columns: 80, + }); + const first = await pollGen2Terminal(workspaceId, userId, sessionId, 0); + + await sendGen2TerminalInput(workspaceId, userId, sessionId, "x"); + const next = await pollGen2Terminal( + workspaceId, + userId, + sessionId, + first.nextSequence, + ); + + expect(next.chunks.map((chunk) => chunk.data).join("")).toContain("x"); + await closeGen2Terminal(workspaceId, userId, sessionId); + }); + it("runs a Codex turn end to end and reduces it into activity", async () => { // The real command builder, the real exec transport, the real decoder and // the real reducer -- only the guest is doubled. diff --git a/apps/web/lib/runtime/fake-guest.ts b/apps/web/lib/runtime/fake-guest.ts index 38060fd90..3878a768d 100644 --- a/apps/web/lib/runtime/fake-guest.ts +++ b/apps/web/lib/runtime/fake-guest.ts @@ -390,7 +390,9 @@ export function handleFakeGuestRequest( if (action === "/resize") return new Response(null, { status: 204 }); if (action === "/poll") { const after = Number(input.after ?? 0); - const chunks = session.chunks.filter((chunk) => chunk.sequence > after); + const chunks = session.chunks.filter( + (chunk) => chunk.sequence >= after, + ); return json({ result: { chunks, diff --git a/services/orchestrator/src/guest.rs b/services/orchestrator/src/guest.rs index 212ae8995..9831d3772 100644 --- a/services/orchestrator/src/guest.rs +++ b/services/orchestrator/src/guest.rs @@ -922,7 +922,7 @@ impl GuestService { while output .chunks .front() - .is_some_and(|chunk| chunk.sequence <= request.after) + .is_some_and(|chunk| chunk.sequence < request.after) { if let Some(chunk) = output.chunks.pop_front() { output.buffered_bytes = output.buffered_bytes.saturating_sub(chunk.data.len()); @@ -2001,7 +2001,7 @@ impl GuestService { while output .chunks .front() - .is_some_and(|chunk| chunk.sequence <= request.after) + .is_some_and(|chunk| chunk.sequence < request.after) { if let Some(chunk) = output.chunks.pop_front() { output.buffered_bytes = output.buffered_bytes.saturating_sub(chunk.data.len()); @@ -3176,7 +3176,7 @@ mod tests { assert!(TerminalUser::parse(passwd, "codev-shell").is_none()); } - use std::process::Command; + use std::{process::Command, sync::atomic::AtomicBool}; use tempfile::tempdir; @@ -3868,6 +3868,97 @@ mod tests { assert_eq!(close.status, 200); } + #[test] + fn terminal_poll_keeps_a_chunk_at_the_next_sequence_cursor() { + let directory = tempdir().expect("tempdir"); + let service = GuestService::new(directory.path()).expect("service"); + let start = service.handle("POST", "/v1/terminals", br#"{"rows":24,"columns":80}"#); + assert_eq!(start.status, 200); + let start_body: serde_json::Value = + serde_json::from_slice(&start.body).expect("terminal start"); + let session_id = start_body["sessionId"].as_str().expect("session id"); + let session = service.terminal(session_id).expect("terminal session"); + let cursor = { + let mut output = session.output.lock().expect("terminal output lock"); + let cursor = output.next_sequence; + output.chunks.push_back(TerminalChunk { + sequence: cursor, + data: "boundary-output".into(), + }); + output.next_sequence += 1; + cursor + }; + + let poll = service.handle( + "POST", + &format!("/v1/terminals/{session_id}/poll"), + serde_json::to_vec(&serde_json::json!({ + "after": cursor, + "waitMilliseconds": 0, + })) + .expect("poll request") + .as_slice(), + ); + assert_eq!(poll.status, 200); + let result: TerminalPollResponse = + serde_json::from_slice(&poll.body).expect("terminal poll"); + assert!( + result.chunks.iter().any(|chunk| { + chunk.sequence == cursor && chunk.data.contains("boundary-output") + }) + ); + + let close = service.handle("DELETE", &format!("/v1/terminals/{session_id}"), b""); + assert_eq!(close.status, 200); + } + + #[test] + fn codex_poll_keeps_a_chunk_at_the_next_sequence_cursor() { + let directory = tempdir().expect("tempdir"); + let service = GuestService::new(directory.path()).expect("service"); + let session_id = "codex-cursor-test"; + let cursor = 1; + let bytes = b"boundary-output".to_vec(); + service.codex_execs.lock().expect("codex exec lock").insert( + session_id.into(), + Arc::new(CodexExecSession { + output: Mutex::new(CodexExecOutput { + chunks: VecDeque::from([CodexExecRawChunk { + sequence: cursor, + data: bytes.clone(), + }]), + buffered_bytes: bytes.len(), + next_sequence: cursor + 1, + reader_closed: false, + exit_code: None, + codex_auth_cache_json: None, + }), + output_changed: Condvar::new(), + cancel_requested: AtomicBool::new(false), + }), + ); + + let poll = service.handle( + "POST", + &format!("/v1/codex-execs/{session_id}/poll"), + serde_json::to_vec(&serde_json::json!({ + "after": cursor, + "waitMilliseconds": 0, + })) + .expect("poll request") + .as_slice(), + ); + assert_eq!(poll.status, 200); + let result: CodexExecPollResponse = + serde_json::from_slice(&poll.body).expect("Codex exec poll"); + assert!(result.chunks.iter().any(|chunk| { + chunk.sequence == cursor + && BASE64 + .decode(&chunk.data_base64) + .is_ok_and(|data| data == bytes) + })); + } + #[test] fn terminal_start_reclaims_oldest_when_at_capacity() { let directory = tempdir().expect("tempdir"); diff --git a/services/orchestrator/src/model.rs b/services/orchestrator/src/model.rs index 2c8ca88e7..4ff8d60b2 100644 --- a/services/orchestrator/src/model.rs +++ b/services/orchestrator/src/model.rs @@ -344,6 +344,7 @@ pub struct TerminalResizeRequest { #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct TerminalPollRequest { + /// Exclusive cursor: the response includes chunks at this sequence and later. #[serde(default)] pub after: u64, #[serde(default)] @@ -361,6 +362,7 @@ pub struct TerminalChunk { #[serde(rename_all = "camelCase")] pub struct TerminalPollResponse { pub chunks: Vec, + /// Pass as `after` in the next poll; a chunk at this value may arrive later. pub next_sequence: u64, pub exited: bool, pub exit_code: Option, @@ -389,6 +391,7 @@ pub struct CodexExecStartRequest { #[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct CodexExecPollRequest { + /// Exclusive cursor: the response includes chunks at this sequence and later. #[serde(default)] pub after: u64, #[serde(default)] @@ -410,6 +413,7 @@ pub struct CodexExecChunk { #[serde(rename_all = "camelCase")] pub struct CodexExecPollResponse { pub chunks: Vec, + /// Pass as `after` in the next poll; a chunk at this value may arrive later. pub next_sequence: u64, pub exited: bool, pub exit_code: Option, From eca956f78336e6b267875e3d1f79d8f0839a0ae7 Mon Sep 17 00:00:00 2001 From: Yousef Abdelhadi Date: Sun, 27 Sep 2026 03:02:52 -0400 Subject: [PATCH 4/4] perf(runtime): reuse private Claude VMs briefly --- .../claude-runtime-execution.test.ts | 17 +++ .../lib/providers/claude-runtime-execution.ts | 15 ++- apps/web/lib/runtime/orchestrator-sandbox.ts | 11 ++ apps/web/lib/runtime/orchestrator.ts | 1 + docs/SECURITY.md | 5 + .../orchestrator/src/backend/firecracker.rs | 104 +++++++++++++++- services/orchestrator/src/backend/mod.rs | 117 +++++++++++++++++- services/orchestrator/src/http_api.rs | 10 ++ 8 files changed, 271 insertions(+), 9 deletions(-) diff --git a/apps/web/lib/providers/claude-runtime-execution.test.ts b/apps/web/lib/providers/claude-runtime-execution.test.ts index 53f20a2ba..9c822a86a 100644 --- a/apps/web/lib/providers/claude-runtime-execution.test.ts +++ b/apps/web/lib/providers/claude-runtime-execution.test.ts @@ -8,6 +8,7 @@ const mocks = vi.hoisted(() => ({ poll: vi.fn(), close: vi.fn(), snapshot: vi.fn(), + park: vi.fn(), destroy: vi.fn(), })); vi.mock("./claude-connection-session", () => ({ @@ -25,8 +26,10 @@ vi.mock("../runtime/orchestrator", () => ({ pollCodexExecInSandbox: mocks.poll, closeCodexExecInSandbox: mocks.close, snapshotWorkspace: mocks.snapshot, + parkSandbox: mocks.park, destroySandbox: mocks.destroy, })); +import { OrchestratorError } from "../runtime/orchestrator-request"; import { encodeClaudeRuntimeReference } from "./claude-runtime-reference"; import { startClaudeExecution, @@ -101,6 +104,20 @@ describe("private Claude execution", () => { ); await cleanupClaudeExecution(id); expect(mocks.snapshot).toHaveBeenCalledWith(profileId, "a".repeat(40)); + expect(mocks.park).toHaveBeenCalledWith(profileId); + expect(mocks.destroy).not.toHaveBeenCalled(); + expect(mocks.release).toHaveBeenCalledWith(connectionId, leaseUntil); + }); + it("uses the checkpoint and destroys the VM while the old host lacks parking", async () => { + const id = await startClaudeExecution( + "sender", + "sonnet", + "Hello", + "request", + ); + mocks.park.mockRejectedValueOnce(new OrchestratorError("Not found", 404)); + await cleanupClaudeExecution(id); + expect(mocks.snapshot).toHaveBeenCalledOnce(); expect(mocks.destroy).toHaveBeenCalledWith(profileId); expect(mocks.release).toHaveBeenCalledWith(connectionId, leaseUntil); }); diff --git a/apps/web/lib/providers/claude-runtime-execution.ts b/apps/web/lib/providers/claude-runtime-execution.ts index 9e905f557..4fd2e3344 100644 --- a/apps/web/lib/providers/claude-runtime-execution.ts +++ b/apps/web/lib/providers/claude-runtime-execution.ts @@ -25,8 +25,10 @@ import { pollCodexExecInSandbox, closeCodexExecInSandbox, snapshotWorkspace, + parkSandbox, destroySandbox, } from "../runtime/orchestrator"; +import { OrchestratorError } from "../runtime/orchestrator-request"; // Official CLI aliases resolve on the signed-in account; do not claim specific // dated API models are available to every subscription. @@ -324,8 +326,19 @@ export async function cleanupClaudeExecution(id: string) { if (connection?.id === run.connectionId) { const sandbox = await getSandbox(run.reference.profileId); await snapshotWorkspace(run.reference.profileId, sandbox.headSha); + try { + await parkSandbox(run.reference.profileId); + } catch (error) { + // During a staggered rollout, an older host lacks the park route. + // The checkpoint is already durable, so use the prior cleanup path. + if (!(error instanceof OrchestratorError && error.status === 404)) { + throw error; + } + await destroySandbox(run.reference.profileId); + } + } else { + await destroySandbox(run.reference.profileId); } - await destroySandbox(run.reference.profileId); } await releaseClaudeSubscriptionExecution(run.connectionId, run.leaseUntil); } diff --git a/apps/web/lib/runtime/orchestrator-sandbox.ts b/apps/web/lib/runtime/orchestrator-sandbox.ts index 8a59ccfba..d64ca86d7 100644 --- a/apps/web/lib/runtime/orchestrator-sandbox.ts +++ b/apps/web/lib/runtime/orchestrator-sandbox.ts @@ -94,3 +94,14 @@ export async function touchSandbox(workspaceId: string) { .object({ sandbox: sandboxInstanceSchema }) .parse(await response.json()).sandbox; } + +/** Park a checkpointed ephemeral VM for a short, reclaimable idle window. */ +export async function parkSandbox(workspaceId: string) { + const response = await orchestratorRequest( + "POST", + `/v1/sandboxes/${workspaceId}/park`, + ); + return z + .object({ sandbox: sandboxInstanceSchema }) + .parse(await response.json()).sandbox; +} diff --git a/apps/web/lib/runtime/orchestrator.ts b/apps/web/lib/runtime/orchestrator.ts index a6ee6fc0a..02db3f0f4 100644 --- a/apps/web/lib/runtime/orchestrator.ts +++ b/apps/web/lib/runtime/orchestrator.ts @@ -61,6 +61,7 @@ export { resumeSandbox, discardSandboxSnapshot, touchSandbox, + parkSandbox, } from "./orchestrator-sandbox"; export { startSandboxTerminal, diff --git a/docs/SECURITY.md b/docs/SECURITY.md index d959e87f5..87b490c55 100644 --- a/docs/SECURITY.md +++ b/docs/SECURITY.md @@ -46,6 +46,11 @@ CoDev separates the Vercel control plane from untrusted Firecracker guests. treats an already-missing sandbox as success. - Quotas bound active workspaces, queued turns, daily turns, terminal sessions, publication size, and control-plane request rates. +- A connected Claude subscription keeps its own credential-bound microVM for + three minutes after a reply. Each reply checkpoints that VM before releasing + the execution lease. An idle VM yields its host slot immediately when new + active work needs capacity; its private checkpoint remains available for the + same connection. Disconnect destroys the VM and its checkpoint. ## Design-partner feedback diff --git a/services/orchestrator/src/backend/firecracker.rs b/services/orchestrator/src/backend/firecracker.rs index 859e6f2ef..31bd9606e 100644 --- a/services/orchestrator/src/backend/firecracker.rs +++ b/services/orchestrator/src/backend/firecracker.rs @@ -43,6 +43,7 @@ use crate::{ }; const GUEST_PORT: u32 = 52; +const EPHEMERAL_PARK_TTL: Duration = Duration::from_secs(3 * 60); /// Firecracker reserves CIDs 0-2, so guest CIDs start here and the slot index /// is recoverable as `guest_cid - GUEST_CID_BASE`. @@ -197,6 +198,7 @@ struct RunningMachine { jail_dir: PathBuf, slot: u32, reap_on_expiry: bool, + parked: AtomicBool, hibernate_on_idle: bool, hibernating: AtomicBool, reaper_exempt_requests: AtomicUsize, @@ -597,10 +599,20 @@ impl FirecrackerBackend { )); } if let Some(machine) = self.machines.read().await.get(&request.workspace_id) { + if request.ephemeral && machine.reap_on_expiry { + let mut instance = machine.instance.write().expect("machine lock"); + instance.expires_at = request.expires_at; + instance.last_activity_at = Utc::now(); + machine.parked.store(false, Ordering::Release); + return Ok(instance.clone()); + } return Ok(machine.instance.read().expect("machine lock").clone()); } if self.machines.read().await.len() >= self.config.max_sandboxes { - return Err(RuntimeError::CapacityExceeded); + self.reclaim_parked_slot().await?; + if self.machines.read().await.len() >= self.config.max_sandboxes { + return Err(RuntimeError::CapacityExceeded); + } } let machine = Arc::new(self.prepare_and_start(&request).await?); @@ -612,6 +624,65 @@ impl FirecrackerBackend { Ok(instance) } + async fn reclaim_parked_slot(&self) -> Result<()> { + let candidate = { + let machines = self.machines.read().await; + machines + .iter() + .filter(|(_, machine)| { + machine.parked.load(Ordering::Acquire) && Arc::strong_count(machine) == 1 + }) + .min_by_key(|(_, machine)| { + machine.instance.read().expect("machine lock").expires_at + }) + .map(|(id, machine)| (id.clone(), machine.clone())) + }; + if let Some((workspace_id, machine)) = candidate { + self.stop_machine(machine).await?; + self.machines.write().await.remove(&workspace_id); + info!(%workspace_id, "reclaimed parked Firecracker sandbox slot"); + } + Ok(()) + } + + /// Only a checkpointed credential-bound VM can enter the short idle window. + pub async fn park(&self, workspace_id: &str) -> Result { + let _guard = self.provision.lock().await; + let machine = self + .machines + .read() + .await + .get(workspace_id) + .cloned() + .ok_or(RuntimeError::SandboxNotFound)?; + if !machine.reap_on_expiry { + return Err(RuntimeError::BadRequest( + "only ephemeral sandboxes can be parked".into(), + )); + } + if !matches!( + self.snapshot_metadata(workspace_id).await?, + Some(( + _, + MicroVmSnapshotMetadata { + kind: SnapshotKind::FullMachine, + .. + } + )) + ) { + return Err(RuntimeError::Conflict( + "sandbox must have a complete VM checkpoint before parking".into(), + )); + } + let mut instance = machine.instance.write().expect("machine lock"); + let now = Utc::now(); + instance.last_activity_at = now; + instance.expires_at = + now + chrono::Duration::from_std(EPHEMERAL_PARK_TTL).map_err(RuntimeError::internal)?; + machine.parked.store(true, Ordering::Release); + Ok(instance.clone()) + } + pub async fn get(&self, workspace_id: &str) -> Result { let machine = self.machine(workspace_id).await?; Ok(machine.instance.read().expect("machine lock").clone()) @@ -1156,7 +1227,9 @@ impl FirecrackerBackend { .await .map_err(RuntimeError::internal)?; for (name, source) in files { - link_or_copy(&source, &staging.join(name)).await?; + // This VM remains live after the checkpoint. A hard link + // would let later writes corrupt the durable snapshot. + clone_or_copy(&source, &staging.join(name)).await?; } let metadata = serde_json::to_vec(&MicroVmSnapshotMetadata { head_sha: head_sha.to_owned(), @@ -1168,10 +1241,27 @@ impl FirecrackerBackend { .await .map_err(RuntimeError::internal)?; let destination = self.snapshot_dir(&machine.workspace_id()); - remove_directory_if_present(&destination).await?; - fs::rename(&staging, &destination) - .await - .map_err(RuntimeError::internal)?; + let previous = self.previous_snapshot_dir(&machine.workspace_id()); + let had_previous = match fs::metadata(&destination).await { + Ok(_) => { + remove_directory_if_present(&previous).await?; + fs::rename(&destination, &previous) + .await + .map_err(RuntimeError::internal)?; + true + } + Err(error) if error.kind() == ErrorKind::NotFound => false, + Err(error) => return Err(RuntimeError::internal(error)), + }; + if let Err(error) = fs::rename(&staging, &destination).await { + if had_previous { + let _ = fs::rename(&previous, &destination).await; + } + return Err(RuntimeError::internal(error)); + } + if let Err(error) = remove_directory_if_present(&previous).await { + warn!(workspace_id = %machine.workspace_id(), %error, "could not remove prior VM checkpoint"); + } Ok::<(), RuntimeError>(()) } .await; @@ -1187,6 +1277,7 @@ impl FirecrackerBackend { snapshot_ms = started_at.elapsed().as_millis() as u64, "firecracker snapshot persisted" ); + api.resume().await?; Ok(()) } @@ -1512,6 +1603,7 @@ impl FirecrackerBackend { jail_dir, slot, reap_on_expiry: request.ephemeral, + parked: AtomicBool::new(false), hibernate_on_idle: request.hibernate_on_idle, hibernating: AtomicBool::new(false), reaper_exempt_requests: AtomicUsize::new(0), diff --git a/services/orchestrator/src/backend/mod.rs b/services/orchestrator/src/backend/mod.rs index e1edb505d..d9a860260 100644 --- a/services/orchestrator/src/backend/mod.rs +++ b/services/orchestrator/src/backend/mod.rs @@ -249,6 +249,14 @@ impl Backend { } } + pub async fn park(&self, workspace_id: &str) -> Result { + match self { + Self::Fake(backend) => backend.park(workspace_id), + #[cfg(target_os = "linux")] + Self::Firecracker(backend) => backend.park(workspace_id).await, + } + } + pub async fn destroy(&self, workspace_id: &str) -> Result<()> { match self { Self::Fake(backend) => backend.destroy(workspace_id), @@ -769,6 +777,7 @@ pub type SharedBackend = Arc; pub struct FakeBackend { instances: RwLock>, ephemeral: RwLock>, + parked: RwLock>, host_shutdown_pending: std::sync::atomic::AtomicBool, max: usize, } @@ -784,6 +793,7 @@ impl FakeBackend { Self { instances: RwLock::new(HashMap::new()), ephemeral: RwLock::new(HashSet::new()), + parked: RwLock::new(HashSet::new()), host_shutdown_pending: std::sync::atomic::AtomicBool::new(false), max: MAX_ACTIVE_SESSIONS, } @@ -821,6 +831,10 @@ impl FakeBackend { !ephemeral.contains(workspace_id) || instance.expires_at > now }); ephemeral.retain(|workspace_id| instances.contains_key(workspace_id)); + self.parked + .write() + .expect("fake backend lock") + .retain(|workspace_id| instances.contains_key(workspace_id)); before - instances.len() } @@ -834,11 +848,35 @@ impl FakeBackend { "Firecracker host is shutting down".into(), )); } - if let Some(instance) = instances.get(&request.workspace_id) { + if let Some(instance) = instances.get_mut(&request.workspace_id) { + if request.ephemeral + && self + .ephemeral + .read() + .expect("fake backend lock") + .contains(&request.workspace_id) + { + instance.expires_at = request.expires_at; + instance.last_activity_at = Utc::now(); + self.parked + .write() + .expect("fake backend lock") + .remove(&request.workspace_id); + } return Ok(instance.clone()); } if instances.len() >= self.max { - return Err(RuntimeError::CapacityExceeded); + let mut parked = self.parked.write().expect("fake backend lock"); + if let Some(id) = parked.iter().next().cloned() { + parked.remove(&id); + instances.remove(&id); + self.ephemeral + .write() + .expect("fake backend lock") + .remove(&id); + } else { + return Err(RuntimeError::CapacityExceeded); + } } let now = Utc::now(); let instance = Instance { @@ -880,6 +918,31 @@ impl FakeBackend { Ok(instance.clone()) } + fn park(&self, workspace_id: &str) -> Result { + let mut instances = self.instances.write().expect("fake backend lock"); + if !self + .ephemeral + .read() + .expect("fake backend lock") + .contains(workspace_id) + { + return Err(RuntimeError::BadRequest( + "only ephemeral sandboxes can be parked".into(), + )); + } + let instance = instances + .get_mut(workspace_id) + .ok_or(RuntimeError::SandboxNotFound)?; + let now = Utc::now(); + instance.last_activity_at = now; + instance.expires_at = now + Duration::minutes(3); + self.parked + .write() + .expect("fake backend lock") + .insert(workspace_id.to_owned()); + Ok(instance.clone()) + } + fn destroy(&self, workspace_id: &str) -> Result<()> { let removed = self .instances @@ -894,6 +957,10 @@ impl FakeBackend { .write() .expect("fake backend lock") .remove(workspace_id); + self.parked + .write() + .expect("fake backend lock") + .remove(workspace_id); Ok(()) } @@ -1253,4 +1320,50 @@ mod tests { )); assert!(backend.get("live-workspace").await.is_ok()); } + + #[tokio::test] + async fn parked_ephemeral_reuses_its_vm_and_yields_capacity() { + let backend = Backend::fake(); + let mut request = create_request("private-profile"); + request.ephemeral = true; + let first = backend.create(request.clone()).await.expect("first boot"); + let parked = backend.park("private-profile").await.expect("park"); + let idle = (parked.expires_at - Utc::now()).num_seconds(); + assert!((175..=180).contains(&idle)); + + request.expires_at = Utc::now() + Duration::minutes(10); + let reused = backend.create(request).await.expect("reuse warm VM"); + assert_eq!(reused.id, first.id); + assert!(reused.expires_at > parked.expires_at); + + backend.park("private-profile").await.expect("park again"); + for index in 1..MAX_ACTIVE_SESSIONS { + backend + .create(create_request(&format!("busy-{index}"))) + .await + .expect("fill other slots"); + } + backend + .create(create_request("new-active")) + .await + .expect("new work reclaims parked slot"); + assert!(matches!( + backend.get("private-profile").await, + Err(RuntimeError::SandboxNotFound) + )); + assert_eq!(backend.active_count().await, MAX_ACTIVE_SESSIONS); + } + + #[tokio::test] + async fn durable_workspace_cannot_be_parked() { + let backend = Backend::fake(); + backend + .create(create_request("durable")) + .await + .expect("create durable workspace"); + assert!(matches!( + backend.park("durable").await, + Err(RuntimeError::BadRequest(_)) + )); + } } diff --git a/services/orchestrator/src/http_api.rs b/services/orchestrator/src/http_api.rs index 3dc9021ff..8c7be92f2 100644 --- a/services/orchestrator/src/http_api.rs +++ b/services/orchestrator/src/http_api.rs @@ -62,6 +62,7 @@ pub fn router(backend: SharedBackend, ide: IdeBackend) -> Router { ) .route("/v1/sandboxes/{workspace_id}/resume", post(resume_sandbox)) .route("/v1/sandboxes/{workspace_id}/activity", post(touch_sandbox)) + .route("/v1/sandboxes/{workspace_id}/park", post(park_sandbox)) .route( "/v1/sandboxes/{workspace_id}/ide", post(start_ide).get(get_ide).delete(stop_ide), @@ -326,6 +327,15 @@ async fn touch_sandbox( Ok(Json(serde_json::json!({ "sandbox": instance }))) } +async fn park_sandbox( + State(backend): State, + Path(workspace_id): Path, +) -> Result> { + validate_workspace_id(&workspace_id)?; + let instance = backend.park(&workspace_id).await?; + Ok(Json(serde_json::json!({ "sandbox": instance }))) +} + async fn destroy_sandbox( State(backend): State, Path(workspace_id): Path,