From 066176351ebd52b71fd98689d78d6e05e4467bdf Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 23 Aug 2026 06:06:46 -0700 Subject: [PATCH] fix(gateway): avoid failed worker cleanup blocking startup (#128103) * fix(gateway): defer failed worker cleanup after readiness * test(gateway): preserve periodic worker recovery ordering * fix(gateway): track post-ready session retirement * fix(gateway): defer failed teardown past readiness --- ...r-worker-placement-reconcile-guard.test.ts | 98 ++++++ ...server-worker-placement-reconcile-guard.ts | 6 +- ...r-worker-placement-recovery-events.test.ts | 1 + ...r-worker-placement-startup-cleanup.test.ts | 281 ++++++++++++++++++ .../server-worker-placement-startup.test.ts | 153 +++++++--- .../server-worker-placement-startup.ts | 8 +- .../placement-dispatch-coordinator.ts | 2 +- .../placement-dispatch-failure.ts | 1 + .../placement-dispatch-recovery.ts | 18 +- .../placement-dispatch-test-harness.ts | 3 + ...ker-turn-launcher-failure-recovery.test.ts | 1 + 11 files changed, 529 insertions(+), 43 deletions(-) create mode 100644 src/gateway/server-worker-placement-reconcile-guard.test.ts create mode 100644 src/gateway/server-worker-placement-startup-cleanup.test.ts diff --git a/src/gateway/server-worker-placement-reconcile-guard.test.ts b/src/gateway/server-worker-placement-reconcile-guard.test.ts new file mode 100644 index 000000000000..d028ba86716c --- /dev/null +++ b/src/gateway/server-worker-placement-reconcile-guard.test.ts @@ -0,0 +1,98 @@ +import { describe, expect, it, vi } from "vitest"; +import { installWorkerPlacementReconcileGuard } from "./server-worker-placement-reconcile-guard.js"; + +const localClaim = { + owner: "local" as const, + claimId: "claim-cleanup", + runId: "run-cleanup", + generation: 1, + ownerEpoch: null, +}; + +function createFailedPlacementGuard(params: { + turnClaim: typeof localClaim | null; + activeOwnerEpoch: number | null; + destroyRequestedAtMs: number | null; +}) { + let guard: + | ((environmentId: string, reconcileCore: () => Promise) => Promise) + | undefined; + const resumeProvisioning = vi.fn(); + installWorkerPlacementReconcileGuard({ + placements: { + list: () => [ + { + sessionId: "session-cleanup", + state: "failed", + environmentId: "worker-cleanup", + turnClaim: params.turnClaim, + activeOwnerEpoch: params.activeOwnerEpoch, + }, + ], + } as never, + environments: { + get: (environmentId: string) => ({ + environmentId, + state: "provisioning", + destroyRequestedAtMs: params.destroyRequestedAtMs, + }), + installReconcileEnvironmentGuard: (installed: typeof guard) => { + guard = installed; + return async () => {}; + }, + } as never, + dispatch: { resumeProvisioning } as never, + isStopping: () => false, + }); + if (!guard) { + throw new Error("worker placement reconciliation guard was not installed"); + } + return { guard, resumeProvisioning }; +} + +describe("worker placement reconciliation teardown authority", () => { + it("allows destruction only after its exact failed owner has released all authority", async () => { + const { guard, resumeProvisioning } = createFailedPlacementGuard({ + turnClaim: null, + activeOwnerEpoch: null, + destroyRequestedAtMs: 1, + }); + const reconcileCore = vi.fn(async () => {}); + + await guard("worker-cleanup", reconcileCore); + + expect(reconcileCore).toHaveBeenCalledOnce(); + expect(resumeProvisioning).not.toHaveBeenCalled(); + }); + + it.each([ + { + reason: "a retained local turn claim", + turnClaim: localClaim, + activeOwnerEpoch: null, + destroyRequestedAtMs: 1, + }, + { + reason: "an active owner epoch", + turnClaim: null, + activeOwnerEpoch: 2, + destroyRequestedAtMs: 1, + }, + { + reason: "no durable destruction request", + turnClaim: null, + activeOwnerEpoch: null, + destroyRequestedAtMs: null, + }, + ])("keeps failed-placement cleanup fenced with $reason", async ({ reason: _, ...params }) => { + const { guard, resumeProvisioning } = createFailedPlacementGuard(params); + const reconcileCore = vi.fn(async () => {}); + + await expect(guard("worker-cleanup", reconcileCore)).rejects.toThrow( + "provisioning owner is failed", + ); + + expect(reconcileCore).not.toHaveBeenCalled(); + expect(resumeProvisioning).not.toHaveBeenCalled(); + }); +}); diff --git a/src/gateway/server-worker-placement-reconcile-guard.ts b/src/gateway/server-worker-placement-reconcile-guard.ts index af18f863dc75..3c182cb7a5a1 100644 --- a/src/gateway/server-worker-placement-reconcile-guard.ts +++ b/src/gateway/server-worker-placement-reconcile-guard.ts @@ -29,7 +29,11 @@ export function installWorkerPlacementReconcileGuard(params: { owner && (environment?.state === "requested" || environment?.state === "provisioning" || - environment?.state === "bootstrapping") + environment?.state === "bootstrapping") && + (owner.state !== "failed" || + owner.turnClaim !== null || + owner.activeOwnerEpoch !== null || + environment.destroyRequestedAtMs === null) ) { throw new Error(`Worker environment ${environmentId} provisioning owner is ${owner.state}`); } diff --git a/src/gateway/server-worker-placement-recovery-events.test.ts b/src/gateway/server-worker-placement-recovery-events.test.ts index 2fdd73e99f18..5fde4bd836ad 100644 --- a/src/gateway/server-worker-placement-recovery-events.test.ts +++ b/src/gateway/server-worker-placement-recovery-events.test.ts @@ -97,6 +97,7 @@ describe("worker placement recovery session events", () => { throw new Error("worker placement runtime did not start"); } + await vi.dynamicImportSettled(); await vi.advanceTimersByTimeAsync(60_000); expect(context.broadcastToConnIds).not.toHaveBeenCalled(); expect(readSessionsMutationVersion(context)).toBe(initialMutationVersion); diff --git a/src/gateway/server-worker-placement-startup-cleanup.test.ts b/src/gateway/server-worker-placement-startup-cleanup.test.ts new file mode 100644 index 000000000000..8b74702b55ef --- /dev/null +++ b/src/gateway/server-worker-placement-startup-cleanup.test.ts @@ -0,0 +1,281 @@ +import { describe, expect, it, vi } from "vitest"; +import { createDeferredCore } from "../shared/deferred.js"; + +const runtimeFactoryMocks = vi.hoisted(() => ({ + createDiskSpace: vi.fn(), + createSessionEvidenceResolver: vi.fn(), +})); + +vi.mock("./server-worker-placement-session-evidence.js", () => ({ + createWorkerPlacementSessionEvidenceResolver: runtimeFactoryMocks.createSessionEvidenceResolver, +})); + +vi.mock("./worker-environments/placement-disk-space.js", async (importOriginal) => { + const actual = + await importOriginal(); + return { + ...actual, + createWorkerPlacementDiskSpaceMonitor: runtimeFactoryMocks.createDiskSpace, + }; +}); + +import { createGatewayWorkerPlacementRuntime } from "./server-worker-placement-startup.js"; +import { seedActivePlacement } from "./worker-environments/placement-dispatch-test-fixtures.js"; +import { createWorkerSessionPlacementStore } from "./worker-environments/placement-store.js"; +import * as workerEnvironmentSupport from "./worker-environments/service.test-support.js"; + +describe("worker placement startup cleanup ownership", () => { + workerEnvironmentSupport.setupWorkerEnvironmentServiceSuite(); + + it.each([ + { retainedAuthority: "a local turn claim", claimLocalTurn: true }, + { retainedAuthority: "an active owner epoch", claimLocalTurn: false }, + ])( + "never reaches worker providers during startup while a failed placement retains $retainedAuthority", + async ({ claimLocalTurn }) => { + const environmentId = "worker-startup-fenced"; + const provision = vi.fn(async () => ({ + leaseId: "lease-startup-fenced", + ssh: workerEnvironmentSupport.SSH_ENDPOINT, + })); + const inspect = vi.fn(async () => ({ status: "active" as const })); + const destroy = vi.fn(async () => {}); + const environments = workerEnvironmentSupport.createService( + workerEnvironmentSupport.createProvider({ provision, inspect, destroy }), + ); + const requestedEnvironment = workerEnvironmentSupport.testState.store.createIntent({ + environmentId, + providerId: "fake", + profileId: "development", + profileSnapshot: { settings: { region: "test" } }, + provisionOperationId: "provision:startup-fenced", + }); + workerEnvironmentSupport.testState.store.transition({ + environmentId, + from: requestedEnvironment.state, + to: "provisioning", + }); + workerEnvironmentSupport.testState.store.requestDestroy({ + environmentId, + state: "provisioning", + }); + const placements = createWorkerSessionPlacementStore({ + database: workerEnvironmentSupport.testState.stateDb, + now: () => workerEnvironmentSupport.testState.nowMs, + }); + let failed; + if (claimLocalTurn) { + const identity = { + sessionId: "session-startup-fenced", + sessionKey: "agent:main:startup-fenced", + agentId: "main", + }; + placements.claimTurn({ + ...identity, + owner: { kind: "local" }, + claimId: "startup-fenced-local-claim", + runId: "startup-fenced-local-run", + }); + const requested = placements.startDispatch({ ...identity, executionMode: "remote-exec" }); + failed = placements.fail({ + sessionId: requested.sessionId, + expectedGeneration: requested.generation, + recoveryError: "startup worker placement failed before its local claim was released", + }); + // Current dispatch binds environments only after its local turn drains; preserve a + // crash-era terminal environment reference in SQLite without fabricating record shapes. + workerEnvironmentSupport.testState.stateDb.db + .prepare("UPDATE worker_session_placements SET environment_id = ? WHERE session_id = ?") + .run(environmentId, failed.sessionId); + failed = placements.get(failed.sessionId); + } else { + const active = seedActivePlacement(placements, { + environmentId, + ownerEpoch: 1, + executionMode: "remote-exec", + }); + if (active.state !== "active") { + throw new Error("startup authority fixture did not produce an active placement"); + } + const draining = placements.startDrain({ + sessionId: active.sessionId, + environmentId, + ownerEpoch: active.activeOwnerEpoch, + expectedGeneration: active.generation, + }); + const reconciling = placements.startReconcile({ + sessionId: draining.sessionId, + environmentId, + ownerEpoch: active.activeOwnerEpoch, + expectedGeneration: draining.generation, + }); + failed = placements.fail({ + sessionId: reconciling.sessionId, + expectedGeneration: reconciling.generation, + recoveryError: "startup worker placement failed before its owner epoch was released", + }); + } + expect(failed).toMatchObject({ + state: "failed", + environmentId, + activeOwnerEpoch: claimLocalTurn ? null : 1, + turnClaim: claimLocalTurn ? expect.objectContaining({ owner: "local" }) : null, + }); + runtimeFactoryMocks.createSessionEvidenceResolver.mockResolvedValue(async () => "current"); + runtimeFactoryMocks.createDiskSpace.mockReturnValue({ + read: vi.fn(), + version: vi.fn(() => 0), + sweep: vi.fn().mockResolvedValue(undefined), + }); + const runtime = createGatewayWorkerPlacementRuntime({ + placements, + environments, + gatewayNamespace: "gateway-test", + revokeSessionAuthority: vi.fn(), + warn: vi.fn(), + }); + const sidecar = await runtime.startRuntime({ + isClosePreludeStarted: () => false, + registerSidecar: vi.fn(), + unregisterSidecar: vi.fn(), + }); + try { + expect(sidecar).not.toBeNull(); + await environments.reconcileOnce(); + expect(provision).not.toHaveBeenCalled(); + expect(inspect).not.toHaveBeenCalled(); + expect(destroy).not.toHaveBeenCalled(); + expect(workerEnvironmentSupport.testState.store.get(environmentId)).toMatchObject({ + state: "provisioning", + leaseId: null, + destroyRequestedAtMs: expect.any(Number), + }); + } finally { + await sidecar?.stop(); + } + }, + ); + + it("becomes ready before failed-placement lease adoption and drains its exact teardown", async () => { + const environmentId = "worker-startup-indeterminate"; + const operationId = "provision:startup-indeterminate"; + const releaseProvision = createDeferredCore(); + const provisionStarted = createDeferredCore(); + const provision = vi.fn(async () => { + provisionStarted.resolve(); + await releaseProvision.promise; + return { leaseId: "lease-startup-adopted", ssh: workerEnvironmentSupport.SSH_ENDPOINT }; + }); + const destroy = vi.fn(async () => {}); + const environments = workerEnvironmentSupport.createService( + workerEnvironmentSupport.createProvider({ provision, destroy }), + ); + const intent = workerEnvironmentSupport.testState.store.createIntent({ + environmentId, + providerId: "fake", + profileId: "development", + profileSnapshot: { settings: { region: "test" } }, + provisionOperationId: operationId, + }); + workerEnvironmentSupport.testState.store.transition({ + environmentId, + from: intent.state, + to: "provisioning", + }); + workerEnvironmentSupport.testState.store.requestDestroy({ + environmentId, + state: "provisioning", + }); + const placements = createWorkerSessionPlacementStore({ + database: workerEnvironmentSupport.testState.stateDb, + now: () => workerEnvironmentSupport.testState.nowMs, + }); + const requested = placements.startDispatch({ + sessionId: "session-startup-indeterminate", + sessionKey: "agent:main:startup-indeterminate", + agentId: "main", + executionMode: "remote-exec", + }); + const provisioning = placements.transition({ + sessionId: requested.sessionId, + from: "requested", + to: "provisioning", + expectedGeneration: requested.generation, + patch: { environmentId }, + }); + const failed = placements.fail({ + sessionId: provisioning.sessionId, + expectedGeneration: provisioning.generation, + recoveryError: "startup worker placement failed", + }); + expect(failed).toMatchObject({ + state: "failed", + environmentId, + turnClaim: null, + activeOwnerEpoch: null, + }); + runtimeFactoryMocks.createSessionEvidenceResolver.mockResolvedValue(async () => "current"); + runtimeFactoryMocks.createDiskSpace.mockReturnValue({ + read: vi.fn(), + version: vi.fn(() => 0), + sweep: vi.fn().mockResolvedValue(undefined), + }); + const runtime = createGatewayWorkerPlacementRuntime({ + placements, + environments, + gatewayNamespace: "gateway-test", + revokeSessionAuthority: vi.fn(), + warn: vi.fn(), + }); + let registeredSidecar: { stop: () => Promise } | undefined; + const starting = runtime.startRuntime({ + isClosePreludeStarted: () => false, + registerSidecar: (sidecar) => { + registeredSidecar = sidecar; + }, + unregisterSidecar: vi.fn(), + }); + + try { + const first = await Promise.race([ + starting.then((sidecar) => ({ kind: "ready" as const, sidecar })), + provisionStarted.promise.then(() => ({ kind: "provision" as const })), + ]); + expect(first.kind).toBe("ready"); + if (first.kind !== "ready" || !first.sidecar) { + throw new Error("worker startup did not reach readiness before provider adoption"); + } + await provisionStarted.promise; + expect(provision).toHaveBeenCalledWith({ region: "test" }, operationId, undefined); + expect(workerEnvironmentSupport.testState.store.get(environmentId)).toMatchObject({ + state: "provisioning", + leaseId: null, + destroyRequestedAtMs: expect.any(Number), + }); + + let stopped = false; + const stopping = first.sidecar.stop().then(() => { + stopped = true; + }); + await Promise.resolve(); + expect(stopped).toBe(false); + expect(destroy).not.toHaveBeenCalled(); + + releaseProvision.resolve(); + await stopping; + + expect(destroy).toHaveBeenCalledWith({ + leaseId: "lease-startup-adopted", + profile: { region: "test" }, + }); + expect(workerEnvironmentSupport.testState.store.get(environmentId)).toMatchObject({ + state: "destroyed", + leaseId: "lease-startup-adopted", + }); + } finally { + releaseProvision.resolve(); + await starting.catch(() => undefined); + await registeredSidecar?.stop(); + } + }); +}); diff --git a/src/gateway/server-worker-placement-startup.test.ts b/src/gateway/server-worker-placement-startup.test.ts index 7b64452824b0..fc0b23809296 100644 --- a/src/gateway/server-worker-placement-startup.test.ts +++ b/src/gateway/server-worker-placement-startup.test.ts @@ -212,6 +212,7 @@ describe("worker placement startup health lifetime", () => { }); expect(sidecar).not.toBeNull(); + expect(reconcileActive).not.toHaveBeenCalled(); expect(diskSpace.sweep).toHaveBeenCalledOnce(); await vi.advanceTimersByTimeAsync(60_000); expect(reconcileActive).toHaveBeenCalledOnce(); @@ -236,10 +237,89 @@ describe("worker placement startup health lifetime", () => { } }); - it("drains deferred startup session evidence before stopping environments", async () => { - const evidence = createDeferredCore<"current">(); - runtimeFactoryMocks.resolveSessionEvidence.mockImplementation(async () => evidence.promise); - runtimeFactoryMocks.createSessionEvidenceResolver.mockResolvedValue( + it.each(["provisioning", "active"] as const)( + "waits for %s placement authority recovery before exposing readiness", + async (state) => { + const releaseRecovery = createDeferredCore(); + const reconcile = vi.fn(async () => await releaseRecovery.promise); + runtimeFactoryMocks.createDiskSpace.mockReturnValue({ + read: vi.fn(), + version: vi.fn(() => 0), + sweep: vi.fn().mockResolvedValue(undefined), + }); + runtimeFactoryMocks.createDispatch.mockReturnValue({ + dispatch: vi.fn(), + forceDestroyEnvironment: vi.fn(), + reclaim: vi.fn(), + reconcile, + reconcileActive: vi.fn().mockResolvedValue(undefined), + }); + runtimeFactoryMocks.createSessionEvidenceResolver.mockResolvedValue(async () => "current"); + const placement = { + sessionId: `session-startup-${state}`, + sessionKey: `agent:main:startup-${state}`, + agentId: "main", + state, + generation: 1, + environmentId: `worker-startup-${state}`, + activeOwnerEpoch: state === "active" ? 1 : null, + turnClaim: null, + }; + const environments = { + installReconcileEnvironmentGuard: vi.fn(() => vi.fn()), + start: vi.fn(), + stop: vi.fn().mockResolvedValue(undefined), + }; + const runtime = createGatewayWorkerPlacementRuntime({ + placements: { + workspaceResultInstanceId: () => "gateway-test", + get: () => placement, + list: () => [placement], + retireSessionPlacement: vi.fn(), + pruneOrphanedWorkspaceReconciliations: () => [], + listWorkspaceReconciliationOwners: () => [], + listPendingWorkspaceResults: () => [], + } as never, + environments: environments as never, + gatewayNamespace: "gateway-test", + revokeSessionAuthority: vi.fn(), + warn: vi.fn(), + }); + const starting = runtime.startRuntime({ + isClosePreludeStarted: () => false, + registerSidecar: vi.fn(), + unregisterSidecar: vi.fn(), + }); + + try { + await vi.waitFor(() => expect(reconcile).toHaveBeenCalledWith("startup")); + expect(environments.start).not.toHaveBeenCalled(); + + let ready = false; + void starting.then(() => { + ready = true; + }); + await Promise.resolve(); + expect(ready).toBe(false); + + releaseRecovery.resolve(); + const sidecar = await starting; + expect(sidecar).not.toBeNull(); + expect(environments.start).toHaveBeenCalledOnce(); + await sidecar?.stop(); + } finally { + releaseRecovery.resolve(); + await starting.catch(() => undefined); + } + }, + ); + + it("immediately retires absent sessions after readiness and drains retirement on stop", async () => { + const evidence = createDeferredCore<"absent">(); + const reconcileActive = vi.fn().mockResolvedValue(undefined); + const retireSessionPlacement = vi.fn(); + runtimeFactoryMocks.resolveSessionEvidence.mockImplementationOnce(async () => evidence.promise); + runtimeFactoryMocks.createSessionEvidenceResolver.mockResolvedValueOnce( runtimeFactoryMocks.resolveSessionEvidence, ); runtimeFactoryMocks.createDiskSpace.mockReturnValue({ @@ -252,7 +332,7 @@ describe("worker placement startup health lifetime", () => { forceDestroyEnvironment: vi.fn(), reclaim: vi.fn(), reconcile: vi.fn().mockResolvedValue(undefined), - reconcileActive: vi.fn().mockResolvedValue(undefined), + reconcileActive, }); const placement = { sessionId: "session-startup", @@ -286,7 +366,7 @@ describe("worker placement startup health lifetime", () => { workspaceResultInstanceId: () => "gateway-test", get: () => placement, list: () => [placement], - retireSessionPlacement: vi.fn(), + retireSessionPlacement, pruneOrphanedWorkspaceReconciliations: () => [], listWorkspaceReconciliationOwners: () => [], listPendingWorkspaceResults: () => [], @@ -296,39 +376,40 @@ describe("worker placement startup health lifetime", () => { revokeSessionAuthority: vi.fn(), warn: vi.fn(), }); - let closeStarted = false; let sidecar: { stop: () => Promise } | undefined; const unregisterSidecar = vi.fn(); - const starting = runtime.startRuntime({ - isClosePreludeStarted: () => closeStarted, - registerSidecar: (registered) => { - sidecar = registered; - }, - unregisterSidecar, - }); - await vi.waitFor(() => expect(runtimeFactoryMocks.resolveSessionEvidence).toHaveBeenCalled()); - closeStarted = true; - const stopping = sidecar?.stop(); - const repeatedStop = sidecar?.stop(); - if (!stopping || !repeatedStop) { - throw new Error("startup did not register its placement sidecar"); - } - let repeatedStopSettled = false; - void repeatedStop.then(() => { - repeatedStopSettled = true; - }); + try { + const starting = runtime.startRuntime({ + isClosePreludeStarted: () => false, + registerSidecar: (registered) => { + sidecar = registered; + }, + unregisterSidecar, + }); + await expect(starting).resolves.toBe(sidecar); + expect(runtimeFactoryMocks.resolveSessionEvidence).toHaveBeenCalledOnce(); + expect(reconcileActive).not.toHaveBeenCalled(); + expect(retireSessionPlacement).not.toHaveBeenCalled(); - await Promise.resolve(); - expect(repeatedStop).toBe(stopping); - expect(repeatedStopSettled).toBe(false); - expect(environments.stopNodeEnrollmentWaits).toHaveBeenCalledOnce(); - expect(environments.stop).not.toHaveBeenCalled(); - evidence.resolve("current"); - await expect(starting).resolves.toBeNull(); - await Promise.all([stopping, repeatedStop]); - expect(environments.stop).toHaveBeenCalledOnce(); - expect(unregisterSidecar).toHaveBeenCalledOnce(); - expect(unregisterSidecar).toHaveBeenCalledWith(sidecar); + const stopping = sidecar?.stop(); + const repeatedStop = sidecar?.stop(); + if (!stopping || !repeatedStop) { + throw new Error("startup did not register its placement sidecar"); + } + + await Promise.resolve(); + expect(repeatedStop).toBe(stopping); + expect(environments.stopNodeEnrollmentWaits).toHaveBeenCalledOnce(); + expect(environments.stop).not.toHaveBeenCalled(); + evidence.resolve("absent"); + await Promise.all([stopping, repeatedStop]); + expect(retireSessionPlacement).toHaveBeenCalledOnce(); + expect(environments.stop).toHaveBeenCalledOnce(); + expect(unregisterSidecar).not.toHaveBeenCalled(); + } finally { + evidence.resolve("absent"); + await sidecar?.stop(); + } }); it("retries worker environment cleanup after a failed stop attempt", async () => { diff --git a/src/gateway/server-worker-placement-startup.ts b/src/gateway/server-worker-placement-startup.ts index 7cf5e29e8519..f80498f16fb2 100644 --- a/src/gateway/server-worker-placement-startup.ts +++ b/src/gateway/server-worker-placement-startup.ts @@ -623,8 +623,7 @@ export function createGatewayWorkerPlacementRuntime( return await stopBeforeReady(); } const startupReconcile = (async () => { - await dispatchService.reconcile(); - await sessionRetirement.reconcile(); + await dispatchService.reconcile("startup"); await reconcilePublications(); })(); placementReconcile.current = startupReconcile; @@ -647,6 +646,11 @@ export function createGatewayWorkerPlacementRuntime( if (hooks.isClosePreludeStarted()) { return await stopBeforeReady(); } + void trackOperation( + placementReconcile, + sessionRetirement.reconcile(), + "Worker placement reconcile sweep failed", + ); void sweepDiskSpace(); placementReconcileInterval = setInterval( sweepActivePlacements, diff --git a/src/gateway/worker-environments/placement-dispatch-coordinator.ts b/src/gateway/worker-environments/placement-dispatch-coordinator.ts index f43c647fb7f0..bcdaca8a7ea7 100644 --- a/src/gateway/worker-environments/placement-dispatch-coordinator.ts +++ b/src/gateway/worker-environments/placement-dispatch-coordinator.ts @@ -175,7 +175,7 @@ export function coordinateWorkerPlacementDispatch( }, reclaim: async (request, authorize) => await runExclusivePlacementOperation(() => service.reclaim(request, authorize)), - reconcile: () => runReconciliation(service.reconcile), + reconcile: (mode) => runReconciliation(() => service.reconcile(mode)), reconcileActive: (environmentId) => environmentId === undefined ? runReconciliation(() => service.reconcileActive()) diff --git a/src/gateway/worker-environments/placement-dispatch-failure.ts b/src/gateway/worker-environments/placement-dispatch-failure.ts index d563137aa23e..bf52cbe11073 100644 --- a/src/gateway/worker-environments/placement-dispatch-failure.ts +++ b/src/gateway/worker-environments/placement-dispatch-failure.ts @@ -77,6 +77,7 @@ export type WorkerDispatchEnvironmentService = Pick< | "createFromProfileSnapshot" | "destroy" | "get" + | "reconcileEnvironment" | "reconcileOnce" | "startTunnel" | "stopTunnel" diff --git a/src/gateway/worker-environments/placement-dispatch-recovery.ts b/src/gateway/worker-environments/placement-dispatch-recovery.ts index 54e4fd9fd164..ad7c4baa9010 100644 --- a/src/gateway/worker-environments/placement-dispatch-recovery.ts +++ b/src/gateway/worker-environments/placement-dispatch-recovery.ts @@ -121,8 +121,17 @@ export function createPlacementRecoveryActions(deps: PlacementRecoveryDeps) { } }; - const reconcile = async (): Promise => { - await environments.reconcileOnce(); + const reconcile = async (mode?: "startup"): Promise => { + if (mode === "startup") { + // Readiness fences live owners; unowned teardown remains in the service-owned sweep. + for (const { environmentId, state } of placements.listForReconcile()) { + if (environmentId && state !== "failed" && state !== "reclaimed") { + await environments.reconcileEnvironment(environmentId); + } + } + } else { + await environments.reconcileOnce(); + } const pendingResultOwners = await recoverPendingWorkspaceResults(deps, true); const journalOwners = blockingWorkspaceJournalSessions(placements); const moveOwners = (await deps.recoverPlacementMoves?.()) ?? new Set(); @@ -172,7 +181,10 @@ export function createPlacementRecoveryActions(deps: PlacementRecoveryDeps) { continue; } if (isFailedPlacement(placement)) { - await failure.retryFailedTeardown(placement); + // Terminal cleanup never gates readiness; tracked post-start owners resume it safely. + if (mode !== "startup") { + await failure.retryFailedTeardown(placement); + } continue; } const error = new Error(`Worker dispatch interrupted in ${placement.state}`); diff --git a/src/gateway/worker-environments/placement-dispatch-test-harness.ts b/src/gateway/worker-environments/placement-dispatch-test-harness.ts index 03e795e1b7d2..3bfe5f9834ac 100644 --- a/src/gateway/worker-environments/placement-dispatch-test-harness.ts +++ b/src/gateway/worker-environments/placement-dispatch-test-harness.ts @@ -367,6 +367,9 @@ export function createHarness( reconcileOnce: vi.fn(async () => { log.push("environment:reconcile"); }), + reconcileEnvironment: vi.fn(async () => { + log.push("environment:reconcile"); + }), }; const service = createWorkerPlacementDispatchService({ placements, diff --git a/src/gateway/worker-environments/worker-turn-launcher-failure-recovery.test.ts b/src/gateway/worker-environments/worker-turn-launcher-failure-recovery.test.ts index 9343fff6e0f1..ee8043a64cee 100644 --- a/src/gateway/worker-environments/worker-turn-launcher-failure-recovery.test.ts +++ b/src/gateway/worker-environments/worker-turn-launcher-failure-recovery.test.ts @@ -203,6 +203,7 @@ describe("worker turn launcher failure recovery", () => { throw new Error("unexpected inherited worker environment creation"); }), reconcileOnce, + reconcileEnvironment: vi.fn(), }; const workspaceOperations = createWorkerWorkspaceOperationCoordinator(); const dispatch = createWorkerPlacementDispatchService({