diff --git a/src/gateway/server-methods/environments.test.ts b/src/gateway/server-methods/environments.test.ts index a1d9fb2af213..7b5c64fae6a2 100644 --- a/src/gateway/server-methods/environments.test.ts +++ b/src/gateway/server-methods/environments.test.ts @@ -33,9 +33,10 @@ type TestWorkerService = { function mockContext( workerEnvironmentService?: TestWorkerService, reconcileActive: (environmentId?: string) => Promise = vi.fn(async () => {}), - forceDestroyEnvironment: (environmentId: string) => Promise = vi.fn(async () => - workerRecord({ state: "destroyed" }), - ), + forceDestroyEnvironment: ( + environmentId: string, + onCleanupError?: (error: unknown) => void, + ) => Promise = vi.fn(async () => workerRecord({ state: "destroyed" })), ) { return { logGateway: { @@ -123,7 +124,10 @@ async function callEnvironmentMethod( options: { service?: TestWorkerService; reconcileActive?: (environmentId?: string) => Promise; - forceDestroyEnvironment?: (environmentId: string) => Promise; + forceDestroyEnvironment?: ( + environmentId: string, + onCleanupError?: (error: unknown) => void, + ) => Promise; } = {}, ) { const respond = vi.fn(); @@ -462,9 +466,7 @@ describe("environment gateway methods", () => { it("durably abandons placement ownership before forced destruction", async () => { const service = workerService(); - const forceDestroyEnvironment = vi.fn( - async (environmentId: string) => await service.destroy(environmentId), - ); + const forceDestroyEnvironment = vi.fn(async () => workerRecord({ state: "destroyed" })); const [ok, payload] = await callEnvironmentMethod( "environments.destroy", @@ -474,11 +476,45 @@ describe("environment gateway methods", () => { expect(ok).toBe(true); expect(payload).toMatchObject({ worker: { state: "destroyed" } }); - expect(forceDestroyEnvironment).toHaveBeenCalledExactlyOnceWith("worker-1"); - expect(service.destroy).toHaveBeenCalledExactlyOnceWith("worker-1"); + expect(forceDestroyEnvironment).toHaveBeenCalledExactlyOnceWith( + "worker-1", + expect.any(Function), + ); + expect(service.destroy).not.toHaveBeenCalled(); expect(service.destroyUnattached).not.toHaveBeenCalled(); }); + it("logs best-effort forced teardown errors without failing the call", async () => { + const service = workerService(); + const forceDestroyEnvironment = vi.fn( + async (_environmentId: string, onCleanupError?: (error: unknown) => void) => { + onCleanupError?.(new Error("provider stop remains pending")); + return workerRecord({ state: "destroying" }); + }, + ); + const context = mockContext( + service, + vi.fn(async () => {}), + forceDestroyEnvironment, + ); + const respond = vi.fn(); + + await environmentsHandlers["environments.destroy"]?.({ + params: { environmentId: "worker-1", force: true }, + respond, + context, + } as never); + + expect(respond).toHaveBeenCalledWith( + true, + expect.objectContaining({ worker: expect.objectContaining({ state: "destroying" }) }), + undefined, + ); + expect(context.logGateway.warn).toHaveBeenCalledWith( + "worker environment forced teardown cleanup failed: Error: provider stop remains pending", + ); + }); + it("reconciles active placements before returning destroyed worker state", async () => { const service = workerService(); const reconcileActive = vi.fn(async () => {}); diff --git a/src/gateway/server-methods/environments.ts b/src/gateway/server-methods/environments.ts index 6911078beca2..69d61c99d580 100644 --- a/src/gateway/server-methods/environments.ts +++ b/src/gateway/server-methods/environments.ts @@ -223,7 +223,11 @@ export const environmentsHandlers: GatewayRequestHandlers = { throw new Error("cloud worker placement control is unavailable"); } const destroyed = params.force - ? await placementService!.forceDestroyEnvironment!(params.environmentId) + ? await placementService!.forceDestroyEnvironment!(params.environmentId, (error) => { + context.logGateway.warn( + `worker environment forced teardown cleanup failed: ${formatForLog(error)}`, + ); + }) : await service.destroyUnattached(params.environmentId); // Destruction is authoritative. Project the dead worker into its owning // placement before returning, or immediate session deletion stays fenced. diff --git a/src/gateway/server-worker-placement-startup.ts b/src/gateway/server-worker-placement-startup.ts index 0d51ae54613c..254c353a3c5d 100644 --- a/src/gateway/server-worker-placement-startup.ts +++ b/src/gateway/server-worker-placement-startup.ts @@ -127,8 +127,10 @@ function coordinateWorkerPlacementDispatch( }; return { dispatch: async (request) => await runPlacementOperation(() => service.dispatch(request)), - forceDestroyEnvironment: (environmentId) => - runExclusivePlacementOperation(() => service.forceDestroyEnvironment(environmentId)), + forceDestroyEnvironment: (environmentId, onCleanupError) => + runExclusivePlacementOperation(() => + service.forceDestroyEnvironment(environmentId, onCleanupError), + ), reclaim: async (request) => await runPlacementOperation(() => service.reclaim(request)), reconcile: () => runReconciliation(service.reconcile), reconcileActive: () => runReconciliation(service.reconcileActive), diff --git a/src/gateway/worker-environments/placement-dispatch-recovery.ts b/src/gateway/worker-environments/placement-dispatch-recovery.ts index 4d24a0ab4679..20b93c88f48a 100644 --- a/src/gateway/worker-environments/placement-dispatch-recovery.ts +++ b/src/gateway/worker-environments/placement-dispatch-recovery.ts @@ -42,6 +42,24 @@ function isFailedPlacement( return placement.state === "failed"; } +function blockingWorkspaceJournalSessions( + placements: PlacementRecoveryDeps["placements"], +): Set { + const sessions = new Set(); + for (const owner of placements.listWorkspaceReconciliationOwners()) { + const placement = placements.get(owner.sessionId); + if ( + (placement?.state === "active" || placement?.state === "draining") && + placement.environmentId === owner.environmentId && + placement.activeOwnerEpoch === owner.ownerEpoch && + placement.generation === owner.placementGeneration + ) { + sessions.add(owner.sessionId); + } + } + return sessions; +} + export function createPlacementRecoveryActions(deps: PlacementRecoveryDeps) { const { environments, failure, placements } = deps; @@ -163,9 +181,7 @@ export function createPlacementRecoveryActions(deps: PlacementRecoveryDeps) { const reconcile = async (): Promise => { await environments.reconcileOnce(); const pendingResultOwners = await recoverPendingWorkspaceResults(deps, true); - const journalOwners = new Set( - placements.listWorkspaceReconciliationOwners().map((owner) => owner.sessionId), - ); + const journalOwners = blockingWorkspaceJournalSessions(placements); for (const placement of placements.listForReconcile()) { if (journalOwners.has(placement.sessionId) || pendingResultOwners.has(placement.sessionId)) { continue; @@ -204,9 +220,7 @@ export function createPlacementRecoveryActions(deps: PlacementRecoveryDeps) { const reconcileActive = async (environmentId?: string): Promise => { await environments.reconcileOnce(); const pendingResultOwners = await recoverPendingWorkspaceResults(deps, false); - const journalOwners = new Set( - placements.listWorkspaceReconciliationOwners().map((owner) => owner.sessionId), - ); + const journalOwners = blockingWorkspaceJournalSessions(placements); for (const placement of placements.listForReconcile()) { if (journalOwners.has(placement.sessionId) || pendingResultOwners.has(placement.sessionId)) { continue; diff --git a/src/gateway/worker-environments/placement-dispatch-test-harness.ts b/src/gateway/worker-environments/placement-dispatch-test-harness.ts index 11f13da12ab0..5402e645563f 100644 --- a/src/gateway/worker-environments/placement-dispatch-test-harness.ts +++ b/src/gateway/worker-environments/placement-dispatch-test-harness.ts @@ -16,7 +16,10 @@ import { } from "./placement-dispatch-test-fixtures.js"; import { createWorkerPlacementDispatchService } from "./placement-dispatch.js"; import type { WorkerTunnelHandle } from "./tunnel.js"; -import { createWorkerWorkspaceOperationCoordinator } from "./workspace-operation-coordinator.js"; +import { + createWorkerWorkspaceOperationCoordinator, + type WorkerWorkspaceOperationCoordinator, +} from "./workspace-operation-coordinator.js"; export function createHarness( placementStore: PlacementStore, @@ -38,6 +41,8 @@ export function createHarness( workspacePath?: string; priorWorkspaceResultConflict?: { paths: string[]; stagedResultRef: string }; reconcileConflictPaths?: string[]; + workspaceOperations?: WorkerWorkspaceOperationCoordinator; + destroyFailureState?: "draining" | "destroying"; } = {}, ) { const reconciledManifestRef = MANIFEST_REF.replaceAll("b", "c"); @@ -54,10 +59,12 @@ export function createHarness( }; const placements: WorkerDispatchPlacementStore = { get: (sessionId) => placementStore.get(sessionId), - loadWorkspaceReconciliation: (owner) => placementStore.loadWorkspaceReconciliation(owner), + loadWorkspaceReconciliation: (owner, loadOptions) => + placementStore.loadWorkspaceReconciliation(owner, loadOptions), beginWorkspaceReconciliation: (owner, journal) => placementStore.beginWorkspaceReconciliation(owner, journal), - abortWorkspaceReconciliation: (owner) => placementStore.abortWorkspaceReconciliation(owner), + abortWorkspaceReconciliation: (owner, abortOptions) => + placementStore.abortWorkspaceReconciliation(owner, abortOptions), listWorkspaceReconciliationOwners: () => placementStore.listWorkspaceReconciliationOwners(), listPendingWorkspaceResults: () => placementStore.listPendingWorkspaceResults(), workspaceResultInstanceId: () => placementStore.workspaceResultInstanceId(), @@ -266,6 +273,13 @@ export function createHarness( destroy: vi.fn(async () => { log.push("teardown:destroy"); if (options.destroyFails) { + if (options.destroyFailureState) { + currentEnvironment = { + ...attached, + state: options.destroyFailureState, + tunnelStatus: "stopped", + }; + } throw new Error("destroy pending"); } const destroyed = destroyedEnvironment((currentEnvironment?.ownerEpoch ?? 1) + 1); @@ -279,7 +293,7 @@ export function createHarness( const service = createWorkerPlacementDispatchService({ placements, environments, - workspaceOperations: createWorkerWorkspaceOperationCoordinator(), + workspaceOperations: options.workspaceOperations ?? createWorkerWorkspaceOperationCoordinator(), runLocalBarrier: async ({ startDispatch }) => { log.push("barrier"); const placement = startDispatch(); diff --git a/src/gateway/worker-environments/placement-dispatch.ts b/src/gateway/worker-environments/placement-dispatch.ts index ed975555c06c..21aac670d913 100644 --- a/src/gateway/worker-environments/placement-dispatch.ts +++ b/src/gateway/worker-environments/placement-dispatch.ts @@ -1,6 +1,7 @@ import { randomUUID } from "node:crypto"; import { createPlacementFailureActions, + isUnavailableEnvironment, type WorkerActivationBarrier, type WorkerActiveDispatchPlacement, type WorkerDispatchEnvironmentService, @@ -457,14 +458,28 @@ export function createWorkerPlacementDispatchService(options: WorkerPlacementDis return { dispatch, - forceDestroyEnvironment: (environmentId: string) => + forceDestroyEnvironment: (environmentId: string, onCleanupError?: (error: unknown) => void) => options.workspaceOperations.run(environmentId, async () => { await forceAbandonWorkerEnvironment({ placements, environmentId, resolveWorkspacePath: options.resolveWorkspacePath, + onCleanupError, }); - return await environments.destroy(environmentId); + try { + return await environments.destroy(environmentId); + } catch (error) { + const current = environments.get(environmentId); + if (!current || !isUnavailableEnvironment(current)) { + throw error; + } + try { + onCleanupError?.(error); + } catch { + // Reporting cannot overturn the durable placement/environment fences. + } + return current; + } }), reclaim, reconcile: recovery.reconcile, diff --git a/src/gateway/worker-environments/placement-force-abandon.test.ts b/src/gateway/worker-environments/placement-force-abandon.test.ts index 80e3219bfbe8..daa247006967 100644 --- a/src/gateway/worker-environments/placement-force-abandon.test.ts +++ b/src/gateway/worker-environments/placement-force-abandon.test.ts @@ -1,7 +1,8 @@ +import { createHash } from "node:crypto"; import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; -import { afterEach, beforeEach, describe, expect, it } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { closeOpenClawStateDatabaseForTest, openOpenClawStateDatabase, @@ -92,4 +93,150 @@ describe("forced worker environment abandonment", () => { }); expect(store.listPendingWorkspaceResults()).toEqual([]); }); + + it("deletes a stale journal without replaying it into the current workspace", async () => { + const store = createWorkerSessionPlacementStore({ database, now: () => 1_000 }); + const { environmentId } = createDispatchEnvironmentFixtures(); + const active = seedActivePlacement(store, { environmentId, ownerEpoch: 2 }); + if (active.state !== "active") { + throw new Error("active placement fixture was not active"); + } + const owner = { + sessionId: active.sessionId, + environmentId: active.environmentId, + ownerEpoch: active.activeOwnerEpoch, + placementGeneration: active.generation, + }; + store.beginWorkspaceReconciliation(owner, { + version: 1, + temporaryNonce: "b".repeat(32), + baseManifestRef: active.workspaceBaseManifestRef, + currentManifestRef: `sha256:${"c".repeat(64)}`, + baseEntries: [], + appliedEntries: [], + baseTree: "f".repeat(40), + basePackSha256: createHash("sha256").update("").digest("hex"), + basePack: Buffer.alloc(0), + }); + const draining = store.startDrain({ + sessionId: active.sessionId, + environmentId: active.environmentId, + ownerEpoch: active.activeOwnerEpoch, + expectedGeneration: active.generation, + }); + if (draining.state !== "draining") { + throw new Error("draining placement fixture was not draining"); + } + store.startReconcile({ + sessionId: draining.sessionId, + environmentId: draining.environmentId, + ownerEpoch: draining.activeOwnerEpoch, + expectedGeneration: draining.generation, + }); + const resolveWorkspacePath = vi.fn(async () => root); + + await forceAbandonWorkerEnvironment({ + placements: store, + environmentId, + resolveWorkspacePath, + }); + + expect(resolveWorkspacePath).not.toHaveBeenCalled(); + expect(store.listWorkspaceReconciliationOwners()).toEqual([]); + expect(store.get(REQUEST.sessionId)).toMatchObject({ state: "failed" }); + }); + + it("retains a current journal when its best-effort rollback fails", async () => { + const store = createWorkerSessionPlacementStore({ database, now: () => 1_000 }); + const { environmentId } = createDispatchEnvironmentFixtures(); + const active = seedActivePlacement(store, { environmentId, ownerEpoch: 2 }); + if (active.state !== "active") { + throw new Error("active placement fixture was not active"); + } + const owner = { + sessionId: active.sessionId, + environmentId: active.environmentId, + ownerEpoch: active.activeOwnerEpoch, + placementGeneration: active.generation, + }; + store.beginWorkspaceReconciliation(owner, { + version: 1, + temporaryNonce: "c".repeat(32), + baseManifestRef: active.workspaceBaseManifestRef, + currentManifestRef: `sha256:${"d".repeat(64)}`, + baseEntries: [], + appliedEntries: [], + baseTree: "f".repeat(40), + basePackSha256: createHash("sha256").update("").digest("hex"), + basePack: Buffer.alloc(0), + }); + const onCleanupError = vi.fn(); + + const resolveWorkspacePath = vi.fn(async () => { + throw new Error("workspace temporarily unavailable"); + }); + + for (let attempt = 0; attempt < 2; attempt += 1) { + await forceAbandonWorkerEnvironment({ + placements: store, + environmentId, + resolveWorkspacePath, + onCleanupError, + }); + } + + expect(store.get(REQUEST.sessionId)).toMatchObject({ state: "failed" }); + expect(store.listWorkspaceReconciliationOwners()).toEqual([owner]); + expect(resolveWorkspacePath).toHaveBeenCalledTimes(2); + expect(onCleanupError).toHaveBeenCalledWith( + expect.objectContaining({ message: "workspace temporarily unavailable" }), + ); + }); + + it("retains a current journal when loading it fails", async () => { + const store = createWorkerSessionPlacementStore({ database, now: () => 1_000 }); + const { environmentId } = createDispatchEnvironmentFixtures(); + const active = seedActivePlacement(store, { environmentId, ownerEpoch: 2 }); + if (active.state !== "active") { + throw new Error("active placement fixture was not active"); + } + const owner = { + sessionId: active.sessionId, + environmentId: active.environmentId, + ownerEpoch: active.activeOwnerEpoch, + placementGeneration: active.generation, + }; + store.beginWorkspaceReconciliation(owner, { + version: 1, + temporaryNonce: "d".repeat(32), + baseManifestRef: active.workspaceBaseManifestRef, + currentManifestRef: `sha256:${"e".repeat(64)}`, + baseEntries: [], + appliedEntries: [], + baseTree: "f".repeat(40), + basePackSha256: createHash("sha256").update("").digest("hex"), + basePack: Buffer.alloc(0), + }); + const onCleanupError = vi.fn(); + vi.spyOn(store, "loadWorkspaceReconciliation").mockImplementation(() => { + throw new Error("journal temporarily unreadable"); + }); + + const resolveWorkspacePath = vi.fn(async () => root); + for (let attempt = 0; attempt < 2; attempt += 1) { + await forceAbandonWorkerEnvironment({ + placements: store, + environmentId, + resolveWorkspacePath, + onCleanupError, + }); + } + + expect(store.get(REQUEST.sessionId)).toMatchObject({ state: "failed" }); + expect(store.listWorkspaceReconciliationOwners()).toEqual([owner]); + expect(resolveWorkspacePath).not.toHaveBeenCalled(); + expect(onCleanupError).toHaveBeenCalledWith( + expect.objectContaining({ message: "journal temporarily unreadable" }), + ); + }); }); diff --git a/src/gateway/worker-environments/placement-force-abandon.ts b/src/gateway/worker-environments/placement-force-abandon.ts index 36347ce4b320..e0e56b7f8c4b 100644 --- a/src/gateway/worker-environments/placement-force-abandon.ts +++ b/src/gateway/worker-environments/placement-force-abandon.ts @@ -14,16 +14,29 @@ async function tryResolveWorkspacePath( agentId: string; }) => Promise, placement: { sessionId: string; sessionKey: string; agentId: string }, + onCleanupError?: (error: unknown) => void, ): Promise { try { return await resolveWorkspacePath(placement); - } catch { + } catch (error) { // Forced teardown is the last-resort state owner. If the session/worktree is // already gone, skip local repair/ref cleanup and still release the claim. + reportCleanupError(onCleanupError, error); return undefined; } } +function reportCleanupError( + onCleanupError: ((error: unknown) => void) | undefined, + error: unknown, +): void { + try { + onCleanupError?.(error); + } catch { + // Cleanup reporting cannot overturn a committed forced abandonment. + } +} + export async function forceAbandonWorkerEnvironment(params: { placements: WorkerDispatchPlacementStore; environmentId: string; @@ -32,51 +45,66 @@ export async function forceAbandonWorkerEnvironment(params: { sessionKey: string; agentId: string; }) => Promise; + onCleanupError?: (error: unknown) => void; }): Promise { const { environmentId, placements } = params; const recoveryError = "Cloud worker result abandoned by forced operator teardown"; - for (const owner of placements.listWorkspaceReconciliationOwners()) { - if (owner.environmentId !== environmentId) { - continue; - } + const journalOwners = params.placements + .listWorkspaceReconciliationOwners() + .filter((owner) => owner.environmentId === environmentId); + const journalCleanups: Array<{ + owner: (typeof journalOwners)[number]; + placement: { sessionId: string; sessionKey: string; agentId: string }; + journal: NonNullable>; + }> = []; + const retainedJournalSessions = new Set(); + for (const owner of journalOwners) { const placement = placements.get(owner.sessionId); + const isCurrentOwner = + (placement?.state === "active" || placement?.state === "draining") && + placement.generation === owner.placementGeneration; + const isForceFailedOwner = + placement?.state === "failed" && + placement.recoveryError.startsWith(recoveryError) && + placement.generation > owner.placementGeneration; if ( - (placement?.state !== "active" && placement?.state !== "draining") || - placement.environmentId !== owner.environmentId || - placement.activeOwnerEpoch !== owner.ownerEpoch || - placement.generation !== owner.placementGeneration + placement && + (isCurrentOwner || isForceFailedOwner) && + placement.environmentId === owner.environmentId && + placement.activeOwnerEpoch === owner.ownerEpoch ) { - throw new Error(`Forced teardown found a stale workspace journal: ${owner.sessionId}`); - } - const journal = placements.loadWorkspaceReconciliation(owner); - if (journal) { - const root = await tryResolveWorkspacePath(params.resolveWorkspacePath, placement); - if (root) { - await recoverWorkerWorkspaceReconciliation({ root, journal }); + try { + const journal = placements.loadWorkspaceReconciliation( + owner, + isForceFailedOwner ? { allowFailedOwner: true } : undefined, + ); + if (journal) { + journalCleanups.push({ owner, placement, journal }); + } + } catch (error) { + reportCleanupError(params.onCleanupError, error); + retainedJournalSessions.add(owner.sessionId); } - placements.abortWorkspaceReconciliation(owner); } } + const stagedResultCleanups: Array<{ + placement: { sessionId: string; sessionKey: string; agentId: string }; + refs: string[]; + }> = []; for (const pending of placements.listPendingWorkspaceResults()) { if (pending.environmentId === environmentId) { const placement = placements.get(pending.sessionId); - if (!placement) { - if (pending.stagedResultRef) { - throw new Error( - `Forced teardown found a staged result without a placement: ${pending.sessionId}`, - ); - } - } else { - const root = await tryResolveWorkspacePath(params.resolveWorkspacePath, placement); - if (root) { - const finalRef = pending.stagedResultRef ?? workerWorkspaceResultRef(pending.claimId); - const refs = [finalRef, preparedWorkerWorkspaceResultRef(finalRef)]; - for (const stagedResultRef of refs) { - if (await hasWorkerWorkspaceResultRef({ root, stagedResultRef })) { - await deleteStagedWorkerWorkspaceResult({ root, stagedResultRef }); - } - } - } + if ( + (placement?.state === "active" || placement?.state === "draining") && + placement.environmentId === pending.environmentId && + placement.activeOwnerEpoch === pending.ownerEpoch && + placement.generation === pending.placementGeneration + ) { + const finalRef = pending.stagedResultRef ?? workerWorkspaceResultRef(pending.claimId); + stagedResultCleanups.push({ + placement, + refs: [finalRef, preparedWorkerWorkspaceResultRef(finalRef)], + }); } placements.abandonWorkspaceResult(pending); } @@ -110,4 +138,46 @@ export async function forceAbandonWorkerEnvironment(params: { }); } } + + // The durable fence is now closed. Filesystem rollback and ref cleanup are + // useful hygiene, but a changed or missing workspace must not revive it. + for (const cleanup of journalCleanups) { + if (cleanup.journal.appliedManifestRef) { + continue; + } + try { + const root = await params.resolveWorkspacePath(cleanup.placement); + await recoverWorkerWorkspaceReconciliation({ root, journal: cleanup.journal }); + } catch (error) { + reportCleanupError(params.onCleanupError, error); + retainedJournalSessions.add(cleanup.owner.sessionId); + } + } + // Placement failure is durable before journal removal. A crash during the + // best-effort rollback therefore leaves a fenced placement and retriable journal. + for (const owner of journalOwners) { + if (retainedJournalSessions.has(owner.sessionId)) { + continue; + } + placements.abortWorkspaceReconciliation(owner, { force: true }); + } + for (const cleanup of stagedResultCleanups) { + try { + const root = await tryResolveWorkspacePath( + params.resolveWorkspacePath, + cleanup.placement, + params.onCleanupError, + ); + if (!root) { + continue; + } + for (const stagedResultRef of cleanup.refs) { + if (await hasWorkerWorkspaceResultRef({ root, stagedResultRef })) { + await deleteStagedWorkerWorkspaceResult({ root, stagedResultRef }); + } + } + } catch (error) { + reportCleanupError(params.onCleanupError, error); + } + } } diff --git a/src/gateway/worker-environments/placement-force-destroy.test.ts b/src/gateway/worker-environments/placement-force-destroy.test.ts new file mode 100644 index 000000000000..d0131999dc23 --- /dev/null +++ b/src/gateway/worker-environments/placement-force-destroy.test.ts @@ -0,0 +1,164 @@ +import { createHash } from "node:crypto"; +import fs from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { createDeferred } from "../../shared/deferred.js"; +import { + closeOpenClawStateDatabaseForTest, + openOpenClawStateDatabase, + type OpenClawStateDatabase, +} from "../../state/openclaw-state-db.js"; +import { type PlacementStore, REQUEST } from "./placement-dispatch-test-fixtures.js"; +import { createHarness } from "./placement-dispatch-test-harness.js"; +import { createWorkerSessionPlacementStore } from "./placement-store.js"; +import { createWorkerWorkspaceOperationCoordinator } from "./workspace-operation-coordinator.js"; + +describe("forced worker environment destruction", () => { + let root: string; + let database: OpenClawStateDatabase; + let placementStore: PlacementStore; + + beforeEach(async () => { + root = await fs.mkdtemp(path.join(await fs.realpath(os.tmpdir()), "openclaw-force-destroy-")); + database = openOpenClawStateDatabase({ env: { OPENCLAW_STATE_DIR: root } }); + placementStore = createWorkerSessionPlacementStore({ database, now: () => 1_000 }); + }); + + afterEach(async () => { + closeOpenClawStateDatabaseForTest(); + await fs.rm(root, { recursive: true, force: true }); + }); + + it("serializes with workspace work and abandons an applied result fence", async () => { + const workspaceOperations = createWorkerWorkspaceOperationCoordinator(); + const harness = createHarness(placementStore, { workspaceOperations, workspacePath: root }); + await harness.environments.attachSession({ + environmentId: harness.ready.environmentId, + ownerEpoch: harness.ready.ownerEpoch, + sessionId: REQUEST.sessionId, + }); + const active = harness.placements.seedActive(harness.attached.ownerEpoch); + if (active.state !== "active") { + throw new Error("active placement fixture was not active"); + } + const claim = placementStore.claimTurn({ + ...REQUEST, + claimId: "force-destroy-claim", + runId: "force-destroy-run", + owner: { + kind: "worker", + environmentId: active.environmentId, + ownerEpoch: active.activeOwnerEpoch, + }, + }); + placementStore.markWorkspaceResultPending(claim); + const owner = { + sessionId: active.sessionId, + environmentId: active.environmentId, + ownerEpoch: active.activeOwnerEpoch, + placementGeneration: active.generation, + }; + const appliedManifestRef = harness.reconciledManifestRef; + placementStore.beginWorkspaceReconciliation(owner, { + version: 1, + temporaryNonce: "a".repeat(32), + baseManifestRef: active.workspaceBaseManifestRef, + currentManifestRef: appliedManifestRef, + baseEntries: [], + appliedEntries: [], + baseTree: "f".repeat(40), + basePackSha256: createHash("sha256").update("").digest("hex"), + basePack: Buffer.alloc(0), + }); + placementStore.updateWorkspaceBaseManifest({ claim, manifestRef: appliedManifestRef }); + expect(placementStore.loadWorkspaceReconciliation(owner)?.appliedManifestRef).toBe( + appliedManifestRef, + ); + + const releaseWorkspaceOperation = createDeferred(); + const workspaceOperation = workspaceOperations.run(active.environmentId, async () => { + await releaseWorkspaceOperation.promise; + }); + const forceDestroy = harness.service.forceDestroyEnvironment(active.environmentId); + await Promise.resolve(); + expect(harness.environments.destroy).not.toHaveBeenCalled(); + + releaseWorkspaceOperation.resolve(); + await expect(Promise.all([workspaceOperation, forceDestroy])).resolves.toEqual([ + undefined, + expect.objectContaining({ state: "destroyed" }), + ]); + expect(placementStore.get(REQUEST.sessionId)).toMatchObject({ + state: "failed", + turnClaim: null, + recoveryError: "Cloud worker result abandoned by forced operator teardown", + }); + expect(placementStore.listPendingWorkspaceResults()).toEqual([]); + expect(placementStore.listWorkspaceReconciliationOwners()).toEqual([]); + }); + + it.each([ + { failure: "tunnel stop", state: "draining" as const }, + { failure: "provider stop", state: "destroying" as const }, + ])("stays successful after $failure failure", async ({ state }) => { + const harness = createHarness(placementStore, { + destroyFails: true, + destroyFailureState: state, + workspacePath: root, + }); + harness.placements.seedActive(harness.attached.ownerEpoch); + const onCleanupError = vi.fn(); + + await expect( + harness.service.forceDestroyEnvironment(harness.ready.environmentId, onCleanupError), + ).resolves.toMatchObject({ state }); + + expect(harness.placements.current()).toMatchObject({ + state: "failed", + recoveryError: "Cloud worker result abandoned by forced operator teardown", + }); + expect(onCleanupError).toHaveBeenCalledWith( + expect.objectContaining({ message: "destroy pending" }), + ); + }); + + it("retries remote teardown when a failed rollback journal remains", async () => { + const harness = createHarness(placementStore, { + destroyFails: true, + destroyFailureState: "destroying", + failAt: "workspace", + }); + const active = harness.placements.seedActive(harness.attached.ownerEpoch); + if (active.state !== "active") { + throw new Error("active placement fixture was not active"); + } + const owner = { + sessionId: active.sessionId, + environmentId: active.environmentId, + ownerEpoch: active.activeOwnerEpoch, + placementGeneration: active.generation, + }; + placementStore.beginWorkspaceReconciliation(owner, { + version: 1, + temporaryNonce: "e".repeat(32), + baseManifestRef: active.workspaceBaseManifestRef, + currentManifestRef: `sha256:${"e".repeat(64)}`, + baseEntries: [], + appliedEntries: [], + baseTree: "f".repeat(40), + basePackSha256: createHash("sha256").update("").digest("hex"), + basePack: Buffer.alloc(0), + }); + + await expect( + harness.service.forceDestroyEnvironment(active.environmentId), + ).resolves.toMatchObject({ state: "destroying" }); + expect(placementStore.listWorkspaceReconciliationOwners()).toEqual([owner]); + vi.mocked(harness.environments.destroy).mockClear(); + + await harness.service.reconcileActive(active.environmentId); + + expect(harness.environments.destroy).toHaveBeenCalledExactlyOnceWith(active.environmentId); + }); +}); diff --git a/src/gateway/worker-environments/placement-workspace-journal.ts b/src/gateway/worker-environments/placement-workspace-journal.ts index d60400770af2..2882c48d794a 100644 --- a/src/gateway/worker-environments/placement-workspace-journal.ts +++ b/src/gateway/worker-environments/placement-workspace-journal.ts @@ -24,13 +24,25 @@ type WorkerWorkspaceJournalOwner = { placementGeneration: number; }; -function assertJournalOwner(db: DatabaseSync, owner: WorkerWorkspaceJournalOwner) { +function assertJournalOwner( + db: DatabaseSync, + owner: WorkerWorkspaceJournalOwner, + options: { allowFailedOwner?: boolean } = {}, +) { const placement = getRequired(db, owner.sessionId); + const isCurrentOwner = + (placement.state === "active" || placement.state === "draining") && + placement.generation === owner.placementGeneration; + // Forced teardown advances the exact owner to failed before best-effort + // rollback. Admit that state without weakening the manifest checks below. + const isAllowedFailedOwner = + options.allowFailedOwner === true && + placement.state === "failed" && + placement.generation > owner.placementGeneration; if ( - (placement.state !== "active" && placement.state !== "draining") || + (!isCurrentOwner && !isAllowedFailedOwner) || placement.environmentId !== owner.environmentId || - placement.activeOwnerEpoch !== owner.ownerEpoch || - placement.generation !== owner.placementGeneration + placement.activeOwnerEpoch !== owner.ownerEpoch ) { throw new Error(`Cannot reconcile stale worker workspace for session ${owner.sessionId}`); } @@ -79,9 +91,10 @@ export function createPlacementWorkspaceJournalOps(runtime: PlacementStoreRuntim loadWorkspaceReconciliation( owner: WorkerWorkspaceJournalOwner, + options: { allowFailedOwner?: boolean } = {}, ): WorkerWorkspaceReconciliationJournal | undefined { const db = read(); - const placement = assertJournalOwner(db, owner); + const placement = assertJournalOwner(db, owner, options); const row = executeSqliteQuerySync( db, query(db) @@ -154,10 +167,30 @@ export function createPlacementWorkspaceJournalOps(runtime: PlacementStoreRuntim }); }, - abortWorkspaceReconciliation(owner: WorkerWorkspaceJournalOwner): void { + abortWorkspaceReconciliation( + owner: WorkerWorkspaceJournalOwner, + options: { force?: boolean } = {}, + ): void { write((db) => { - assertJournalOwner(db, owner); - clearWorkerWorkspaceReconciliation(db, owner.sessionId); + if (!options.force) { + assertJournalOwner(db, owner); + clearWorkerWorkspaceReconciliation(db, owner.sessionId); + return; + } + // Forced teardown owns this exact durable journal even when placement + // state advanced after a failed recovery sweep. + const result = executeSqliteQuerySync( + db, + query(db) + .deleteFrom("worker_workspace_reconciliations") + .where("session_id", "=", owner.sessionId) + .where("environment_id", "=", owner.environmentId) + .where("owner_epoch", "=", owner.ownerEpoch) + .where("placement_generation", "=", owner.placementGeneration), + ); + if (result.numAffectedRows !== 1n) { + throw new Error(`Worker workspace journal changed for ${owner.sessionId}`); + } }); }, }; diff --git a/src/gateway/worker-environments/service-contract.ts b/src/gateway/worker-environments/service-contract.ts index 242fc5c4e950..2e4fac7f8710 100644 --- a/src/gateway/worker-environments/service-contract.ts +++ b/src/gateway/worker-environments/service-contract.ts @@ -52,6 +52,9 @@ export type WorkerPlacementDispatchContract = { reclaim?( request: WorkerPlacementReclaimRequest, ): Promise>; - forceDestroyEnvironment?(environmentId: string): Promise; + forceDestroyEnvironment?( + environmentId: string, + onCleanupError?: (error: unknown) => void, + ): Promise; reconcileActive?(environmentId?: string): Promise; };