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
This commit is contained in:
Peter Steinberger
2026-08-23 06:06:46 -07:00
committed by GitHub
parent 2aa4b694e1
commit 066176351e
11 changed files with 529 additions and 43 deletions
@@ -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<void>) => Promise<void>)
| 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();
});
});
@@ -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}`);
}
@@ -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);
@@ -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<typeof import("./worker-environments/placement-disk-space.js")>();
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<void> } | 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();
}
});
});
@@ -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<void> } | 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 () => {
@@ -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,
@@ -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())
@@ -77,6 +77,7 @@ export type WorkerDispatchEnvironmentService = Pick<
| "createFromProfileSnapshot"
| "destroy"
| "get"
| "reconcileEnvironment"
| "reconcileOnce"
| "startTunnel"
| "stopTunnel"
@@ -121,8 +121,17 @@ export function createPlacementRecoveryActions(deps: PlacementRecoveryDeps) {
}
};
const reconcile = async (): Promise<void> => {
await environments.reconcileOnce();
const reconcile = async (mode?: "startup"): Promise<void> => {
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<string>();
@@ -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}`);
@@ -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,
@@ -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({