diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 3835a395ec..08864af759 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -2,6 +2,8 @@ ## [Unreleased] +- Fixed supervisor recovery replacing live, load-slow session workers and interrupting their in-flight work. + ## [0.7.0] - 2026-08-05 ### Breaking Changes diff --git a/packages/coding-agent/src/modes/daemon/daemon-supervisor.ts b/packages/coding-agent/src/modes/daemon/daemon-supervisor.ts index de7210c79d..f7ddb1d4b9 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-supervisor.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-supervisor.ts @@ -2409,6 +2409,21 @@ export class DaemonSupervisor { return; } this.log(`Could not adopt worker ${worker.descriptor.workerId}: ${String(error)}`); + const processAlive = isProcessAlive(worker.descriptor.pid); + const observedProcessStartId = processAlive ? getProcessStartId(worker.descriptor.pid) : undefined; + if ( + worker.descriptor.processStartId !== undefined && + processAlive && + (observedProcessStartId === undefined || observedProcessStartId === worker.descriptor.processStartId) + ) { + worker.descriptor.lifecycle = "recovering"; + worker.descriptor.lastError = error instanceof Error ? error.message : String(error); + this.persistWorker(worker); + void this.recoverWorker(worker).catch((recoveryError) => + this.log(`Could not recover worker ${worker.descriptor.workerId}: ${String(recoveryError)}`), + ); + return; + } await this.recoverWorker(worker); } } @@ -2717,8 +2732,10 @@ export class DaemonSupervisor { return worker.recovery; } worker.recovery = (async () => { - for (const [retryIndex, retryDelay] of WORKER_RETRY_DELAYS_MS.entries()) { + let keepRetryingLiveWorker = false; + for (const retryDelay of WORKER_RETRY_DELAYS_MS) { await delay(retryDelay); + keepRetryingLiveWorker = false; if (this.isWorkerRecoveryCancelled(worker)) { return; } @@ -2756,22 +2773,25 @@ export class DaemonSupervisor { await this.assertRecoveryAllowed(); worker.client?.close(); worker.client = undefined; - if (retryIndex < WORKER_RETRY_DELAYS_MS.length - 1) { - throw error; - } + // A verified live worker can be load-slow rather than dead. Keep probing + // instead of relaunching it and dropping its in-flight operations. + keepRetryingLiveWorker = + worker.descriptor.processStartId !== undefined && + observedProcessStartId === worker.descriptor.processStartId; + throw error; } } if ( processAlive && (worker.descriptor.processStartId === undefined || observedProcessStartId === undefined) ) { + keepRetryingLiveWorker = + worker.descriptor.processStartId !== undefined && observedProcessStartId === undefined; throw new Error( `Cannot safely replace live session worker ${worker.descriptor.workerId} without a verified process identity`, ); } - const safeToKillWorkerProcess = - processAlive && processIdentityMatches && worker.descriptor.processStartId !== undefined; - await this.recoverUncertainWorkerOperations(worker, safeToKillWorkerProcess); + await this.recoverUncertainWorkerOperations(worker, false); if (this.isWorkerRecoveryCancelled(worker)) { return; } @@ -2794,6 +2814,15 @@ export class DaemonSupervisor { this.persistWorker(worker); } } + if (keepRetryingLiveWorker) { + worker.descriptor.lifecycle = "recovering"; + this.persistWorker(worker); + this.deferWorkerRecovery( + worker, + new Error(worker.descriptor.lastError ?? "Live session worker did not answer recovery probes"), + ); + return; + } try { await this.assertRecoveryAllowed(); } catch { diff --git a/packages/coding-agent/test/daemon-supervisor-monitor.test.ts b/packages/coding-agent/test/daemon-supervisor-monitor.test.ts index c05afc018b..dccabf5185 100644 --- a/packages/coding-agent/test/daemon-supervisor-monitor.test.ts +++ b/packages/coding-agent/test/daemon-supervisor-monitor.test.ts @@ -6,6 +6,7 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import type { AgentMessage } from "@earendil-works/pi-agent-core"; import { afterEach, describe, expect, it, vi } from "vitest"; +import { getProcessStartId } from "../src/core/session-lease.js"; import type { DaemonSocketClient } from "../src/modes/daemon/active-session-state.js"; import { CommandRecoveryJournal } from "../src/modes/daemon/command-recovery-journal.js"; import { DaemonCatalogClient } from "../src/modes/daemon/daemon-catalog-process.js"; @@ -1401,6 +1402,140 @@ describe("daemon worker supervisor monitoring", () => { expect(worker.descriptor.lifecycle).toBe("failed"); }); + it("continues startup after a verified live worker fails its initial adoption probe", async () => { + type AdoptionWorker = { + descriptor: { + workerId: string; + pid: number; + processStartId?: string; + rootActiveSessionId: string; + lifecycle?: string; + lastError?: string; + }; + }; + type AdoptionHarness = { + connectWorker: ReturnType; + subscribeWorker: ReturnType; + refreshWorkerSummaries: ReturnType; + recoverWorker: ReturnType; + persistWorker: ReturnType; + log: ReturnType; + assertRecoveryAllowed: ReturnType; + adoptOrRecoverWorker(worker: AdoptionWorker): Promise; + }; + const worker: AdoptionWorker = { + descriptor: { + workerId: "worker-live-unreachable", + pid: process.pid, + processStartId: getProcessStartId(process.pid), + rootActiveSessionId: "active-1", + }, + }; + const pendingRecovery = new Promise(() => {}); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + connectWorker: vi.fn(async () => {}), + subscribeWorker: vi.fn(async () => {}), + refreshWorkerSummaries: vi.fn(async () => { + throw new Error("Timed out waiting for daemon worker response to list"); + }), + recoverWorker: vi.fn(() => pendingRecovery), + persistWorker: vi.fn(), + log: vi.fn(), + assertRecoveryAllowed: vi.fn(async () => {}), + }) as AdoptionHarness; + + await expect(supervisor.adoptOrRecoverWorker(worker)).resolves.toBeUndefined(); + + expect(supervisor.recoverWorker).toHaveBeenCalledWith(worker); + expect(worker.descriptor.lifecycle).toBe("recovering"); + }); + + it.each([ + { name: "after repeated probe timeouts", identityUnavailable: false, expectedConnections: 4 }, + { name: "when its identity is temporarily unavailable", identityUnavailable: true, expectedConnections: 1 }, + ])("keeps retrying a verified live worker $name", async ({ identityUnavailable, expectedConnections }) => { + vi.useFakeTimers(); + type RecoveryWorker = { + descriptor: { + workerId: string; + pid: number; + processStartId?: string; + rootActiveSessionId: string; + createCommand: { type: "create" }; + lifecycle?: string; + consecutiveFailures: number; + lastFailureAt?: string; + lastError?: string; + }; + intentionalStop: boolean; + stopRevision: number; + recovery?: Promise; + client?: { close(): void }; + }; + type RecoveryHarness = { + workers: Map; + shuttingDown: boolean; + connectWorker: ReturnType; + subscribeWorker: ReturnType; + refreshWorkerSummaries: ReturnType; + recoverUncertainWorkerOperations: ReturnType; + launchWorker: ReturnType; + persistWorker: ReturnType; + syncAgentPeers: ReturnType; + broadcastHeartbeatsChanged: ReturnType; + log: ReturnType; + assertRecoveryAllowed: ReturnType; + recoverWorker(worker: RecoveryWorker): Promise; + }; + const worker: RecoveryWorker = { + descriptor: { + workerId: "worker-live-unreachable", + pid: process.pid, + processStartId: getProcessStartId(process.pid), + rootActiveSessionId: "active-1", + createCommand: { type: "create" }, + consecutiveFailures: 0, + }, + intentionalStop: false, + stopRevision: 0, + }; + workerLaunchTestState.forceMissingProcessStartId = identityUnavailable; + const timeout = new Error("Timed out connecting to daemon session worker"); + let remainingTimeouts = identityUnavailable ? 0 : 3; + const connectWorker = vi.fn(async () => { + if (remainingTimeouts-- > 0) { + throw timeout; + } + }); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + workers: new Map([[worker.descriptor.workerId, worker]]), + shuttingDown: false, + connectWorker, + subscribeWorker: vi.fn(async () => {}), + refreshWorkerSummaries: vi.fn(async () => {}), + recoverUncertainWorkerOperations: vi.fn(async () => {}), + launchWorker: vi.fn(async () => worker), + persistWorker: vi.fn(() => { + if (identityUnavailable && worker.descriptor.consecutiveFailures === 3) { + workerLaunchTestState.forceMissingProcessStartId = false; + } + }), + syncAgentPeers: vi.fn(async () => {}), + broadcastHeartbeatsChanged: vi.fn(), + log: vi.fn(), + assertRecoveryAllowed: vi.fn(async () => {}), + }) as RecoveryHarness; + + const recovery = supervisor.recoverWorker(worker); + await vi.runAllTimersAsync(); + await recovery; + + expect(supervisor.connectWorker).toHaveBeenCalledTimes(expectedConnections); + expect(supervisor.recoverUncertainWorkerOperations).not.toHaveBeenCalled(); + expect(supervisor.launchWorker).not.toHaveBeenCalled(); + expect(worker.descriptor.lifecycle).toBe("ready"); + }); + it("keeps a recovered worker ready when peer synchronization fails", async () => { vi.useFakeTimers(); type RecoveryWorker = {