diff --git a/server/src/__tests__/heartbeat-process-recovery.test.ts b/server/src/__tests__/heartbeat-process-recovery.test.ts index df7ffae7d9..af0c89a8bd 100644 --- a/server/src/__tests__/heartbeat-process-recovery.test.ts +++ b/server/src/__tests__/heartbeat-process-recovery.test.ts @@ -575,6 +575,72 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => { return { environmentId, leaseId }; } + it("does not reap active adapter executions started by another heartbeat service instance", async () => { + let releaseAdapter: (() => void) | null = null; + const adapterStarted = new Promise((resolve) => { + mockAdapterExecute.mockImplementationOnce(async () => { + resolve(); + await new Promise((release) => { + releaseAdapter = release; + }); + return { + exitCode: 0, + signal: null, + timedOut: false, + errorMessage: null, + summary: "Remote run completed.", + provider: "test", + model: "test-model", + }; + }); + }); + + const { runId, wakeupRequestId } = await seedRunFixture({ + adapterType: "openclaw_gateway", + agentStatus: "idle", + runStatus: "queued", + processPid: null, + processGroupId: null, + includeIssue: false, + }); + const executorHeartbeat = heartbeatService(db); + const reaperHeartbeat = heartbeatService(db); + + await executorHeartbeat.resumeQueuedRuns(); + await Promise.race([ + adapterStarted, + new Promise((_, reject) => { + setTimeout(() => reject(new Error("Timed out waiting for adapter execution to start")), 3_000); + }), + ]); + + await db + .update(heartbeatRuns) + .set({ + updatedAt: new Date("2026-03-19T00:00:00.000Z"), + }) + .where(eq(heartbeatRuns.id, runId)); + + const result = await reaperHeartbeat.reapOrphanedRuns({ staleThresholdMs: 1 }); + expect(result).toEqual({ reaped: 0, runIds: [] }); + + const activeRun = await reaperHeartbeat.getRun(runId); + expect(activeRun?.status).toBe("running"); + expect(activeRun?.errorCode).toBeNull(); + + const wakeup = await db + .select() + .from(agentWakeupRequests) + .where(eq(agentWakeupRequests.id, wakeupRequestId)) + .then((rows) => rows[0] ?? null); + expect(wakeup?.status).toBe("claimed"); + + if (!releaseAdapter) throw new Error("Adapter release handle was not captured"); + releaseAdapter(); + const settledRun = await waitForRunToSettle(executorHeartbeat, runId, 5_000); + expect(settledRun?.status).toBe("succeeded"); + }); + async function seedStrandedIssueFixture(input: { status: "todo" | "in_progress"; runStatus: "failed" | "timed_out" | "cancelled" | "succeeded"; diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index c530f09513..739b4dc47e 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -409,6 +409,9 @@ const SESSIONED_LOCAL_ADAPTERS = new Set([ "opencode_local", "pi_local", ]); +// Routes and the scheduler construct separate heartbeatService instances, but +// they must agree on in-process adapter executions when reaping stale runs. +const activeRunExecutions = new Set(); const INLINE_BASE64_IMAGE_DATA_RE = /("type":"image","source":\{"type":"base64","data":")([A-Za-z0-9+/=]{1024,})(")/g; type RuntimeConfigSecretResolver = Pick< @@ -3556,7 +3559,6 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {}) environmentRuntime, }); const workspaceOperationsSvc = workspaceOperationService(db); - const activeRunExecutions = new Set(); const liveRunExecutions = { has(id: string) { return runningProcesses.has(id) || activeRunExecutions.has(id);