diff --git a/packages/coding-agent/.changes/daemon-recovery-hardening.md b/packages/coding-agent/.changes/daemon-recovery-hardening.md new file mode 100644 index 0000000000..82751a6f92 --- /dev/null +++ b/packages/coding-agent/.changes/daemon-recovery-hardening.md @@ -0,0 +1,2 @@ +- Fixed daemon startup and recovery to preserve slow live processes and fail closed after socket lock loss. +- Recovery never signals a live worker process it cannot verify as its own: a persistently failing live worker parks as failed with its process left running (reclaimed automatically by the next fresh create once its identity is verified or it exits). The one deliberate exception is replacing an authenticated pre-roster worker during adoption. A live worker that stays silent through ten probe rounds (~2.5 minutes) also parks as failed instead of probing forever. diff --git a/packages/coding-agent/src/cli/daemon-launch.ts b/packages/coding-agent/src/cli/daemon-launch.ts index 819f0503a4..0d57425064 100644 --- a/packages/coding-agent/src/cli/daemon-launch.ts +++ b/packages/coding-agent/src/cli/daemon-launch.ts @@ -64,10 +64,19 @@ async function canConnectToDaemon(socketPath: string, timeoutMs: number): Promis type DaemonVersionProbe = | { status: "absent" } | { status: "current"; hello: DaemonHello } - | { status: "stale"; hello?: DaemonHello }; + | { status: "stale"; hello: DaemonHello } + | { status: "unresponsive" }; + +function isCurrentDaemonHello(hello: DaemonHello): boolean { + return ( + hello.protocol.version === DAEMON_PROTOCOL_VERSION && + hello.schemaId === DAEMON_SCHEMA_ID && + hello.appVersion === VERSION + ); +} /** Connect to a running daemon and check whether it matches this client's protocol and app version. */ -export async function probeDaemonVersion(socketPath: string): Promise { +export async function probeDaemonVersion(socketPath: string, helloTimeoutMs = 2000): Promise { let client: DaemonClient | undefined; for (const timeoutMs of [250, 2000]) { const candidate = new DaemonClient(socketPath); @@ -83,11 +92,8 @@ export async function probeDaemonVersion(socketPath: string): Promise { +type StaleDaemonDisposition = "current" | "stopped" | "busy"; + +async function shutdownStaleDaemonIfNotBusy(socketPath: string): Promise { const client = new DaemonClient(socketPath); - let connected = false; - let hasBusySessions = false; - let loadedSessionCount = 0; try { await client.connect(1000); - connected = true; - try { - const result = await queryActiveDaemonSessions(client, { includeClientOwned: true }); - loadedSessionCount = result.sessions.length; - hasBusySessions = - result.busyClientOwnedSessionCount !== 0 || result.sessions.some((summary) => isSessionBusy(summary)); - } catch { - // Couldn't confirm idleness: treat as busy rather than risk interrupting work. - hasBusySessions = true; - } } catch { - // Couldn't reach it to inspect; don't send a blind shutdown, just verify below. - } finally { client.close(); + return (await waitForDaemonGone(socketPath)) ? "stopped" : "busy"; } - if (!connected) { - return waitForDaemonGone(socketPath); + let loadedSessionCount = 0; + let hasBusySessions = true; + try { + const result = await queryActiveDaemonSessions(client, { includeClientOwned: true }); + loadedSessionCount = result.sessions.length; + hasBusySessions = + result.busyClientOwnedSessionCount !== 0 || result.sessions.some((summary) => isSessionBusy(summary)); + } catch { + // An unresponsive daemon is not safe to replace. + } + + const hello = client.hello; + if (hello && isCurrentDaemonHello(hello)) { + client.close(); + logDaemonLaunch(`daemon on ${socketPath} finished starting while staleness was being checked; reusing it`); + return "current"; } if (hasBusySessions) { + client.close(); logDaemonLaunch(`refusing to replace stale daemon on ${socketPath}: busy session(s) present`); - return false; + return "busy"; } logDaemonLaunch( `replacing stale daemon on ${socketPath} (idle): ${loadedSessionCount} loaded session(s) will reload`, ); - return shutdownDaemonAndWait(socketPath); + return (await shutdownConnectedDaemonAndWait(client, socketPath, 5000, hello)) ? "stopped" : "busy"; } async function ensureDaemonRunning(socketPath: string, spawnCwd?: string): Promise { - const probe = await probeDaemonVersion(socketPath); + const probeStartedAt = Date.now(); + let probe = await probeDaemonVersion(socketPath); + if (probe.status === "unresponsive") { + const remainingStartupMs = Math.max(1, DAEMON_STARTUP_TIMEOUT_MS - (Date.now() - probeStartedAt)); + probe = await probeDaemonVersion(socketPath, remainingStartupMs); + } if (probe.status === "current") { return; } + if (probe.status === "unresponsive") { + throw new Error( + `Prime Agent daemon on ${socketPath} accepted connections but did not finish startup within ${DAEMON_STARTUP_TIMEOUT_MS / 1000} seconds. ` + + `It was left running to avoid interrupting active work. + +Run: +${formatCurrentCliCommand(["shutdown", "--force"])} + +Then retry the original command.`, + ); + } if (probe.status === "stale") { - const stopped = await shutdownStaleDaemonIfNotBusy(socketPath); - if (!stopped) throw new StaleDaemonError(socketPath, probe.hello); + const disposition = await shutdownStaleDaemonIfNotBusy(socketPath); + if (disposition === "current") return; + if (disposition === "busy") throw new StaleDaemonError(socketPath, probe.hello); } const entrypoint = process.argv[1]; diff --git a/packages/coding-agent/src/cli/daemon-update-restart.ts b/packages/coding-agent/src/cli/daemon-update-restart.ts index 6afcae9fdb..f436f2f6d1 100644 --- a/packages/coding-agent/src/cli/daemon-update-restart.ts +++ b/packages/coding-agent/src/cli/daemon-update-restart.ts @@ -324,11 +324,18 @@ function coordinatorRecordPath(registryDir: string, socketPath: string): string async function withCoordinatorRegistryGuard(registryDir: string, action: () => T | Promise): Promise { mkdirSync(registryDir, { recursive: true, mode: 0o700 }); + let compromisedError: Error | undefined; + const assertGuardHeld = () => { + if (compromisedError) throw new Error(`Coordinator registry guard was compromised: ${compromisedError.message}`); + }; const release = await lockfile.lock(registryDir, { realpath: false, lockfilePath: resolve(registryDir, ".guard"), stale: COORDINATOR_REGISTRY_LOCK_STALE_MS, update: COORDINATOR_REGISTRY_LOCK_UPDATE_MS, + onCompromised: (error) => { + compromisedError ??= error; + }, retries: { retries: COORDINATOR_REGISTRY_LOCK_RETRIES, factor: 1, @@ -337,9 +344,13 @@ async function withCoordinatorRegistryGuard(registryDir: string, action: () = }, }); try { - return await action(); + assertGuardHeld(); + const result = await action(); + assertGuardHeld(); + return result; } finally { - await release(); + if (compromisedError) await release().catch(() => undefined); + else await release(); } } diff --git a/packages/coding-agent/src/core/auth-storage.ts b/packages/coding-agent/src/core/auth-storage.ts index a72a5d7c1c..389bc64f85 100644 --- a/packages/coding-agent/src/core/auth-storage.ts +++ b/packages/coding-agent/src/core/auth-storage.ts @@ -126,11 +126,23 @@ export class FileAuthStorageBackend implements AuthStorageBackend { const maxAttempts = 10; const delayMs = 20; let lastError: unknown; + let compromisedError: Error | undefined; for (let attempt = 1; attempt <= maxAttempts; attempt++) { try { - return lockfile.lockSync(path, { realpath: false }); + const release = lockfile.lockSync(path, { + realpath: false, + onCompromised: (error) => { + compromisedError ??= error; + }, + }); + if (compromisedError) { + release(); + throw compromisedError; + } + return release; } catch (error) { + if (compromisedError) throw compromisedError; const code = typeof error === "object" && error !== null && "code" in error ? String((error as { code?: unknown }).code) @@ -211,11 +223,8 @@ export class FileAuthStorageBackend implements AuthStorageBackend { return result; } finally { if (release) { - try { - await release(); - } catch { - // Ignore unlock errors when lock is compromised. - } + if (lockCompromised) await release().catch(() => undefined); + else await release(); } } } diff --git a/packages/coding-agent/src/core/cron-jobs.ts b/packages/coding-agent/src/core/cron-jobs.ts index 3c26ed275c..c9dede0476 100644 --- a/packages/coding-agent/src/core/cron-jobs.ts +++ b/packages/coding-agent/src/core/cron-jobs.ts @@ -1499,15 +1499,26 @@ function withCronJobsStateLocks(paths: readonly string[], action: () => T): T for (const path of [...new Set(paths)].sort()) { mkdirSync(dirname(path), { recursive: true, mode: 0o700 }); let release: (() => void) | undefined; + let lockCompromised = false; for (let attempt = 0; attempt < 100; attempt++) { try { release = lockSync(path, { realpath: false, lockfilePath: `${path}.lock`, stale: 30_000, + onCompromised: () => { + lockCompromised = true; + }, }); + if (lockCompromised) { + release(); + throw new Error(`Cron jobs lock compromised: ${path}`); + } break; } catch (error) { + if (lockCompromised) { + throw new Error(`Cron jobs lock compromised: ${path}`); + } if ((error as NodeJS.ErrnoException).code !== "ELOCKED" || attempt === 99) { throw error; } diff --git a/packages/coding-agent/src/core/session-lease.ts b/packages/coding-agent/src/core/session-lease.ts index 6c4e2975cf..7f89449c67 100644 --- a/packages/coding-agent/src/core/session-lease.ts +++ b/packages/coding-agent/src/core/session-lease.ts @@ -187,12 +187,16 @@ function isLeaseOwnerAlive(owner: SessionLeaseOwner): boolean { function withLeaseGuard(directory: string, action: () => T): T { let release: (() => void) | undefined; + let guardCompromised = false; for (let attempt = 0; attempt < 100; attempt++) { try { release = lockSync(directory, { realpath: false, lockfilePath: `${directory}.guard`, stale: 5000, + onCompromised: () => { + guardCompromised = true; + }, }); break; } catch (error) { @@ -208,10 +212,24 @@ function withLeaseGuard(directory: string, action: () => T): T { if (!release) { throw new Error(`Could not coordinate session lease: ${directory}`); } + const assertGuardHeld = () => { + if (guardCompromised) throw new Error(`Session lease guard was compromised: ${directory}`); + }; try { - return action(); + assertGuardHeld(); + const result = action(); + assertGuardHeld(); + return result; } finally { - release(); + if (guardCompromised) { + try { + release(); + } catch { + // The compromised guard no longer owns a lock that can be safely released. + } + } else { + release(); + } } } diff --git a/packages/coding-agent/src/core/settings-manager.ts b/packages/coding-agent/src/core/settings-manager.ts index f35e0da29d..e2b7ee8e77 100644 --- a/packages/coding-agent/src/core/settings-manager.ts +++ b/packages/coding-agent/src/core/settings-manager.ts @@ -236,11 +236,23 @@ export class FileSettingsStorage implements SettingsStorage { const maxAttempts = 10; const delayMs = 20; let lastError: unknown; + let compromisedError: Error | undefined; for (let attempt = 1; attempt <= maxAttempts; attempt++) { try { - return lockfile.lockSync(path, { realpath: false }); + const release = lockfile.lockSync(path, { + realpath: false, + onCompromised: (error) => { + compromisedError ??= error; + }, + }); + if (compromisedError) { + release(); + throw compromisedError; + } + return release; } catch (error) { + if (compromisedError) throw compromisedError; const code = typeof error === "object" && error !== null && "code" in error ? String((error as { code?: unknown }).code) diff --git a/packages/coding-agent/src/modes/daemon/daemon-catalog-entry.ts b/packages/coding-agent/src/modes/daemon/daemon-catalog-entry.ts new file mode 100644 index 0000000000..a8186211be --- /dev/null +++ b/packages/coding-agent/src/modes/daemon/daemon-catalog-entry.ts @@ -0,0 +1,7 @@ +import { runDaemonCatalogProcess } from "./daemon-catalog-process.js"; + +runDaemonCatalogProcess().catch((error: unknown) => { + const message = error instanceof Error ? error.message : String(error); + process.stderr.write(`Prime Agent daemon catalog failed: ${message}\n`); + process.exit(1); +}); diff --git a/packages/coding-agent/src/modes/daemon/daemon-catalog-process.ts b/packages/coding-agent/src/modes/daemon/daemon-catalog-process.ts index 9a69e1ebff..17453aaabb 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-catalog-process.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-catalog-process.ts @@ -1,11 +1,34 @@ import { type ChildProcess, spawn } from "node:child_process"; import { randomUUID } from "node:crypto"; +import { existsSync } from "node:fs"; +import { createRequire } from "node:module"; +import { join, sep } from "node:path"; +import { fileURLToPath } from "node:url"; import { createCliSubprocessEnv, createCliSubprocessLaunchSpec } from "../../cli/subprocess-launch.js"; +import { getPackageDir, isBunBinary } from "../../config.js"; import type { DeleteSessionFileResult } from "../../core/session-file-actions.js"; import { deleteSessionFile } from "../../core/session-file-actions.js"; import { readSessionInfo, type SessionInfo, SessionManager } from "../../core/session-manager.js"; export const DAEMON_CATALOG_ROLE_ENV = "PRIME_AGENT_INTERNAL_DAEMON_CATALOG"; +const DAEMON_CATALOG_START_TIMEOUT_MS = 30_000; + +export function isDaemonCatalogSourcePath(modulePath: string, packageDir: string): boolean { + return modulePath.startsWith(`${join(packageDir, "src")}${sep}`); +} + +function resolveDaemonCatalogEntrypoint(): string { + const packageDir = getPackageDir(); + const sourceEntrypoint = join(packageDir, "src", "modes", "daemon", "daemon-catalog-entry.ts"); + const compiledEntrypoint = join(packageDir, "dist", "modes", "daemon", "daemon-catalog-entry.js"); + const runningFromSource = isDaemonCatalogSourcePath(fileURLToPath(import.meta.url), packageDir); + const candidates = runningFromSource + ? [sourceEntrypoint, compiledEntrypoint] + : [compiledEntrypoint, sourceEntrypoint]; + const entrypoint = candidates.find((candidate) => existsSync(candidate)); + if (entrypoint) return entrypoint; + throw new Error("Cannot locate the daemon catalog entrypoint"); +} interface SessionInfoWire extends Omit { created: string; @@ -319,10 +342,26 @@ export class DaemonCatalogClient { } private async spawnCatalog(): Promise { - const launch = createCliSubprocessLaunchSpec(["--version"]); - const child = spawn(launch.command, launch.args, { + let command: string; + let args: string[]; + let environment = createCliSubprocessEnv({ ...process.env, [DAEMON_CATALOG_ROLE_ENV]: "1" }); + if (isBunBinary) { + const launch = createCliSubprocessLaunchSpec(["--version"]); + command = launch.command; + args = launch.args; + } else { + const catalogEntry = resolveDaemonCatalogEntrypoint(); + const execArgs = catalogEntry.endsWith(".ts") + ? [...process.execArgv, "--import", createRequire(import.meta.url).resolve("tsx")] + : process.execArgv; + const launch = createCliSubprocessLaunchSpec([], undefined, execArgs, catalogEntry); + command = launch.command; + args = launch.args; + environment = createCliSubprocessEnv(environment, catalogEntry, execArgs); + } + const child = spawn(command, args, { cwd: process.cwd(), - env: createCliSubprocessEnv({ ...process.env, [DAEMON_CATALOG_ROLE_ENV]: "1" }), + env: environment, stdio: ["ignore", "ignore", "ignore", "ipc"], }); this.child = child; @@ -341,11 +380,12 @@ export class DaemonCatalogClient { } child.kill("SIGKILL"); rejectReady(error); - }, 5000); + }, DAEMON_CATALOG_START_TIMEOUT_MS); const cleanup = () => { clearTimeout(timeout); child.off("message", onMessage); child.off("error", onError); + child.off("exit", onExit); }; const onMessage = (value: unknown) => { if (isCatalogOutbound(value) && value.type === "ready") { @@ -357,8 +397,12 @@ export class DaemonCatalogClient { cleanup(); rejectReady(error); }; + const onExit = (code: number | null, signal: NodeJS.Signals | null) => { + onError(new Error(`Daemon catalog exited during startup (${signal ?? code ?? "unknown"})`)); + }; child.on("message", onMessage); child.once("error", onError); + child.once("exit", onExit); }); } diff --git a/packages/coding-agent/src/modes/daemon/daemon-socket.ts b/packages/coding-agent/src/modes/daemon/daemon-socket.ts index bed8a705ac..99a3a15105 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-socket.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-socket.ts @@ -13,21 +13,53 @@ const DAEMON_SOCKET_RELEASE_POLL_MS = 25; const DAEMON_SOCKET_LOCK_STALE_MS = 5000; const DAEMON_SOCKET_LOCK_UPDATE_MS = 1000; +type DaemonSocketCompromiseListener = (error: Error) => void; + export class DaemonSocketPathLease { private released = false; + private compromisedError?: Error; + private readonly compromiseListeners = new Set(); constructor( readonly socketPath: string, private readonly releaseLock: () => Promise, ) {} - async release(): Promise { - if (this.released) { - return; + get compromise(): Error | undefined { + return this.compromisedError; + } + + onCompromised(listener: DaemonSocketCompromiseListener): () => void { + if (this.compromisedError) { + this.notifyCompromiseListener(listener, this.compromisedError); + return () => {}; } + this.compromiseListeners.add(listener); + return () => this.compromiseListeners.delete(listener); + } + + recordCompromise(error: Error): void { + if (this.compromisedError) return; + this.compromisedError = error; + const listeners = [...this.compromiseListeners]; + this.compromiseListeners.clear(); + for (const listener of listeners) this.notifyCompromiseListener(listener, error); + } + + async release(): Promise { + if (this.released) return; this.released = true; + this.compromiseListeners.clear(); await this.releaseLock(); } + + private notifyCompromiseListener(listener: DaemonSocketCompromiseListener, error: Error): void { + try { + listener(error); + } catch { + // A lease callback must not rethrow from proper-lockfile's refresh callback. + } + } } export interface DaemonSocketIdentity { @@ -47,10 +79,16 @@ export async function acquireDaemonSocketPathLease(socketPath: string): Promise< if (process.platform === "win32") { return undefined; } + let lease: DaemonSocketPathLease | undefined; + let pendingCompromise: Error | undefined; const releaseLock = await lockfile.lock(socketPath, { realpath: false, stale: DAEMON_SOCKET_LOCK_STALE_MS, update: DAEMON_SOCKET_LOCK_UPDATE_MS, + onCompromised: (error) => { + if (lease) lease.recordCompromise(error); + else pendingCompromise = error; + }, retries: { retries: 600, factor: 1, @@ -58,7 +96,9 @@ export async function acquireDaemonSocketPathLease(socketPath: string): Promise< maxTimeout: DAEMON_SOCKET_RELEASE_POLL_MS, }, }); - return new DaemonSocketPathLease(socketPath, releaseLock); + lease = new DaemonSocketPathLease(socketPath, releaseLock); + if (pendingCompromise) lease.recordCompromise(pendingCompromise); + return lease; } export async function prepareDaemonSocketPath(socketPath: string, lease?: DaemonSocketPathLease): Promise { @@ -69,7 +109,8 @@ export async function prepareDaemonSocketPath(socketPath: string, lease?: Daemon } if (lease) { assertSocketLease(socketPath, lease); - await prepareUnixDaemonSocketPath(socketPath); + assertSocketLeaseHeld(socketPath, lease); + await prepareUnixDaemonSocketPath(socketPath, lease); return; } if (!existsSync(socketPath)) { @@ -80,13 +121,13 @@ export async function prepareDaemonSocketPath(socketPath: string, lease?: Daemon } const ownedLease = await acquireDaemonSocketPathLease(socketPath); try { - await prepareUnixDaemonSocketPath(socketPath); + await prepareUnixDaemonSocketPath(socketPath, ownedLease); } finally { await ownedLease?.release(); } } -async function prepareUnixDaemonSocketPath(socketPath: string): Promise { +async function prepareUnixDaemonSocketPath(socketPath: string, lease?: DaemonSocketPathLease): Promise { if (!existsSync(socketPath)) { return; } @@ -131,6 +172,7 @@ async function prepareUnixDaemonSocketPath(socketPath: string): Promise { } } + if (lease) assertSocketLeaseHeld(socketPath, lease); unlinkSync(socketPath); } @@ -159,6 +201,9 @@ export function cleanupDaemonSocketPath( } if (lease) { assertSocketLease(socketPath, lease); + if (lease.compromise) { + return; + } try { cleanupUnixDaemonSocketPath(socketPath, expectedIdentity); } catch { @@ -167,16 +212,28 @@ export function cleanupDaemonSocketPath( return; } let releaseLock: (() => void) | undefined; + let lockCompromised = false; try { releaseLock = lockfile.lockSync(socketPath, { realpath: false, stale: DAEMON_SOCKET_LOCK_STALE_MS, update: DAEMON_SOCKET_LOCK_UPDATE_MS, + onCompromised: () => { + lockCompromised = true; + }, retries: 0, }); } catch { return; } + if (lockCompromised) { + try { + releaseLock(); + } catch { + // Best effort release. + } + return; + } try { cleanupUnixDaemonSocketPath(socketPath, expectedIdentity); } catch { @@ -213,6 +270,12 @@ function assertSocketLease(socketPath: string, lease: DaemonSocketPathLease): vo } } +function assertSocketLeaseHeld(socketPath: string, lease: DaemonSocketPathLease): void { + if (lease.compromise) { + throw new Error(`Daemon socket lease for ${socketPath} was compromised: ${lease.compromise.message}`); + } +} + export function defaultDaemonSocketDir(): string { const suffix = typeof process.getuid === "function" ? String(process.getuid()) : "user"; return join(tmpdir(), `prime-agent-${suffix}`); diff --git a/packages/coding-agent/src/modes/daemon/daemon-supervisor-ownership.ts b/packages/coding-agent/src/modes/daemon/daemon-supervisor-ownership.ts index 18fbd43f76..599fcd75db 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-supervisor-ownership.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-supervisor-ownership.ts @@ -7,6 +7,7 @@ import { realpathSync, renameSync, rmSync, + statSync, writeFileSync, } from "node:fs"; import { homedir } from "node:os"; @@ -356,11 +357,16 @@ function readLegacyOwnersForSocket( async function withDaemonSupervisorRegistryGuard(registryDir: string, action: () => T | Promise): Promise { mkdirSync(registryDir, { recursive: true, mode: 0o700 }); + const guardPath = resolve(registryDir, ".guard"); + let compromisedError: Error | undefined; const release = await lockfile.lock(registryDir, { realpath: false, - lockfilePath: resolve(registryDir, ".guard"), + lockfilePath: guardPath, stale: REGISTRY_LOCK_STALE_MS, update: REGISTRY_LOCK_UPDATE_MS, + onCompromised: (error) => { + compromisedError ??= error; + }, retries: { retries: REGISTRY_LOCK_RETRIES, factor: 1, @@ -368,10 +374,44 @@ async function withDaemonSupervisorRegistryGuard(registryDir: string, action: maxTimeout: REGISTRY_LOCK_RETRY_MS, }, }); + // Compromise detection is timer-driven and cannot preempt a synchronous stall: a stalled action's + // writes may already be on disk when a successor reclaims the stale guard. The guard directory's + // inode is the ownership identity (a steal is rmdir+mkdir), checked synchronously where the timer + // cannot run; when the inode is unobservable, only timer-driven detection applies. + const guardIno = (() => { + try { + return statSync(guardPath, { bigint: true }).ino; + } catch { + return undefined; + } + })(); + const guardStolen = () => { + if (guardIno === undefined) return false; + try { + return statSync(guardPath, { bigint: true }).ino !== guardIno; + } catch { + return true; + } + }; + const assertGuardHeld = () => { + if (compromisedError) + throw new Error(`Daemon supervisor registry guard was compromised: ${compromisedError.message}`); + if (guardStolen()) + throw new Error("Daemon supervisor registry guard was compromised: the guard lock changed hands"); + }; try { - return await action(); + assertGuardHeld(); + const result = await action(); + assertGuardHeld(); + return result; } finally { - await release(); + if (compromisedError) { + await release().catch(() => undefined); + } else if (!guardStolen()) { + await release(); + } + // A stolen-but-undetected guard is never released: that would delete the successor's lock. + // The abandoned updater notices the foreign mtime on its next tick and cleans itself up. } } diff --git a/packages/coding-agent/src/modes/daemon/daemon-supervisor.ts b/packages/coding-agent/src/modes/daemon/daemon-supervisor.ts index fb6549bde1..48260a7f18 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-supervisor.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-supervisor.ts @@ -123,7 +123,11 @@ import { isDaemonShutdownAdmissionActive, waitForDaemonStartupFence, } from "./daemon-supervisor-ownership.js"; -import { DaemonWorkerClient } from "./daemon-worker-client.js"; +import { + DaemonWorkerAuthenticationError, + DaemonWorkerClient, + DaemonWorkerProbeTimeoutError, +} from "./daemon-worker-client.js"; import { DAEMON_WORKER_ACTIVE_SESSION_ID_ENV, DAEMON_WORKER_INSTANCE_ID_ENV, @@ -180,6 +184,8 @@ const UPDATE_RESTART_WORKER_REQUEST_TIMEOUT_MS = 90_000; const UPDATE_RESTART_PREPARE_DEADLINE_MS = 100_000; const WORKER_RETRY_DELAYS_MS = [250, 1000, 5000] as const; const DEFERRED_RECOVERY_RECHECK_MS = 5000; +// ~2.5 minutes of probing: each round is one 5s defer recheck plus a ~11s three-delay probe pass. +const MAX_DEFERRED_RECOVERY_ROUNDS = 10; const STOP_FINALIZATION_RECHECK_MS = 250; const STOP_FINALIZATION_SIGKILL_GRACE_MS = 5000; const STOP_FINALIZATION_RETRY_MS = 5000; @@ -332,6 +338,8 @@ interface ResidentWorker { peerTransportCapable?: boolean; /** In-flight replacement connection during authentication; an allowed frame source alongside client. */ pendingClient?: DaemonWorkerClient; + /** Consecutive defer->probe rounds against a live-but-silent worker; bounded by MAX_DEFERRED_RECOVERY_ROUNDS. */ + deferredRecoveryRounds?: number; /** Bumped per applied roster frame; a summaries pull that straddles a frame must not gap-fill. */ rosterEpoch?: number; rosterApplyChain?: Promise; @@ -426,6 +434,10 @@ function isSupervisorRecoveryCancelled(error: unknown): boolean { return isSupervisorShutdownAdmissionCancelled(error) || isSupervisorGenerationStale(error); } +function isDaemonWorkerProbeTimeout(error: unknown): boolean { + return error instanceof DaemonWorkerProbeTimeoutError; +} + function isSupervisorShutdownAdmissionCancelled(error: unknown): boolean { return ( error instanceof SupervisorRecoveryCancelledError || @@ -642,9 +654,11 @@ export class DaemonSupervisor { private ownsSocketPath = false; private socketIdentity?: DaemonSocketIdentity; private socketLease?: DaemonSocketPathLease; + private socketLeaseCompromise?: Error; private ownership?: Awaited>; private cleanupPromise?: Promise; private shuttingDown = false; + private startupComplete = false; private updateRestartPhase?: "draining" | "fencing" | "prepared"; private readonly mutationDrain = new MutationDrainLatch(); private readonly clients = new Set(); @@ -713,7 +727,10 @@ export class DaemonSupervisor { throw new Error("Daemon supervisor config is missing agentDir"); } this.socketLease = await acquireDaemonSocketPathLease(this.socketPath); + this.socketLease?.onCompromised((error) => this.handleSocketLeaseCompromised(error)); + this.assertSocketLeaseHeld(); await waitForDaemonStartupFence(this.socketPath); + this.assertSocketLeaseHeld(); this.ownership = await acquireDaemonSupervisorOwnership({ socketPath: this.socketPath, descriptorDir: this.descriptorDir, @@ -721,6 +738,7 @@ export class DaemonSupervisor { generation: this.generation, appVersion: VERSION, }); + this.assertSocketLeaseHeld(); await prepareDaemonSocketPath(this.socketPath, this.socketLease); mkdirSync(this.descriptorDir, { recursive: true, mode: 0o700 }); @@ -734,6 +752,7 @@ export class DaemonSupervisor { this.server = createServer((socket) => this.handleConnection(socket)); await this.listen(); + this.assertSocketLeaseHeld(); this.socketIdentity = getDaemonSocketIdentity(this.socketPath); if (process.platform !== "win32" && !this.socketIdentity) { throw new Error(`Could not capture daemon socket identity: ${this.socketPath}`); @@ -755,6 +774,7 @@ export class DaemonSupervisor { this.log(`Migrated ${migratedJobs} scheduled jobs into session artifacts`); } await this.catalog.start().catch((error) => this.log(`Could not start daemon catalog: ${String(error)}`)); + this.assertSocketLeaseHeld(); await this.seedRosterLedger(); let adoptionFailure: unknown; let adoptionFailed = false; @@ -779,7 +799,10 @@ export class DaemonSupervisor { this.scheduleIdleEvictionSweep(); this.rosterWatchdogTimer = setInterval(() => this.sweepRosterStaleness(), ROSTER_WATCHDOG_INTERVAL_MS); this.rosterWatchdogTimer.unref(); + this.assertSocketLeaseHeld(); await this.ownership.updatePhase("owner"); + this.assertSocketLeaseHeld(); + this.startupComplete = true; this.log(`Prime Agent daemon supervisor ${this.generation} listening on ${this.socketPath}`); this.markReady(); } catch (error) { @@ -1045,8 +1068,14 @@ export class DaemonSupervisor { await ownership.assertCurrent(); } - private async assertRecoveryAllowed(): Promise { + private async assertServingCurrentOwnership(): Promise { + this.assertSupervisorServing(); await this.assertCurrentOwnership(); + this.assertSupervisorServing(); + } + + private async assertRecoveryAllowed(): Promise { + await this.assertServingCurrentOwnership(); if (await isDaemonShutdownAdmissionActive()) { throw new SupervisorRecoveryCancelledError("Daemon shutdown admission cancelled worker recovery"); } @@ -1468,6 +1497,12 @@ export class DaemonSupervisor { } private async handleLine(client: DaemonSocketClient, line: string): Promise { + try { + this.assertSupervisorServing(); + } catch (error) { + this.write(client, failure(salvageDaemonCommandId(line), "dispatch", error, serializeDaemonError(error))); + return; + } let preParsed: ReturnType; try { preParsed = this.parseCommandAndRegisterPromptAdmission(client, line); @@ -1522,7 +1557,7 @@ export class DaemonSupervisor { } try { - await waitForPromptAdmission(this.assertCurrentOwnership(), parsedAdmission?.controller.signal); + await waitForPromptAdmission(this.assertServingCurrentOwnership(), parsedAdmission?.controller.signal); } catch (error) { if (parsedAdmission) this.deletePromptAdmission(parsedAdmission); this.write(client, failure(command.id, command.type, error)); @@ -1562,6 +1597,19 @@ export class DaemonSupervisor { this.write(client, failure(command.id, command.type, "Daemon is preparing an update restart")); return; } + if (mutation && !UPDATE_RESTART_DRAIN_COMMANDS.has(command.type)) { + const idleEvictionFence = this.idleEvictionFence; + if (idleEvictionFence) { + await idleEvictionFence; + try { + await this.assertServingCurrentOwnership(); + } catch (error) { + if (parsedAdmission) this.deletePromptAdmission(parsedAdmission); + this.write(client, failure(command.id, command.type, error, serializeDaemonError(error))); + return; + } + } + } if (journalIdentity) { const admitted = this.commandJournal.begin(journalIdentity.clientId, journalIdentity.commandId, command.type); if (admitted.status === "complete") { @@ -1582,10 +1630,6 @@ export class DaemonSupervisor { } } - if (mutation && !UPDATE_RESTART_DRAIN_COMMANDS.has(command.type)) { - const idleEvictionFence = this.idleEvictionFence; - if (idleEvictionFence) await idleEvictionFence; - } // Attach is intentionally read-only and is not fence-gated. If eviction wins // the race, attach fails cleanly with "Session worker is not connected" and // the client retries through the saved-session path instead of mutating state. @@ -1937,6 +1981,7 @@ export class DaemonSupervisor { worker.descriptor.archiveOnStop = undefined; worker.descriptor.lifecycle = "recovering"; worker.descriptor.consecutiveFailures = 0; + worker.deferredRecoveryRounds = 0; this.persistWorker(worker); await this.recoverWorker(worker); if (this.workers.get(worker.descriptor.workerId)?.descriptor.lifecycle !== "ready") { @@ -2648,7 +2693,7 @@ export class DaemonSupervisor { return false; } worker.intentionalStop = true; - await this.recoverUncertainWorkerOperations(worker, false); + await this.recoverUncertainWorkerOperations(worker); this.invalidateWorkerSessionInputPauses(worker, "Session worker stopped while input was paused"); this.workers.delete(worker.descriptor.workerId); this.flipWorkerRosterEntriesInactive(worker); @@ -2908,6 +2953,7 @@ export class DaemonSupervisor { await this.assertRecoveryAllowed(); worker.descriptor.lifecycle = "ready"; worker.descriptor.consecutiveFailures = 0; + worker.deferredRecoveryRounds = 0; worker.descriptor.lastError = undefined; this.persistWorker(worker); if (!worker.descriptor.ownerClientId) { @@ -3017,13 +3063,17 @@ export class DaemonSupervisor { } catch (error) { lastError = error; client.close(); - if (isSupervisorRecoveryCancelled(error) || error instanceof PreRosterWorkerError) { + if ( + isSupervisorRecoveryCancelled(error) || + error instanceof PreRosterWorkerError || + error instanceof DaemonWorkerAuthenticationError + ) { throw error; } await delay(25); } } - throw new Error(`Timed out connecting to daemon session worker: ${String(lastError)}`); + throw new DaemonWorkerProbeTimeoutError(`Timed out connecting to daemon session worker: ${String(lastError)}`); } private async subscribeWorker(worker: ResidentWorker, activeSessionId: string): Promise { @@ -3095,6 +3145,7 @@ export class DaemonSupervisor { await this.assertRecoveryAllowed(); worker.descriptor.lifecycle = "ready"; worker.descriptor.consecutiveFailures = 0; + worker.deferredRecoveryRounds = 0; this.persistWorker(worker); this.broadcastHeartbeatsChanged(); } catch (error) { @@ -3114,6 +3165,16 @@ export class DaemonSupervisor { this.log(`Could not restart pre-roster worker ${worker.descriptor.workerId}: ${String(restartError)}`); } } + const identityNow = this.processIdentity(worker.descriptor.pid, worker.descriptor.processStartId); + if (isDaemonWorkerProbeTimeout(error) && (identityNow === "current" || identityNow === "unknown")) { + 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); } } @@ -3127,9 +3188,10 @@ export class DaemonSupervisor { worker.descriptor.processStartId = observedProcessStartId; } const identity = () => this.processIdentity(worker.descriptor.pid, worker.descriptor.processStartId); - const initialIdentity = identity(); - await this.recoverUncertainWorkerOperations(worker, initialIdentity === "current"); - if (initialIdentity === "current") { + // The one deliberate kill of a live worker: it authenticated as ours and predates the roster + // protocol, so this identity-verified upgrade replaces it. No await between recheck and signal. + if (identity() === "current") { + signalProcessGroupOrProcess(worker.descriptor.pid, "SIGKILL"); // SIGKILL is uninterceptable; this wait only covers kernel teardown of the old process and socket. const killDeadline = Date.now() + 1000; while (identity() === "current" && Date.now() < killDeadline) { @@ -3138,6 +3200,8 @@ export class DaemonSupervisor { } const finalIdentity = identity(); if (finalIdentity !== "gone" && finalIdentity !== "replaced") { + // A live or unverifiable survivor parks failed with no destructive cleanup: interruption + // marking and orphan reaping must never run against a possibly-active worker. worker.descriptor.lifecycle = "failed"; worker.descriptor.lastError = `Pre-roster worker process ${worker.descriptor.pid} is still running and cannot be replaced safely`; this.persistWorker(worker); @@ -3145,6 +3209,7 @@ export class DaemonSupervisor { this.log(`Kept pre-roster worker ${worker.descriptor.workerId} failed: ${worker.descriptor.lastError}`); return; } + await this.recoverUncertainWorkerOperations(worker); if (this.isWorkerRecoveryCancelled(worker)) { return; } @@ -3224,6 +3289,19 @@ export class DaemonSupervisor { if (worker.deferredRecovery) { return; } + // A live-but-silent worker must not probe forever: park it failed (user-visible through the + // roster's failed status) and keep its process alive for a manual retry_worker. + worker.deferredRecoveryRounds = (worker.deferredRecoveryRounds ?? 0) + 1; + if (worker.deferredRecoveryRounds > MAX_DEFERRED_RECOVERY_ROUNDS) { + worker.descriptor.lifecycle = "failed"; + worker.descriptor.lastError = `Live session worker did not answer recovery probes for ${MAX_DEFERRED_RECOVERY_ROUNDS} rounds: ${disconnectError.message}`; + this.persistWorker(worker); + this.markWorkerRosterEntries(worker, "failed"); + this.log( + `Worker ${worker.descriptor.workerId} is unresponsive; parked failed after ${MAX_DEFERRED_RECOVERY_ROUNDS} probe rounds`, + ); + return; + } worker.deferredRecovery = this.resumeDeferredWorkerRecovery(worker, disconnectError).finally(() => { worker.deferredRecovery = undefined; }); @@ -3455,19 +3533,20 @@ export class DaemonSupervisor { return worker.recovery; } worker.recovery = (async () => { - for (const [retryIndex, retryDelay] of WORKER_RETRY_DELAYS_MS.entries()) { + let keepProbingLiveWorker = false; + for (const retryDelay of WORKER_RETRY_DELAYS_MS) { await delay(retryDelay); + keepProbingLiveWorker = false; if (this.isWorkerRecoveryCancelled(worker)) { return; } try { await this.assertRecoveryAllowed(); - const processAlive = isProcessAlive(worker.descriptor.pid); - const observedProcessStartId = processAlive ? getProcessStartId(worker.descriptor.pid) : undefined; - const processIdentityMatches = - worker.descriptor.processStartId === undefined || - observedProcessStartId === worker.descriptor.processStartId; - if (processAlive && processIdentityMatches) { + const identityNow = this.processIdentity(worker.descriptor.pid, worker.descriptor.processStartId); + const identityCompatible = + identityNow === "current" || + (identityNow === "unknown" && worker.descriptor.processStartId === undefined); + if (identityCompatible) { try { await this.connectWorker(worker, 1500); await this.subscribeWorker(worker, worker.descriptor.rootActiveSessionId); @@ -3475,12 +3554,16 @@ export class DaemonSupervisor { if (this.isWorkerRecoveryCancelled(worker)) { return; } - if (worker.descriptor.processStartId === undefined && observedProcessStartId) { - worker.descriptor.processStartId = observedProcessStartId; + if (worker.descriptor.processStartId === undefined) { + const observedProcessStartId = getProcessStartId(worker.descriptor.pid); + if (observedProcessStartId) { + worker.descriptor.processStartId = observedProcessStartId; + } } await this.assertRecoveryAllowed(); worker.descriptor.lifecycle = "ready"; worker.descriptor.consecutiveFailures = 0; + worker.deferredRecoveryRounds = 0; this.persistWorker(worker); this.broadcastHeartbeatsChanged(); return; @@ -3491,31 +3574,28 @@ export class DaemonSupervisor { await this.assertRecoveryAllowed(); worker.client?.close(); worker.client = undefined; - if (retryIndex < WORKER_RETRY_DELAYS_MS.length - 1) { - throw error; - } + // A worker with the same durable process identity may be load-slow. + // Keep probing it instead of replacing live work after a timeout. + keepProbingLiveWorker = isDaemonWorkerProbeTimeout(error) || identityNow !== "current"; + throw error; } } - if ( - processAlive && - (worker.descriptor.processStartId === undefined || observedProcessStartId === undefined) - ) { + if (identityNow === "unknown") { + keepProbingLiveWorker = true; throw new Error( `Cannot safely replace live session worker ${worker.descriptor.workerId} without a verified process identity`, ); } const recoveryCommand = worker.descriptor.ownerClientId ? worker.transientCreateCommand : undefined; if (!recoveryCommand || !worker.launchEnv) { - await this.recoverUncertainWorkerOperations(worker, false); + await this.recoverUncertainWorkerOperations(worker); worker.descriptor.lifecycle = "failed"; worker.descriptor.lastError = "Waiting for a client with fresh runtime context"; this.persistWorker(worker); this.markWorkerRosterEntries(worker, "failed"); return; } - const safeToKillWorkerProcess = - processAlive && processIdentityMatches && worker.descriptor.processStartId !== undefined; - await this.recoverUncertainWorkerOperations(worker, safeToKillWorkerProcess); + await this.recoverUncertainWorkerOperations(worker); if (this.isWorkerRecoveryCancelled(worker)) { return; } @@ -3538,11 +3618,28 @@ export class DaemonSupervisor { this.persistWorker(worker); } } + if (keepProbingLiveWorker) { + try { + await this.assertRecoveryAllowed(); + } catch { + return; + } + 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 { return; } + // Leak-over-kill: a live worker that keeps failing for non-timeout reasons parks failed with + // its process intact. A verified-identity survivor is reclaimed by the next fresh create; + // an unverifiable one waits for exit — killing a pid we cannot verify as ours is worse. worker.descriptor.lifecycle = "failed"; this.persistWorker(worker); this.markWorkerRosterEntries(worker, "failed"); @@ -3562,35 +3659,25 @@ export class DaemonSupervisor { ); } - private async recoverUncertainWorkerOperations(worker: ResidentWorker, killWorkerProcess = true): Promise { + private isWorkerCleanupCancelled(worker: ResidentWorker): boolean { + return ( + this.shuttingDown || + worker.descriptor.stopRequestedAt !== undefined || + this.workers.get(worker.descriptor.workerId) !== worker + ); + } + + private async recoverUncertainWorkerOperations(worker: ResidentWorker): Promise { await this.assertRecoveryAllowed(); - // Re-check at the last synchronous moment: the process can exit in the await gap and the PID recycle. - if ( - killWorkerProcess && - this.processIdentity(worker.descriptor.pid, worker.descriptor.processStartId) === "current" - ) { - signalProcessGroupOrProcess(worker.descriptor.pid, "SIGKILL"); - } - const orphanProcessJournalPath = worker.descriptor.orphanProcessJournalPath; - if (orphanProcessJournalPath) { - try { - for (const orphan of readActiveOrphanProcesses(orphanProcessJournalPath, worker.descriptor.pid)) { - if (!shouldReapOrphanProcess(orphan)) { - continue; - } - killOrphanProcess(orphan.pid); - } - clearOrphanProcessJournal(orphanProcessJournalPath); - } catch (error) { - this.log(`Could not reap orphaned worker resources: ${String(error)}`); - } - } const journal = new WorkerRecoveryJournal(worker.descriptor.recoveryJournalPath); const latest = journal.getLatest(); const uncertain = latest.filter((record) => record.busy); - if (uncertain.length === 0) { - return; + if (uncertain.length > 0) await this.catalog.start(); + await this.assertRecoveryAllowed(); + if (this.isWorkerCleanupCancelled(worker)) { + throw new SupervisorRecoveryCancelledError("Worker recovery was cancelled before destructive cleanup"); } + const interruptedSessions = new Map< string, { activeSessionId: string; sessionFile: string; operations: Set } @@ -3612,7 +3699,11 @@ export class DaemonSupervisor { } interrupted.operations.add(record.operation); } + await this.assertRecoveryAllowed(); + if (this.isWorkerCleanupCancelled(worker)) { + throw new SupervisorRecoveryCancelledError("Worker recovery was cancelled before interruption was recorded"); + } await Promise.all( [...interruptedSessions.values()].map((interrupted) => this.catalog.markInterrupted(interrupted.sessionFile, interrupted.activeSessionId, [ @@ -3621,6 +3712,30 @@ export class DaemonSupervisor { ), ); await this.assertRecoveryAllowed(); + if (this.isWorkerCleanupCancelled(worker)) { + throw new SupervisorRecoveryCancelledError("Worker recovery was cancelled before process cleanup"); + } + const orphanProcessJournalPath = worker.descriptor.orphanProcessJournalPath; + if (orphanProcessJournalPath) { + try { + const orphans = readActiveOrphanProcesses(orphanProcessJournalPath, worker.descriptor.pid); + let reapFailed = false; + let retryIsSafe = true; + for (const orphan of orphans) { + if (orphan.processStartId === undefined) retryIsSafe = false; + if (!shouldReapOrphanProcess(orphan)) { + continue; + } + if (!killOrphanProcess(orphan.pid)) reapFailed = true; + } + if (!reapFailed || !retryIsSafe) clearOrphanProcessJournal(orphanProcessJournalPath); + } catch (error) { + this.log(`Could not reap orphaned worker resources: ${String(error)}`); + } + } + if (uncertain.length === 0) { + return; + } for (const record of latest) { journal.record({ activeSessionId: record.activeSessionId, @@ -3636,7 +3751,6 @@ export class DaemonSupervisor { .join(", ")}`, ); } - private async refreshWorkerSummaries( worker: ResidentWorker, recovery = false, @@ -4548,6 +4662,7 @@ export class DaemonSupervisor { ownedWorker.descriptor.archiveOnStop = undefined; ownedWorker.descriptor.lifecycle = "recovering"; ownedWorker.descriptor.consecutiveFailures = 0; + ownedWorker.deferredRecoveryRounds = 0; this.persistWorker(ownedWorker); await this.recoverWorker(ownedWorker); } @@ -6180,6 +6295,49 @@ export class DaemonSupervisor { this.signalCleanupHandlers.push(() => process.off("exit", exitHandler)); } + private assertSocketLeaseHeld(): void { + const compromise = this.socketLeaseCompromise ?? this.socketLease?.compromise; + if (compromise) throw new Error(`Daemon socket lease was compromised: ${compromise.message}`); + } + + private assertSupervisorServing(): void { + this.assertSocketLeaseHeld(); + if (this.shuttingDown) { + const error = new Error(`Daemon supervisor generation ${this.generation} is shutting down; retry the command`); + Object.assign(error, { code: "supervisor_generation_stale" as const }); + throw error; + } + } + + private fenceSupervisorSocket(): void { + try { + this.server?.close(); + } catch { + // The server may already be closed by a concurrent shutdown. + } + for (const client of this.clients) { + client.detachInput(); + client.socket.destroy(); + } + } + + private handleSocketLeaseCompromised(error: Error): void { + if (this.socketLeaseCompromise) return; + this.socketLeaseCompromise = error; + this.shuttingDown = true; + this.fenceSupervisorSocket(); + const message = `Daemon socket lease was compromised; relinquishing supervisor ownership: ${error.message}`; + try { + this.log(message); + } catch { + console.error(message); + } + if (!this.startupComplete) return; + void this.cleanupSupervisorResources().catch((cleanupError) => + this.reportCleanupFailure("compromised daemon socket lease", cleanupError), + ); + } + private cleanupSocket(): void { if (!this.ownsSocketPath) { return; diff --git a/packages/coding-agent/src/modes/daemon/daemon-worker-client.ts b/packages/coding-agent/src/modes/daemon/daemon-worker-client.ts index 16f4751ce0..8e416ee5ad 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-worker-client.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-worker-client.ts @@ -33,6 +33,11 @@ export type DaemonWorkerFrameListener = (frame: PrivateFrame void; type DaemonHello = Extract; +export class DaemonWorkerAuthenticationError extends Error {} + +/** A probe (hello/response/connect) timed out; recovery treats the worker as live-but-slow, never dead. */ +export class DaemonWorkerProbeTimeoutError extends Error {} + export class DaemonWorkerClient { private socket?: Socket; private channel?: PrivateFramedChannel; @@ -122,7 +127,7 @@ export class DaemonWorkerClient { reject, timeout: setTimeout(() => { this.helloWaiters.delete(waiter); - reject(new Error("Timed out waiting for daemon worker hello")); + reject(new DaemonWorkerProbeTimeoutError("Timed out waiting for daemon worker hello")); }, timeoutMs), }; this.helloWaiters.add(waiter); @@ -164,7 +169,7 @@ export class DaemonWorkerClient { ): Promise> { const response = await this.requestWorker({ type: "worker_auth", token, ...owner }, timeoutMs); if (!response.success) { - throw new Error(response.error); + throw new DaemonWorkerAuthenticationError(response.error); } return response; } @@ -205,7 +210,9 @@ export class DaemonWorkerClient { const response = new Promise((resolve, reject) => { const timeout = setTimeout(() => { this.pending.delete(id); - reject(new Error(`Timed out waiting for daemon worker response to ${command.type}`)); + reject( + new DaemonWorkerProbeTimeoutError(`Timed out waiting for daemon worker response to ${command.type}`), + ); }, timeoutMs); this.pending.set(id, { resolve, reject, timeout }); }); diff --git a/packages/coding-agent/test/daemon-agent-roster.test.ts b/packages/coding-agent/test/daemon-agent-roster.test.ts index d882b0ed19..bd1dcd47bb 100644 --- a/packages/coding-agent/test/daemon-agent-roster.test.ts +++ b/packages/coding-agent/test/daemon-agent-roster.test.ts @@ -15,6 +15,7 @@ import type { SessionSummary } from "../src/modes/daemon/daemon-session-list.js" import { DaemonSupervisor } from "../src/modes/daemon/daemon-supervisor.js"; import type { DaemonWorkerRosterOutbound } from "../src/modes/daemon/daemon-worker-protocol.js"; import { RlmSpawnLedger } from "../src/modes/daemon/rlm-ledger.js"; +import * as childProcessModule from "../src/utils/child-process.js"; type RosterDelta = Extract; @@ -1529,9 +1530,38 @@ describe("review-round regressions", () => { } ).restartPreRosterWorker(worker, undefined); - expect(recoverUncertainWorkerOperations).toHaveBeenCalledWith(worker, false); + // No destructive cleanup against a possibly-live worker. + expect(recoverUncertainWorkerOperations).not.toHaveBeenCalled(); expect(launchWorker).not.toHaveBeenCalled(); expect(worker.descriptor.lifecycle).toBe("failed"); + + // The one deliberate kill of a live worker: identity-verified pre-roster adoption replaces it. + const killed = makeWorker("worker-2"); + Object.assign(killed.descriptor, { pid: 987_654, processStartId: "start-1", createCommand: { type: "create" } }); + const launchReplacement = vi.fn(); + const cleanup = vi.fn(async () => {}); + const signal = vi.spyOn(childProcessModule, "signalProcessGroupOrProcess").mockImplementation(() => {}); + const identities = ["current", "gone"]; + const killSupervisor = makeSupervisor([killed], { + assertRecoveryAllowed: vi.fn(async () => {}), + recoverUncertainWorkerOperations: cleanup, + launchWorker: launchReplacement, + processIdentity: vi.fn(() => (identities.length > 1 ? identities.shift() : identities[0])), + }); + try { + await ( + killSupervisor as unknown as { + restartPreRosterWorker(worker: WorkerFixture, observedProcessStartId?: string): Promise; + } + ).restartPreRosterWorker(killed, "start-1"); + expect(signal).toHaveBeenCalledWith(987_654, "SIGKILL"); + // Destructive cleanup only after the predecessor is confirmed gone. + expect(cleanup).toHaveBeenCalledOnce(); + expect(signal.mock.invocationCallOrder[0]).toBeLessThan(cleanup.mock.invocationCallOrder[0]!); + expect(launchReplacement).toHaveBeenCalled(); + } finally { + signal.mockRestore(); + } }); }); diff --git a/packages/coding-agent/test/daemon-catalog-entry.test.ts b/packages/coding-agent/test/daemon-catalog-entry.test.ts new file mode 100644 index 0000000000..52a5b00ef3 --- /dev/null +++ b/packages/coding-agent/test/daemon-catalog-entry.test.ts @@ -0,0 +1,24 @@ +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { describe, expect, it } from "vitest"; +import { ENV_AGENT_DIR } from "../src/config.js"; +import { DaemonCatalogClient } from "../src/modes/daemon/daemon-catalog-process.js"; + +describe("daemon catalog entrypoint", () => { + it("starts the dedicated catalog process over IPC", async () => { + const agentDir = mkdtempSync(join(tmpdir(), "pa-catalog-entry-")); + const previousAgentDir = process.env[ENV_AGENT_DIR]; + process.env[ENV_AGENT_DIR] = agentDir; + const client = new DaemonCatalogClient(() => {}); + try { + await expect(client.start()).resolves.toBeUndefined(); + await expect(client.list()).resolves.toEqual([]); + } finally { + await client.stop(); + if (previousAgentDir === undefined) delete process.env[ENV_AGENT_DIR]; + else process.env[ENV_AGENT_DIR] = previousAgentDir; + rmSync(agentDir, { recursive: true, force: true }); + } + }, 10_000); +}); diff --git a/packages/coding-agent/test/daemon-catalog-startup.test.ts b/packages/coding-agent/test/daemon-catalog-startup.test.ts new file mode 100644 index 0000000000..a4bc7f2d8f --- /dev/null +++ b/packages/coding-agent/test/daemon-catalog-startup.test.ts @@ -0,0 +1,63 @@ +import type { ChildProcess, SpawnOptions } from "node:child_process"; +import { EventEmitter } from "node:events"; +import { afterEach, describe, expect, it, vi } from "vitest"; + +const spawnState = vi.hoisted(() => ({ + args: [] as string[], + child: undefined as (EventEmitter & { connected: boolean }) | undefined, +})); + +vi.mock("node:child_process", () => ({ + spawn(_command: string, args: string[], _options: SpawnOptions): ChildProcess { + spawnState.args = args; + const child = Object.assign(new EventEmitter(), { + connected: true, + disconnect: vi.fn(), + kill: vi.fn(), + send: vi.fn(), + }); + spawnState.child = child; + return child as unknown as ChildProcess; + }, +})); + +import { DaemonCatalogClient, isDaemonCatalogSourcePath } from "../src/modes/daemon/daemon-catalog-process.js"; + +afterEach(() => { + vi.useRealTimers(); + spawnState.args = []; + spawnState.child = undefined; +}); + +describe("daemon catalog startup", () => { + it("does not mistake an ancestor src directory for the package source tree", () => { + const packageDir = "/usr/src/app/packages/coding-agent"; + + expect(isDaemonCatalogSourcePath(`${packageDir}/dist/modes/daemon/daemon-catalog-process.js`, packageDir)).toBe( + false, + ); + expect(isDaemonCatalogSourcePath(`${packageDir}/src/modes/daemon/daemon-catalog-process.ts`, packageDir)).toBe( + true, + ); + }); + + it("uses the dedicated entrypoint and allows a cold start past five seconds", async () => { + vi.useFakeTimers(); + const client = new DaemonCatalogClient(() => {}); + const starting = client.start(); + + expect(spawnState.args.some((arg) => /daemon-catalog-entry\.(?:js|ts)$/.test(arg))).toBe(true); + await vi.advanceTimersByTimeAsync(6000); + spawnState.child?.emit("message", { type: "ready" }); + await expect(starting).resolves.toBeUndefined(); + }); + + it("rejects immediately when the catalog exits during startup", async () => { + vi.useFakeTimers(); + const client = new DaemonCatalogClient(() => {}); + const starting = client.start(); + + spawnState.child?.emit("exit", 1, null); + await expect(starting).rejects.toThrow(/exited during startup/); + }); +}); diff --git a/packages/coding-agent/test/daemon-launch.test.ts b/packages/coding-agent/test/daemon-launch.test.ts index 1a649d6cc6..26cc049ef3 100644 --- a/packages/coding-agent/test/daemon-launch.test.ts +++ b/packages/coding-agent/test/daemon-launch.test.ts @@ -3,7 +3,7 @@ import { existsSync, mkdirSync, mkdtempSync, rmSync, writeFileSync } from "node: import { createServer, type Server, type Socket } from "node:net"; import { tmpdir } from "node:os"; import { dirname, join } from "node:path"; -import { afterEach, describe, expect, it } from "vitest"; +import { afterEach, describe, expect, it, vi } from "vitest"; import { ensureInteractiveDaemonRunning, probeDaemonVersion, @@ -25,7 +25,10 @@ interface FakeDaemonOptions { protocolVersion?: number; appVersion?: string; schemaId?: string; + firstSchemaId?: string; serverCapabilities?: string[]; + shouldSendHello?: (connectionIndex: number) => boolean; + onConnection?: (connectionIndex: number) => void; onCommand?: (command: { type: string }) => void; } @@ -41,17 +44,25 @@ function send(socket: Socket, message: unknown): void { async function startFakeDaemon(options: FakeDaemonOptions = {}): Promise { const dir = mkdtempSync(join(tmpdir(), "pa-launch-")); const socketPath = join(dir, "d.sock"); + let connectionIndex = 0; const server: Server = createServer((socket) => { + const currentConnectionIndex = connectionIndex++; + options.onConnection?.(currentConnectionIndex); socket.on("error", () => undefined); - send(socket, { - type: "daemon_hello", - socketPath, - protocol: { name: "prime-agent.daemon", version: options.protocolVersion ?? DAEMON_PROTOCOL_VERSION }, - appVersion: options.appVersion, - schemaId: options.schemaId ?? DAEMON_SCHEMA_ID, - clientId: "fake-client", - serverCapabilities: options.serverCapabilities ?? [], - }); + if (options.shouldSendHello?.(currentConnectionIndex) ?? true) { + send(socket, { + type: "daemon_hello", + socketPath, + protocol: { name: "prime-agent.daemon", version: options.protocolVersion ?? DAEMON_PROTOCOL_VERSION }, + appVersion: options.appVersion, + schemaId: + currentConnectionIndex === 0 && options.firstSchemaId + ? options.firstSchemaId + : (options.schemaId ?? DAEMON_SCHEMA_ID), + clientId: "fake-client", + serverCapabilities: options.serverCapabilities ?? [], + }); + } let buffer = ""; socket.on("data", (chunk) => { buffer += chunk.toString(); @@ -273,6 +284,91 @@ describe("ensureInteractiveDaemonRunning", () => { await expect(probe).resolves.toMatchObject({ status: "current" }); }); + it("classifies a connected daemon without hello as unresponsive", async () => { + const daemon = await startFakeDaemon({ shouldSendHello: () => false }); + cleanups.push(daemon.close); + + await expect(probeDaemonVersion(daemon.socketPath, 10)).resolves.toEqual({ status: "unresponsive" }); + }); + + it("reuses a daemon that becomes current inside the startup window", async () => { + vi.useFakeTimers(); + let resolveFirstConnection = () => {}; + let resolveSecondConnection = () => {}; + const firstConnection = new Promise((resolve) => { + resolveFirstConnection = resolve; + }); + const secondConnection = new Promise((resolve) => { + resolveSecondConnection = resolve; + }); + const daemon = await startFakeDaemon({ + appVersion: VERSION, + shouldSendHello: (connectionIndex) => connectionIndex > 0, + onConnection: (connectionIndex) => { + if (connectionIndex === 0) resolveFirstConnection(); + if (connectionIndex === 1) resolveSecondConnection(); + }, + }); + cleanups.push(daemon.close); + + try { + const ensuring = ensureInteractiveDaemonRunning(daemon.socketPath); + await firstConnection; + await vi.advanceTimersByTimeAsync(2000); + await secondConnection; + await expect(ensuring).resolves.toBeUndefined(); + } finally { + vi.useRealTimers(); + } + }); + + it("leaves a connected unresponsive daemon running after the startup window", async () => { + vi.useFakeTimers(); + const connections: Array<() => void> = []; + const connected = [0, 1].map( + (index) => + new Promise((resolve) => { + connections[index] = resolve; + }), + ); + const commands: string[] = []; + const daemon = await startFakeDaemon({ + shouldSendHello: () => false, + onConnection: (connectionIndex) => connections[connectionIndex]?.(), + onCommand: (command) => commands.push(command.type), + }); + cleanups.push(daemon.close); + + try { + const ensuring = ensureInteractiveDaemonRunning(daemon.socketPath); + const rejected = expect(ensuring).rejects.toThrow(/accepted connections but did not finish startup/); + await connected[0]; + await vi.advanceTimersByTimeAsync(2000); + await connected[1]; + await vi.advanceTimersByTimeAsync(30_000); + await rejected; + expect(commands).not.toContain("list"); + expect(commands).not.toContain("shutdown"); + } finally { + vi.useRealTimers(); + } + }); + + it("cancels replacement when a stale-looking daemon becomes current", async () => { + const commands: string[] = []; + const daemon = await startFakeDaemon({ + appVersion: VERSION, + firstSchemaId: "stale-schema", + schemaId: DAEMON_SCHEMA_ID, + onCommand: (command) => commands.push(command.type), + }); + cleanups.push(daemon.close); + + await expect(ensureInteractiveDaemonRunning(daemon.socketPath)).resolves.toBeUndefined(); + expect(commands).toContain("list"); + expect(commands).not.toContain("shutdown"); + }); + it("fails fast with the daemon log tail when the spawned daemon exits during startup", async () => { const dir = mkdtempSync(join(tmpdir(), "pa-launch-startup-crash-")); const entrypoint = join(dir, "crash.mjs"); diff --git a/packages/coding-agent/test/daemon-socket.test.ts b/packages/coding-agent/test/daemon-socket.test.ts index c266906118..8237d68f23 100644 --- a/packages/coding-agent/test/daemon-socket.test.ts +++ b/packages/coding-agent/test/daemon-socket.test.ts @@ -7,6 +7,7 @@ import lockfile from "proper-lockfile"; import { describe, expect, it } from "vitest"; import { cleanupDaemonSocketPath, + DaemonSocketPathLease, defaultDaemonSocketPath, getDaemonSocketIdentity, normalizeSocketPath, @@ -254,3 +255,36 @@ describe("defaultDaemonSocketPath", () => { } }); }); + +describe.skipIf(process.platform === "win32")("DaemonSocketPathLease compromise hardening", () => { + it("records compromise without rethrowing listener failures", () => { + const lease = new DaemonSocketPathLease("/tmp/test.sock", () => Promise.resolve()); + const observed: Error[] = []; + lease.onCompromised(() => { + throw new Error("listener failed"); + }); + lease.onCompromised((error) => observed.push(error)); + + expect(() => lease.recordCompromise(new Error("lock update failed"))).not.toThrow(); + expect(lease.compromise?.message).toBe("lock update failed"); + expect(observed).toHaveLength(1); + }); + + it("does not unlink a successor socket after the old lease is compromised", async () => { + const dir = mkdtempSync(join(tmpdir(), "pa-socket-compromise-")); + const socketPath = join(dir, "daemon.sock"); + const server = createServer(); + try { + await new Promise((resolve) => server.listen(socketPath, resolve)); + const identity = getDaemonSocketIdentity(socketPath); + const lease = new DaemonSocketPathLease(socketPath, () => Promise.resolve()); + lease.recordCompromise(new Error("lock stolen")); + + cleanupDaemonSocketPath(socketPath, identity, lease); + expect(existsSync(socketPath)).toBe(true); + } finally { + await new Promise((resolve) => server.close(() => resolve())); + rmSync(dir, { recursive: true, force: true }); + } + }); +}); diff --git a/packages/coding-agent/test/daemon-supervisor-admission.test.ts b/packages/coding-agent/test/daemon-supervisor-admission.test.ts index 842681c490..98b5089ef5 100644 --- a/packages/coding-agent/test/daemon-supervisor-admission.test.ts +++ b/packages/coding-agent/test/daemon-supervisor-admission.test.ts @@ -219,6 +219,76 @@ describe("daemon supervisor prompt admission ownership", () => { await Promise.all([first, second]); }); + it("journals a successful mutation when graceful shutdown begins during dispatch", async () => { + const commandJournal = { + lookup: vi.fn(() => undefined), + begin: vi.fn(() => ({ status: "new" as const })), + recordResult: vi.fn(), + acknowledge: vi.fn(), + }; + const response = { type: "response", command: "prompt", success: true } as const; + let supervisor: SupervisorHarness; + supervisor = createHarness({ + commandJournal, + findWorker: vi.fn(async () => ({ + worker: { descriptor: { lifecycle: "ready", rootActiveSessionId: "session-1" } }, + summary: { id: "session-1", activeSessionId: "session-1" }, + })), + forwardToWorker: vi.fn(async () => { + (supervisor as unknown as { shuttingDown: boolean }).shuttingDown = true; + return response; + }), + }); + const owner = client("connection-owner"); + const command = createDaemonCommandEnvelope( + { id: "prompt-1", type: "prompt", activeSessionId: "session-1", message: "hello" }, + "prompt-1", + "logical-client", + ); + + await supervisor.handleLine(owner, JSON.stringify(command)); + + expect(commandJournal.recordResult).toHaveBeenCalledWith("logical-client", "prompt-1", response); + expect((supervisor as unknown as { write: ReturnType }).write).toHaveBeenLastCalledWith( + owner, + response, + ); + }); + + it("does not journal a mutation before an eviction-fence ownership recheck", async () => { + const idleEvictionFence = deferred(); + let ownershipChecks = 0; + const commandJournal = { + lookup: vi.fn(() => undefined), + begin: vi.fn(() => ({ status: "new" as const })), + recordResult: vi.fn(), + acknowledge: vi.fn(), + }; + const supervisor = createHarness({ + commandJournal, + assertCurrent: vi.fn(async () => { + ownershipChecks++; + if (ownershipChecks > 1) throw new Error("supervisor ownership changed"); + }), + }); + (supervisor as unknown as { idleEvictionFence: Promise }).idleEvictionFence = idleEvictionFence.promise; + const owner = client("connection-owner"); + const command = createDaemonCommandEnvelope( + { id: "prompt-1", type: "prompt", activeSessionId: "session-1", message: "hello" }, + "prompt-1", + "logical-client", + ); + + const pending = supervisor.handleLine(owner, JSON.stringify(command)); + await waitFor(() => ownershipChecks === 1); + expect(commandJournal.begin).not.toHaveBeenCalled(); + idleEvictionFence.resolve(); + await pending; + + expect(commandJournal.begin).not.toHaveBeenCalled(); + expect(commandJournal.recordResult).not.toHaveBeenCalled(); + }); + it("lets the originating connection cancel before worker lookup starts", async () => { const ready = deferred(); const findWorker = vi.fn(async () => { diff --git a/packages/coding-agent/test/daemon-supervisor-monitor.test.ts b/packages/coding-agent/test/daemon-supervisor-monitor.test.ts index ebc6e7f14a..8d6954ffde 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 * as orphanProcessModule from "../src/core/orphan-process-journal.js"; 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"; @@ -19,7 +20,13 @@ import { success, } from "../src/modes/daemon/daemon-protocol.js"; import type { SessionSummary } from "../src/modes/daemon/daemon-session-list.js"; +import { DaemonSocketPathLease } from "../src/modes/daemon/daemon-socket.js"; import { DaemonSupervisor } from "../src/modes/daemon/daemon-supervisor.js"; +import { + DaemonWorkerAuthenticationError, + DaemonWorkerClient, + DaemonWorkerProbeTimeoutError, +} from "../src/modes/daemon/daemon-worker-client.js"; import { DAEMON_WORKER_STARTUP_GATE_COMMIT, DAEMON_WORKER_SUPERVISOR_SOCKET_ENV, @@ -1096,6 +1103,93 @@ describe("daemon worker supervisor monitoring", () => { } }); + it("fences the startup socket and leaves resource cleanup to the startup failure path after lease compromise", () => { + const cleanupSupervisorResources = vi.fn(async () => {}); + const fenceSupervisorSocket = vi.fn(); + const lease = new DaemonSocketPathLease("/tmp/daemon.sock", async () => {}); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + shuttingDown: false, + startupComplete: false, + socketLease: lease, + cleanupSupervisorResources, + fenceSupervisorSocket, + log: vi.fn(), + }) as { + shuttingDown: boolean; + assertSocketLeaseHeld(): void; + handleSocketLeaseCompromised(error: Error): void; + }; + lease.onCompromised((error) => supervisor.handleSocketLeaseCompromised(error)); + + lease.recordCompromise(new Error("lock refresh failed")); + + expect(supervisor.shuttingDown).toBe(true); + expect(fenceSupervisorSocket).toHaveBeenCalledOnce(); + expect(cleanupSupervisorResources).not.toHaveBeenCalled(); + expect(() => supervisor.assertSocketLeaseHeld()).toThrow(/lease was compromised/); + }); + + it("relinquishes supervisor resources after socket lease compromise even when logging fails", async () => { + const cleanupSupervisorResources = vi.fn(async () => {}); + const consoleError = vi.spyOn(console, "error").mockImplementation(() => {}); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + shuttingDown: false, + startupComplete: true, + cleanupSupervisorResources, + fenceSupervisorSocket: vi.fn(), + log: vi.fn(() => { + throw new Error("log failed"); + }), + reportCleanupFailure: vi.fn(), + }) as { + shuttingDown: boolean; + handleSocketLeaseCompromised(error: Error): void; + }; + const lease = new DaemonSocketPathLease("/tmp/daemon.sock", async () => {}); + lease.onCompromised((error) => supervisor.handleSocketLeaseCompromised(error)); + + try { + lease.recordCompromise(new Error("lock refresh failed")); + await Promise.resolve(); + + expect(supervisor.shuttingDown).toBe(true); + expect(cleanupSupervisorResources).toHaveBeenCalledOnce(); + } finally { + consoleError.mockRestore(); + } + }); + + it("rejects commands immediately after the supervisor is fenced", async () => { + const writes: string[] = []; + const client = { + id: "client-fenced", + socket: { + destroyed: false, + write: vi.fn((chunk: string) => { + writes.push(chunk); + return true; + }), + }, + } as unknown as DaemonSocketClient; + const handleCommand = vi.fn(); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + shuttingDown: true, + generation: "fenced-generation", + socketPath: "/tmp/fenced.sock", + handleCommand, + }) as { + handleLine(target: DaemonSocketClient, line: string): Promise; + }; + + await supervisor.handleLine( + client, + JSON.stringify(createDaemonCommandEnvelope({ type: "list" }, "command-fenced", "client-fenced")), + ); + + expect(writes.join(" ")).toContain("is shutting down"); + expect(handleCommand).not.toHaveBeenCalled(); + }); + it("does not poll a healthy supervisor after the startup check", async () => { vi.useFakeTimers(); let resolveProbe: () => void = () => undefined; @@ -1678,7 +1772,7 @@ describe("daemon worker supervisor monitoring", () => { await recovery; expect(supervisor.connectWorker).not.toHaveBeenCalled(); - expect(supervisor.recoverUncertainWorkerOperations).toHaveBeenCalledWith(worker, false); + expect(supervisor.recoverUncertainWorkerOperations).toHaveBeenCalledWith(worker); expect(supervisor.launchWorker).not.toHaveBeenCalled(); expect(worker.descriptor.lifecycle).toBe("failed"); expect(supervisor.persistWorker).toHaveBeenCalledWith(worker); @@ -1913,7 +2007,7 @@ describe("daemon worker supervisor monitoring", () => { }; await expect(supervisor.reclaimStaleWorkerRegistration(worker)).resolves.toBe(true); - expect(recoverUncertainWorkerOperations).toHaveBeenCalledWith(worker, false); + expect(recoverUncertainWorkerOperations).toHaveBeenCalledWith(worker); expect(deleteWorkerDescriptor).toHaveBeenCalledWith(worker); expect(workers.has(worker.descriptor.workerId)).toBe(false); }); @@ -1943,40 +2037,156 @@ describe("daemon worker supervisor monitoring", () => { expect(stopWorker).toHaveBeenCalledWith(worker, true, true); }); - it("does not relaunch a live worker whose process identity is unknown", async () => { + it("continues startup recovery after live workers time out during adoption", async () => { + type AdoptionWorker = { + descriptor: { + workerId: string; + pid: number; + processStartId?: string; + rootActiveSessionId: string; + lifecycle?: string; + consecutiveFailures: number; + lastError?: string; + }; + }; + const processStartId = getProcessStartId(process.pid); + expect(processStartId).toBeDefined(); + const workers: AdoptionWorker[] = ["slow-verified", "slow-unverified", "healthy"].map((workerId) => ({ + descriptor: { + workerId, + pid: process.pid, + ...(workerId === "slow-unverified" ? {} : { processStartId: processStartId! }), + rootActiveSessionId: `${workerId}-active`, + consecutiveFailures: 0, + }, + })); + const pendingRecovery = new Promise(() => {}); + const recoverWorker = vi.fn(() => pendingRecovery); + const persistWorker = vi.fn(); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + assertRecoveryAllowed: vi.fn(async () => {}), + connectWorker: vi.fn(async () => {}), + subscribeWorker: vi.fn(async () => {}), + refreshWorkerSummaries: vi.fn(async (worker: AdoptionWorker) => { + if (worker.descriptor.workerId.startsWith("slow")) { + throw new DaemonWorkerProbeTimeoutError("Timed out waiting for daemon worker response to list"); + } + }), + recoverWorker, + persistWorker, + broadcastHeartbeatsChanged: vi.fn(), + log: vi.fn(), + }) as { + adoptOrRecoverWorker(worker: AdoptionWorker): Promise; + }; + + await expect( + Promise.all(workers.map((worker) => supervisor.adoptOrRecoverWorker(worker))), + ).resolves.toBeDefined(); + + expect(recoverWorker).toHaveBeenCalledTimes(2); + expect(workers.map((worker) => worker.descriptor.lifecycle)).toEqual(["recovering", "recovering", "ready"]); + expect(persistWorker).toHaveBeenCalledTimes(3); + }); + + it("preserves authentication rejection outside the probe-timeout recovery path", async () => { + const processStartId = getProcessStartId(process.pid); + expect(processStartId).toBeDefined(); + const worker = { + descriptor: { + workerId: "worker-auth-rejected", + pid: process.pid, + processStartId: processStartId!, + rootActiveSessionId: "active-auth", + consecutiveFailures: 0, + socketPath: "/tmp/worker-auth-rejected.sock", + authenticationToken: "stale-token", + }, + client: undefined, + }; + const authError = new DaemonWorkerAuthenticationError( + "Timed out connecting to daemon session worker: invalid token", + ); + const connect = vi.spyOn(DaemonWorkerClient.prototype, "connect").mockResolvedValue(undefined); + const hello = vi.spyOn(DaemonWorkerClient.prototype, "waitForHello").mockResolvedValue({} as never); + const authenticate = vi.spyOn(DaemonWorkerClient.prototype, "authenticateWorker").mockRejectedValue(authError); + const close = vi.spyOn(DaemonWorkerClient.prototype, "close").mockImplementation(() => undefined); + const recoverWorker = vi.fn(async () => undefined); + const persistWorker = vi.fn(); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + assertRecoveryAllowed: vi.fn(async () => undefined), + supervisorAuthenticationClaim: vi.fn(() => ({ + supervisorGeneration: "generation", + supervisorPid: process.pid, + supervisorSocketPath: "/tmp/supervisor.sock", + })), + recoverWorker, + persistWorker, + log: vi.fn(), + }) as { + adoptOrRecoverWorker(target: typeof worker): Promise; + }; + + try { + await supervisor.adoptOrRecoverWorker(worker); + + expect(authenticate).toHaveBeenCalledOnce(); + expect(recoverWorker).toHaveBeenCalledWith(worker); + expect(persistWorker).not.toHaveBeenCalled(); + expect(worker.descriptor).not.toHaveProperty("lifecycle", "recovering"); + } finally { + connect.mockRestore(); + hello.mockRestore(); + authenticate.mockRestore(); + close.mockRestore(); + } + }); + + it("parks an unresponsive worker failed after the bounded probe rounds", () => { + const worker = { + descriptor: { workerId: "worker-stuck", pid: process.pid, rootActiveSessionId: "active-1" }, + deferredRecoveryRounds: 10, + }; + const persistWorker = vi.fn(); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + workers: new Map([[worker.descriptor.workerId, worker]]), + persistWorker, + markWorkerRosterEntries: vi.fn(), + log: vi.fn(), + }) as { deferWorkerRecovery(target: typeof worker, error: Error): void }; + + supervisor.deferWorkerRecovery(worker, new Error("still silent")); + + expect(worker.descriptor).toMatchObject({ lifecycle: "failed" }); + expect((worker as { deferredRecovery?: unknown }).deferredRecovery).toBeUndefined(); + expect(persistWorker).toHaveBeenCalled(); + }); + + it.each([ + { name: "verified", hasProcessIdentity: true, error: "Timed out waiting for daemon worker hello" }, + { name: "identity-unavailable", hasProcessIdentity: false, error: "worker socket unavailable" }, + ])("defers recovery without replacing a live $name worker", async ({ hasProcessIdentity, error }) => { vi.useFakeTimers(); + const processStartId = hasProcessIdentity ? getProcessStartId(process.pid) : undefined; + if (hasProcessIdentity) expect(processStartId).toBeDefined(); 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; - recoverUncertainWorkerOperations: ReturnType; - launchWorker: ReturnType; - persistWorker: ReturnType; - broadcastHeartbeatsChanged: ReturnType; - log: ReturnType; - assertRecoveryAllowed: ReturnType; - recoverWorker(worker: RecoveryWorker): Promise; }; const worker: RecoveryWorker = { descriptor: { - workerId: "worker-unknown-identity", + workerId: `worker-${hasProcessIdentity ? "verified" : "unknown"}-identity`, pid: process.pid, + ...(processStartId ? { processStartId } : {}), rootActiveSessionId: "active-1", createCommand: { type: "create" }, consecutiveFailures: 0, @@ -1988,23 +2198,32 @@ describe("daemon worker supervisor monitoring", () => { workers: new Map([[worker.descriptor.workerId, worker]]), shuttingDown: false, connectWorker: vi.fn(async () => { - throw new Error("worker socket unavailable"); + throw hasProcessIdentity ? new DaemonWorkerProbeTimeoutError(error) : new Error(error); }), recoverUncertainWorkerOperations: vi.fn(async () => {}), launchWorker: vi.fn(async () => worker), persistWorker: vi.fn(), + broadcastHeartbeatsChanged: vi.fn(), + deferWorkerRecovery: vi.fn(), log: vi.fn(), assertRecoveryAllowed: vi.fn(async () => {}), - }) as RecoveryHarness; + }) as { + recoverWorker(target: RecoveryWorker): Promise; + connectWorker: ReturnType; + recoverUncertainWorkerOperations: ReturnType; + launchWorker: ReturnType; + deferWorkerRecovery: ReturnType; + }; const recovery = supervisor.recoverWorker(worker); - await vi.runAllTimersAsync(); + await vi.advanceTimersByTimeAsync(6250); await recovery; expect(supervisor.connectWorker).toHaveBeenCalledTimes(3); expect(supervisor.recoverUncertainWorkerOperations).not.toHaveBeenCalled(); expect(supervisor.launchWorker).not.toHaveBeenCalled(); - expect(worker.descriptor.lifecycle).toBe("failed"); + expect(supervisor.deferWorkerRecovery).toHaveBeenCalledWith(worker, expect.any(Error)); + expect(worker.descriptor.lifecycle).toBe("recovering"); }); it("reports a stop-tombstoned worker as stopping, not ready", () => { @@ -3770,15 +3989,17 @@ describe("daemon worker supervisor monitoring", () => { const markInterrupted = vi.fn(async () => undefined); const kill = vi.spyOn(process, "kill").mockReturnValue(true); const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { - catalog: { markInterrupted }, + workers: new Map([[worker.descriptor.workerId, worker]]), + shuttingDown: false, + catalog: { start: vi.fn(async () => undefined), markInterrupted }, log: vi.fn(), assertRecoveryAllowed: vi.fn(async () => {}), }) as { - recoverUncertainWorkerOperations(worker: RecoveryWorker, killWorkerProcess: boolean): Promise; + recoverUncertainWorkerOperations(worker: RecoveryWorker): Promise; }; try { - await supervisor.recoverUncertainWorkerOperations(worker, false); + await supervisor.recoverUncertainWorkerOperations(worker); expect(kill).not.toHaveBeenCalled(); expect(markInterrupted).toHaveBeenCalledTimes(2); expect(markInterrupted).toHaveBeenCalledWith("/tmp/root.jsonl", "root-active", ["model_stream"]); @@ -3789,34 +4010,153 @@ describe("daemon worker supervisor monitoring", () => { } }); - it("skips the recovery SIGKILL when the pid identity is no longer current", async () => { - const kill = vi.spyOn(process, "kill").mockReturnValue(true); - const makeSupervisorFixture = (identity: "current" | "replaced") => - Object.assign(Object.create(DaemonSupervisor.prototype), { - log: vi.fn(), - assertRecoveryAllowed: vi.fn(async () => {}), - // The verdict a caller computed before awaiting its way in can go stale; - // the signal must re-check identity at the last moment (PIDs recycle). - processIdentity: vi.fn(() => identity), - }) as { - recoverUncertainWorkerOperations( - worker: { descriptor: { workerId: string; pid: number; recoveryJournalPath: string } }, - killWorkerProcess: boolean, - ): Promise; - }; + it.each([ + { name: "identity-bearing", hasProcessIdentity: true, retained: true }, + { name: "PID-only", hasProcessIdentity: false, retained: false }, + ])("retains only a $name orphan journal after a failed reap", async ({ hasProcessIdentity, retained }) => { + const root = mkdtempSync(join(tmpdir(), "prime-supervisor-orphan-retry-test-")); + const orphanJournalPath = join(root, "worker.orphans.jsonl"); + const workerPid = 987_653; + const orphanPid = process.pid; + const processStartId = hasProcessIdentity ? getProcessStartId(orphanPid) : undefined; + if (hasProcessIdentity) expect(processStartId).toBeDefined(); + writeFileSync( + orphanJournalPath, + `${JSON.stringify({ + version: 1, + pid: orphanPid, + ownerPid: workerPid, + ...(processStartId ? { processStartId } : {}), + active: true, + recordedAt: new Date().toISOString(), + })} +`, + ); + const worker = { + descriptor: { + workerId: "worker-orphan-retry", + pid: workerPid, + rootActiveSessionId: "root-active", + recoveryJournalPath: join(root, "worker.recovery.jsonl"), + orphanProcessJournalPath: orphanJournalPath, + }, + }; + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + workers: new Map([[worker.descriptor.workerId, worker]]), + shuttingDown: false, + catalog: { start: vi.fn(async () => undefined), markInterrupted: vi.fn(async () => undefined) }, + log: vi.fn(), + assertRecoveryAllowed: vi.fn(async () => undefined), + }) as { + recoverUncertainWorkerOperations(target: typeof worker): Promise; + }; + const kill = vi.spyOn(orphanProcessModule, "killOrphanProcess").mockReturnValue(false); + + try { + await supervisor.recoverUncertainWorkerOperations(worker); + + if (hasProcessIdentity || process.platform !== "win32") { + expect(kill).toHaveBeenCalledWith(orphanPid); + } else { + expect(kill).not.toHaveBeenCalled(); + } + expect(existsSync(orphanJournalPath)).toBe(retained); + } finally { + kill.mockRestore(); + rmSync(root, { recursive: true, force: true }); + } + }); + + it("skips catalog startup when recovery has no interrupted operations", async () => { + const root = mkdtempSync(join(tmpdir(), "prime-supervisor-empty-recovery-test-")); const worker = { - descriptor: { workerId: "worker-1", pid: 987_654, recoveryJournalPath: join(tmpdir(), "absent.jsonl") }, + descriptor: { + workerId: "worker-empty-recovery", + pid: 987_654, + rootActiveSessionId: "root-active", + recoveryJournalPath: join(root, "worker.recovery.jsonl"), + }, + }; + const catalogStart = vi.fn(async () => { + throw new Error("catalog unavailable"); + }); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + workers: new Map([[worker.descriptor.workerId, worker]]), + shuttingDown: false, + catalog: { start: catalogStart, markInterrupted: vi.fn() }, + assertRecoveryAllowed: vi.fn(async () => undefined), + }) as { + recoverUncertainWorkerOperations(target: typeof worker): Promise; }; try { - await makeSupervisorFixture("replaced").recoverUncertainWorkerOperations(worker, true); + await expect(supervisor.recoverUncertainWorkerOperations(worker)).resolves.toBeUndefined(); + expect(catalogStart).not.toHaveBeenCalled(); + } finally { + rmSync(root, { recursive: true, force: true }); + } + }); + + it("does not reap interrupted worker resources before the catalog is ready", async () => { + const root = mkdtempSync(join(tmpdir(), "prime-supervisor-catalog-readiness-test-")); + const recoveryJournalPath = join(root, "worker.recovery.jsonl"); + const orphanJournalPath = join(root, "worker.orphans.jsonl"); + const workerPid = 987_654; + new WorkerRecoveryJournal(recoveryJournalPath).record({ + activeSessionId: "root-active", + sessionId: "root-session", + sessionFile: "/tmp/root.jsonl", + busy: true, + operation: "model_stream", + }); + writeFileSync( + orphanJournalPath, + `${JSON.stringify({ + version: 1, + pid: process.pid, + ownerPid: workerPid, + processStartId: getProcessStartId(process.pid), + active: true, + recordedAt: new Date().toISOString(), + })} +`, + ); + const worker = { + descriptor: { + workerId: "worker-catalog-blocked", + pid: workerPid, + rootActiveSessionId: "root-active", + recoveryJournalPath, + orphanProcessJournalPath: orphanJournalPath, + }, + intentionalStop: false, + }; + const catalogError = new Error("Timed out starting daemon catalog"); + const kill = vi.spyOn(orphanProcessModule, "killOrphanProcess").mockReturnValue(true); + const supervisor = Object.assign(Object.create(DaemonSupervisor.prototype), { + workers: new Map([[worker.descriptor.workerId, worker]]), + shuttingDown: false, + catalog: { + start: vi.fn(async () => { + throw catalogError; + }), + markInterrupted: vi.fn(), + }, + assertRecoveryAllowed: vi.fn(async () => {}), + }) as { + recoverUncertainWorkerOperations(target: typeof worker): Promise; + }; + + try { + await expect(supervisor.recoverUncertainWorkerOperations(worker)).rejects.toThrow(catalogError); expect(kill).not.toHaveBeenCalled(); - await makeSupervisorFixture("current").recoverUncertainWorkerOperations(worker, true); - expect(kill).toHaveBeenCalled(); + expect(existsSync(orphanJournalPath)).toBe(true); } finally { kill.mockRestore(); + rmSync(root, { recursive: true, force: true }); } }); + it.each([ { name: "malformed data", data: undefined, error: /invalid update manifest/ }, { diff --git a/packages/coding-agent/test/proper-lockfile-compromise.test.ts b/packages/coding-agent/test/proper-lockfile-compromise.test.ts new file mode 100644 index 0000000000..ae57fe805a --- /dev/null +++ b/packages/coding-agent/test/proper-lockfile-compromise.test.ts @@ -0,0 +1,148 @@ +import { existsSync, mkdtempSync, readdirSync, readFileSync, rmSync } from "node:fs"; +import { createServer } from "node:net"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const lockState = vi.hoisted(() => ({ compromiseAsync: false, compromiseSync: false })); + +type LockOptions = { onCompromised?: (error: Error) => void }; +const releaseAsync = vi.fn(async () => {}); +const releaseSync = vi.fn(() => {}); + +vi.mock("proper-lockfile", () => { + const lock = vi.fn(async (_path: string, options?: LockOptions) => { + if (lockState.compromiseAsync) options?.onCompromised?.(new Error("async lock compromised")); + return releaseAsync; + }); + const lockSync = vi.fn((_path: string, options?: LockOptions) => { + if (lockState.compromiseSync) options?.onCompromised?.(new Error("sync lock compromised")); + return releaseSync; + }); + return { default: { lock, lockSync }, lock, lockSync }; +}); + +import { acquireDaemonUpdateRestartCoordinator } from "../src/cli/daemon-update-restart.js"; +import { FileAuthStorageBackend } from "../src/core/auth-storage.js"; +import { AgentCronJobStore } from "../src/core/cron-jobs.js"; +import { acquireSessionLease, SESSION_LEASES_ENABLED_ENV } from "../src/core/session-lease.js"; +import { FileSettingsStorage } from "../src/core/settings-manager.js"; +import { + acquireDaemonSocketPathLease, + cleanupDaemonSocketPath, + prepareDaemonSocketPath, +} from "../src/modes/daemon/daemon-socket.js"; +import { acquireDaemonSupervisorOwnership } from "../src/modes/daemon/daemon-supervisor-ownership.js"; + +const tempDirs: string[] = []; + +beforeEach(() => { + lockState.compromiseAsync = false; + lockState.compromiseSync = false; + releaseAsync.mockClear(); + releaseSync.mockClear(); +}); + +afterEach(() => { + for (const dir of tempDirs.splice(0)) rmSync(dir, { recursive: true, force: true }); +}); + +function tempDir(prefix: string): string { + const dir = mkdtempSync(join(tmpdir(), prefix)); + tempDirs.push(dir); + return dir; +} + +describe("proper-lockfile compromise boundaries", () => { + it("records a socket lease compromise and fails before preparing the socket", async () => { + lockState.compromiseAsync = true; + const socketPath = join(tempDir("pa-lock-socket-"), "daemon.sock"); + const lease = await acquireDaemonSocketPathLease(socketPath); + + expect(lease?.compromise?.message).toBe("async lock compromised"); + await expect(prepareDaemonSocketPath(socketPath, lease)).rejects.toThrow(/was compromised/); + }); + + it.skipIf(process.platform === "win32")( + "does not unlink a socket when the cleanup guard is compromised", + async () => { + lockState.compromiseSync = true; + const socketPath = join(tempDir("pa-lock-cleanup-"), "daemon.sock"); + const server = createServer(); + try { + await new Promise((resolve) => server.listen(socketPath, resolve)); + cleanupDaemonSocketPath(socketPath); + expect(existsSync(socketPath)).toBe(true); + } finally { + await new Promise((resolve) => server.close(() => resolve())); + } + }, + ); + + it("fails the supervisor ownership mutation closed", async () => { + lockState.compromiseAsync = true; + const root = tempDir("pa-lock-owner-"); + const registryDir = join(root, "registry"); + await expect( + acquireDaemonSupervisorOwnership({ + socketPath: join(root, "daemon.sock"), + descriptorDir: join(root, "descriptors"), + agentDir: join(root, "agent"), + generation: "compromised-owner", + appVersion: "test", + registryDir, + }), + ).rejects.toThrow(/registry guard was compromised/); + expect(existsSync(join(registryDir, "compromised-owner.owner"))).toBe(false); + }); + + it("fails the update coordinator mutation closed", async () => { + lockState.compromiseAsync = true; + const root = tempDir("pa-lock-update-"); + const registryDir = join(root, "registry"); + await expect( + acquireDaemonUpdateRestartCoordinator({ + requestId: "request-1", + socketPath: join(root, "daemon.sock"), + statusPath: join(root, "status.json"), + registryDir, + }), + ).rejects.toThrow(/Coordinator registry guard was compromised/); + expect(readdirSync(registryDir).filter((name) => name.endsWith(".json"))).toEqual([]); + }); + + it("does not create a session lease after its guard is compromised", () => { + lockState.compromiseSync = true; + const agentDir = tempDir("pa-lock-session-"); + expect(() => + acquireSessionLease(join(agentDir, "session.jsonl"), agentDir, { [SESSION_LEASES_ENABLED_ENV]: "1" }), + ).toThrow(/Session lease guard was compromised/); + expect(readdirSync(join(agentDir, "session-leases"))).toEqual([]); + }); + + it("does not mutate auth storage after its lock is compromised", () => { + lockState.compromiseSync = true; + const authPath = join(tempDir("pa-lock-auth-"), "auth.json"); + const storage = new FileAuthStorageBackend(authPath); + expect(() => storage.withLock(() => ({ result: undefined, next: "mutated" }))).toThrow(/lock compromised/); + expect(readFileSync(authPath, "utf8")).toBe("{}"); + }); + + it("does not mutate settings after its lock is compromised", () => { + lockState.compromiseSync = true; + const root = tempDir("pa-lock-settings-"); + const storage = new FileSettingsStorage(root, join(root, "agent")); + expect(() => storage.withLock("global", () => "mutated")).toThrow(/lock compromised/); + expect(existsSync(join(root, "agent", "settings.json"))).toBe(false); + }); + + it("does not mutate scheduled jobs after a lock compromise", () => { + lockState.compromiseSync = true; + const root = tempDir("pa-lock-cron-"); + const artifactDir = join(root, "artifact"); + const store = AgentCronJobStore.forSessionArtifacts(); + store.registerSessionArtifact("session-1", artifactDir); + expect(() => store.recoverSessionArtifact("session-1")).toThrow(/Cron jobs lock compromised/); + expect(existsSync(join(artifactDir, "scheduled-jobs.json"))).toBe(false); + }); +});