From a0cbfa5c014078a8132c07eaa09029bd9ab2048b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Oskar=20Kwas=CC=81niewski?= Date: Mon, 28 Sep 2026 21:23:08 +0200 Subject: [PATCH 1/6] fix(daemon): probe a live daemon again before replacing it as unreachable A 500 ms connect probe measures wall-clock time on the client's event loop. A client that stalls past it (large sync parse, GC on a loaded CI host) reads a listening daemon as unreachable, and the takeover kills it, ending every live session. The next command then meets no session: 2 devices match equally, Daemon request timed out, Invalid daemon response. While the recorded daemon process is still ours, probe up to three more times before replacing it. --- .../daemon-client-stalled-probe.test.ts | 166 ++++++++++++++++++ src/daemon-client/daemon-client-lifecycle.ts | 31 +++- 2 files changed, 196 insertions(+), 1 deletion(-) create mode 100644 src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts diff --git a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts new file mode 100644 index 0000000000..a50c62070a --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts @@ -0,0 +1,166 @@ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import net from 'node:net'; +import path from 'node:path'; +import { test, vi } from 'vitest'; +import { runCmdBackground } from '@agent-device/host-kit/command'; +import { computeDaemonCodeSignature } from '@agent-device/host-kit/code-signature'; +import { + isProcessAlive, + readProcessCommand, + readProcessStartTime, + waitForProcessExit, +} from '@agent-device/host-kit/process'; +import { findProjectRoot, readVersion } from '@agent-device/host-kit/version'; +import { resolveDaemonPaths } from '../../daemon-resolution.ts'; +import { sendToDaemon } from '../daemon-client.ts'; +import { + closeLoopbackServer, + listenOnLoopback, + supportsLoopbackBind, +} from '../../__tests__/test-utils/loopback.ts'; +import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; + +// The spawned stand-in's identity is read once and pinned: a second real `ps` can miss its +// deadline under suite load and misclassify the live process as gone. +const { mockReadProcessStartTime, mockReadProcessCommand } = vi.hoisted(() => ({ + mockReadProcessStartTime: vi.fn<(pid: number) => string | null | undefined>(), + mockReadProcessCommand: vi.fn<(pid: number) => string | null | undefined>(), +})); + +vi.mock('@agent-device/host-kit/process', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + readProcessStartTime: (pid: number) => + mockReadProcessStartTime(pid) ?? actual.readProcessStartTime(pid), + readProcessCommand: (pid: number) => + mockReadProcessCommand(pid) ?? actual.readProcessCommand(pid), + }; +}); + +function resolveCurrentDaemonCodeSignature(): string { + const root = findProjectRoot(); + const distPath = path.join(root, 'dist', 'src', 'internal', 'daemon.js'); + const sourcePath = path.join(root, 'src', 'daemon.ts'); + const entryPath = + process.execArgv.includes('--experimental-strip-types') || !fs.existsSync(distPath) + ? sourcePath + : distPath; + return computeDaemonCodeSignature(entryPath, root); +} + +/** Blocks the event loop right after the first socket arms its timeout, as a loaded client does. */ +function stallAfterFirstProbeArms(): { stalled: () => boolean; restore: () => void } { + const originalCreateConnection = net.createConnection; + let stalled = false; + (net as unknown as { createConnection: typeof net.createConnection }).createConnection = (( + ...args: Parameters + ) => { + const socket = originalCreateConnection(...args); + if (stalled) return socket; + stalled = true; + const armTimeout = socket.setTimeout.bind(socket); + socket.setTimeout = ((timeoutMs: number) => { + armTimeout(timeoutMs); + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 700); + return socket; + }) as typeof socket.setTimeout; + return socket; + }) as typeof net.createConnection; + return { + stalled: () => stalled, + restore: () => { + (net as unknown as { createConnection: typeof net.createConnection }).createConnection = + originalCreateConnection; + }, + }; +} + +test('sendToDaemon keeps a live socket daemon whose probe the client stalled past', async (t) => { + if (!(await supportsLoopbackBind())) { + t.skip('loopback listeners are not permitted in this environment'); + return; + } + + const stateDir = mkdtempForTestSync('agent-device-stalled-probe-'); + const root = mkdtempForTestSync('agent-device-stalled-probe-daemon-'); + const daemonDir = path.join(root, 'agent-device', 'dist', 'src', 'internal'); + const daemonScriptPath = path.join(daemonDir, 'daemon.js'); + fs.mkdirSync(daemonDir, { recursive: true }); + fs.writeFileSync(daemonScriptPath, 'setInterval(() => {}, 1000);\n', 'utf8'); + const daemonProcess = runCmdBackground(process.execPath, [daemonScriptPath], { + stdio: 'ignore', + allowFailure: true, + captureOutput: false, + }); + void daemonProcess.wait.catch(() => {}); + const pid = daemonProcess.child.pid; + assert.ok(pid, 'spawned child should have a pid'); + const server = net.createServer((socket) => { + let requestBody = ''; + socket.setEncoding('utf8'); + socket.on('data', (chunk) => { + requestBody += chunk; + if (!requestBody.includes('\n')) return; + socket.end(`${JSON.stringify({ ok: true, data: { via: 'live-daemon' } })}\n`); + }); + }); + let stall: ReturnType | undefined; + + try { + await new Promise((resolve) => setTimeout(resolve, 50)); + const processStartTime = readProcessStartTime(pid) ?? undefined; + const command = readProcessCommand(pid); + if (command === null || processStartTime === undefined) { + t.skip('process command/start inspection is unavailable in this environment'); + return; + } + mockReadProcessStartTime.mockImplementation((queriedPid: number) => + queriedPid === pid ? processStartTime : undefined, + ); + mockReadProcessCommand.mockImplementation((queriedPid: number) => + queriedPid === pid ? command : undefined, + ); + const port = await listenOnLoopback(server); + const paths = resolveDaemonPaths(stateDir); + fs.mkdirSync(paths.baseDir, { recursive: true }); + fs.writeFileSync( + paths.infoPath, + `${JSON.stringify({ + port, + transport: 'socket', + token: 'local-secret', + pid, + version: readVersion(), + codeSignature: resolveCurrentDaemonCodeSignature(), + processStartTime, + })}\n`, + 'utf8', + ); + stall = stallAfterFirstProbeArms(); + + const response = await sendToDaemon({ + session: 'default', + command: 'stalled-probe-smoke', + positionals: [], + flags: { stateDir, daemonTransport: 'socket' }, + meta: { requestId: 'req-stalled-probe' }, + }); + + assert.equal(stall.stalled(), true); + assert.deepEqual(response, { ok: true, data: { via: 'live-daemon' } }); + assert.equal(isProcessAlive(pid), true); + } finally { + stall?.restore(); + mockReadProcessStartTime.mockReset(); + mockReadProcessCommand.mockReset(); + await closeLoopbackServer(server); + if (isProcessAlive(pid)) { + process.kill(pid, 'SIGKILL'); + await waitForProcessExit(pid, 1_500); + } + fs.rmSync(stateDir, { recursive: true, force: true }); + fs.rmSync(root, { recursive: true, force: true }); + } +}); diff --git a/src/daemon-client/daemon-client-lifecycle.ts b/src/daemon-client/daemon-client-lifecycle.ts index a9025f1e1b..e96f3daafa 100644 --- a/src/daemon-client/daemon-client-lifecycle.ts +++ b/src/daemon-client/daemon-client-lifecycle.ts @@ -11,6 +11,7 @@ import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; import { readProcessStartTime } from '@agent-device/host-kit/process'; import { sleep } from '@agent-device/host-kit/retry'; +import { isAgentDeviceDaemonProcess } from '../daemon-process.ts'; import { findUnrecoveredRepairCommitFailure } from '../session-repair-tombstone.ts'; import { resolveDaemonPaths, @@ -76,6 +77,8 @@ type DaemonStartupWaitResult = | { kind: 'timeout' }; const DAEMON_STARTUP_TIMEOUT_MS = 15_000; +const LIVE_DAEMON_PROBE_RETRIES = 3; +const LIVE_DAEMON_PROBE_RETRY_DELAY_MS = 200; const DAEMON_STARTUP_ATTEMPTS = 2; const DAEMON_STARTUP_LOG_TAIL_BYTES = 64_000; const LOOPBACK_BLOCK_LIST = new net.BlockList(); @@ -192,7 +195,7 @@ async function readReusableLocalDaemon(settings: DaemonClientSettings): Promise< const existing = readDaemonInfo(settings.paths.infoPath); if (!existing) return null; - const viaClientTransport = await canConnectReusableDaemon(existing, settings.transportPreference); + const viaClientTransport = await canReachReusableDaemon(existing, settings.transportPreference); const decision = await resolveDaemonTakeover(existing, { viaClientTransport, onAnyAdvertisedTransport: async () => @@ -209,6 +212,32 @@ async function readReusableLocalDaemon(settings: DaemonClientSettings): Promise< return null; } +/** + * A daemon whose process is still the one its metadata recorded is probed again before it can be + * judged unreachable. A probe's budget is wall-clock time on this client's event loop, so a client + * that stalls past it (a large synchronous parse, a GC pause on a loaded host) reads a listening + * daemon as unreachable, and replacing it ends every session the daemon holds. + */ +async function canReachReusableDaemon( + info: DaemonInfo, + preference: DaemonTransportPreference, +): Promise { + if (await canConnectReusableDaemon(info, preference)) return true; + for (let retry = 1; retry <= LIVE_DAEMON_PROBE_RETRIES; retry += 1) { + if (!isAgentDeviceDaemonProcess(info.pid, info.processStartTime)) return false; + await sleep(LIVE_DAEMON_PROBE_RETRY_DELAY_MS); + if (await canConnectReusableDaemon(info, preference)) { + emitDiagnostic({ + level: 'warn', + phase: 'daemon_probe_recovered', + data: { pid: info.pid, retry }, + }); + return true; + } + } + return false; +} + async function canConnectReusableDaemon( info: DaemonInfo, preference: DaemonTransportPreference, From c3fb83054af82a09c84c97ac6c6ff91dc47d3434 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Oskar=20Kwas=CC=81niewski?= Date: Tue, 29 Sep 2026 01:00:44 +0200 Subject: [PATCH 2/6] test(daemon): cover the probe retry with a deterministic missed probe On Linux the stalled-client test's first probe still connects, so the retry never ran under coverage. A mocked miss exercises it on every host; the stall test skips with the reason where the host outruns it. --- .../daemon-client-stalled-probe.test.ts | 89 +++++++++++++++---- 1 file changed, 72 insertions(+), 17 deletions(-) diff --git a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts index a50c62070a..ec6a8574b7 100644 --- a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts +++ b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts @@ -28,6 +28,25 @@ const { mockReadProcessStartTime, mockReadProcessCommand } = vi.hoisted(() => ({ mockReadProcessCommand: vi.fn<(pid: number) => string | null | undefined>(), })); +// Every probe's answer, in order; a test can also make the next one miss. +const { probeAnswers, mockMissNextProbe } = vi.hoisted(() => ({ + probeAnswers: [] as boolean[], + mockMissNextProbe: { value: false }, +})); + +vi.mock('../daemon-client-transport.ts', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + canConnect: async (...args: Parameters) => { + const reachable = mockMissNextProbe.value ? false : await actual.canConnect(...args); + mockMissNextProbe.value = false; + probeAnswers.push(reachable); + return reachable; + }, + }; +}); + vi.mock('@agent-device/host-kit/process', async (importOriginal) => { const actual = await importOriginal(); return { @@ -77,12 +96,20 @@ function stallAfterFirstProbeArms(): { stalled: () => boolean; restore: () => vo }; } -test('sendToDaemon keeps a live socket daemon whose probe the client stalled past', async (t) => { +type LiveStandIn = { stateDir: string; pid: number }; + +/** + * Runs `body` against a daemon.json naming a live process that reads as an agent-device daemon + * and a loopback socket that answers every request, as a running daemon does. + */ +async function withLiveStandIn( + t: { skip: (reason: string) => void }, + body: (standIn: LiveStandIn) => Promise, +): Promise { if (!(await supportsLoopbackBind())) { t.skip('loopback listeners are not permitted in this environment'); return; } - const stateDir = mkdtempForTestSync('agent-device-stalled-probe-'); const root = mkdtempForTestSync('agent-device-stalled-probe-daemon-'); const daemonDir = path.join(root, 'agent-device', 'dist', 'src', 'internal'); @@ -106,7 +133,6 @@ test('sendToDaemon keeps a live socket daemon whose probe the client stalled pas socket.end(`${JSON.stringify({ ok: true, data: { via: 'live-daemon' } })}\n`); }); }); - let stall: ReturnType | undefined; try { await new Promise((resolve) => setTimeout(resolve, 50)); @@ -138,21 +164,10 @@ test('sendToDaemon keeps a live socket daemon whose probe the client stalled pas })}\n`, 'utf8', ); - stall = stallAfterFirstProbeArms(); - - const response = await sendToDaemon({ - session: 'default', - command: 'stalled-probe-smoke', - positionals: [], - flags: { stateDir, daemonTransport: 'socket' }, - meta: { requestId: 'req-stalled-probe' }, - }); - - assert.equal(stall.stalled(), true); - assert.deepEqual(response, { ok: true, data: { via: 'live-daemon' } }); - assert.equal(isProcessAlive(pid), true); + probeAnswers.length = 0; + await body({ stateDir, pid }); } finally { - stall?.restore(); + mockMissNextProbe.value = false; mockReadProcessStartTime.mockReset(); mockReadProcessCommand.mockReset(); await closeLoopbackServer(server); @@ -163,4 +178,44 @@ test('sendToDaemon keeps a live socket daemon whose probe the client stalled pas fs.rmSync(stateDir, { recursive: true, force: true }); fs.rmSync(root, { recursive: true, force: true }); } +} + +function sendSmoke(stateDir: string) { + return sendToDaemon({ + session: 'default', + command: 'stalled-probe-smoke', + positionals: [], + flags: { stateDir, daemonTransport: 'socket' }, + meta: { requestId: 'req-stalled-probe' }, + }); +} + +test('sendToDaemon keeps a live daemon whose first probe missed', async (t) => { + await withLiveStandIn(t, async ({ stateDir, pid }) => { + mockMissNextProbe.value = true; + + const response = await sendSmoke(stateDir); + + assert.deepEqual(probeAnswers.slice(0, 2), [false, true]); + assert.deepEqual(response, { ok: true, data: { via: 'live-daemon' } }); + assert.equal(isProcessAlive(pid), true); + }); +}); + +test('sendToDaemon keeps a live socket daemon whose probe the client stalled past', async (t) => { + await withLiveStandIn(t, async ({ stateDir, pid }) => { + const stall = stallAfterFirstProbeArms(); + try { + const response = await sendSmoke(stateDir); + if (probeAnswers[0] !== false) { + t.skip('this host completed the connect before the stalled timer fired'); + return; + } + assert.equal(stall.stalled(), true); + assert.deepEqual(response, { ok: true, data: { via: 'live-daemon' } }); + assert.equal(isProcessAlive(pid), true); + } finally { + stall.restore(); + } + }); }); From fd4cc19756a472a87d4b1724dad07d355c16ce08 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Oskar=20Kwas=CC=81niewski?= Date: Tue, 29 Sep 2026 02:18:53 +0200 Subject: [PATCH 3/6] fix(daemon): gate the probe retry on liveness, not the ps identity read Under the load that stalls the probe, ps misses its deadline too and the identity read fails closed, ending the retries. The takeover still proves identity before it signals. --- src/daemon-client/daemon-client-lifecycle.ts | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/src/daemon-client/daemon-client-lifecycle.ts b/src/daemon-client/daemon-client-lifecycle.ts index e96f3daafa..f44c3278d4 100644 --- a/src/daemon-client/daemon-client-lifecycle.ts +++ b/src/daemon-client/daemon-client-lifecycle.ts @@ -8,10 +8,9 @@ import type { DaemonRequest, DaemonResponse } from '../daemon/daemon-request.ts' import { runCmdDetachedMonitored, type ExecDetachedExit } from '@agent-device/host-kit/command'; import { shellQuoteIfNeeded } from '@agent-device/kernel/device-shell'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; -import { readProcessStartTime } from '@agent-device/host-kit/process'; +import { isProcessAlive, readProcessStartTime } from '@agent-device/host-kit/process'; import { sleep } from '@agent-device/host-kit/retry'; -import { isAgentDeviceDaemonProcess } from '../daemon-process.ts'; import { findUnrecoveredRepairCommitFailure } from '../session-repair-tombstone.ts'; import { resolveDaemonPaths, @@ -213,10 +212,12 @@ async function readReusableLocalDaemon(settings: DaemonClientSettings): Promise< } /** - * A daemon whose process is still the one its metadata recorded is probed again before it can be - * judged unreachable. A probe's budget is wall-clock time on this client's event loop, so a client - * that stalls past it (a large synchronous parse, a GC pause on a loaded host) reads a listening - * daemon as unreachable, and replacing it ends every session the daemon holds. + * A daemon whose pid is still alive is probed again before it can be judged unreachable. A probe's + * budget is wall-clock time on this client's event loop, so a client that stalls past it (a large + * synchronous parse, a GC pause on a loaded host) reads a listening daemon as unreachable, and + * replacing it ends every session the daemon holds. Liveness is the signal-0 check, not the `ps` + * identity read: under the load that stalls the probe, `ps` misses its deadline too, and the + * takeover still proves identity before it signals anything. */ async function canReachReusableDaemon( info: DaemonInfo, @@ -224,7 +225,7 @@ async function canReachReusableDaemon( ): Promise { if (await canConnectReusableDaemon(info, preference)) return true; for (let retry = 1; retry <= LIVE_DAEMON_PROBE_RETRIES; retry += 1) { - if (!isAgentDeviceDaemonProcess(info.pid, info.processStartTime)) return false; + if (!isProcessAlive(info.pid)) return false; await sleep(LIVE_DAEMON_PROBE_RETRY_DELAY_MS); if (await canConnectReusableDaemon(info, preference)) { emitDiagnostic({ From 3646f1bfdedbe99074bf98256e567ca5f2f68a6d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Oskar=20Kwas=CC=81niewski?= Date: Tue, 29 Sep 2026 10:58:02 +0200 Subject: [PATCH 4/6] fix(daemon): ask client-transport reachability only when it decides the takeover A version or code mismatch decides alone, so the patient probe no longer waits on a daemon about to be replaced anyway. Pins a dead pid's single probe and the recovery diagnostic, and drops the stall test that patched the global createConnection and skipped where the host outran it. --- .../__tests__/daemon-client-lifecycle.test.ts | 2 +- .../daemon-client-stalled-probe.test.ts | 149 ++++++++++++------ .../__tests__/daemon-launch-spec.test.ts | 20 ++- src/daemon-client/daemon-client-lifecycle.ts | 6 +- src/daemon-client/daemon-launch-spec.ts | 5 +- 5 files changed, 125 insertions(+), 57 deletions(-) diff --git a/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts b/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts index 073f789a58..47fdc26b26 100644 --- a/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts +++ b/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts @@ -501,7 +501,7 @@ test('sendToDaemon does not reuse reachable daemon metadata with mismatched vers assert.deepEqual(response, { ok: true, data: { via: 'fresh-daemon' } }); assert.equal(mockRunCmdDetached.mock.calls.length, 1); - assert.deepEqual(staleDaemon.seenPaths, ['GET /health']); + assert.deepEqual(staleDaemon.seenPaths, []); assert.deepEqual(freshDaemon.seenPaths, ['GET /health', 'POST /rpc']); const staleVersion = fixture.version ?? readVersion(); assert.equal( diff --git a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts index ec6a8574b7..de200604a8 100644 --- a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts +++ b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts @@ -28,11 +28,31 @@ const { mockReadProcessStartTime, mockReadProcessCommand } = vi.hoisted(() => ({ mockReadProcessCommand: vi.fn<(pid: number) => string | null | undefined>(), })); -// Every probe's answer, in order; a test can also make the next one miss. -const { probeAnswers, mockMissNextProbe } = vi.hoisted(() => ({ - probeAnswers: [] as boolean[], - mockMissNextProbe: { value: false }, -})); +// Every probe's answer, in order, with the port it asked; a test can also make the next one miss. +const { probeAnswers, probedPorts, mockMissNextProbe, mockEmitDiagnostic, mockSpawnDaemon } = + vi.hoisted(() => ({ + probeAnswers: [] as boolean[], + probedPorts: [] as (number | undefined)[], + mockMissNextProbe: { value: false }, + mockEmitDiagnostic: vi.fn(), + mockSpawnDaemon: vi.fn(), + })); + +vi.mock('@agent-device/host-kit/diagnostics', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + emitDiagnostic: (...args: Parameters) => { + mockEmitDiagnostic(...args); + actual.emitDiagnostic(...args); + }, + }; +}); + +vi.mock('@agent-device/host-kit/command', async (importOriginal) => { + const actual = await importOriginal(); + return { ...actual, runCmdDetachedMonitored: mockSpawnDaemon }; +}); vi.mock('../daemon-client-transport.ts', async (importOriginal) => { const actual = await importOriginal(); @@ -42,6 +62,7 @@ vi.mock('../daemon-client-transport.ts', async (importOriginal) => { const reachable = mockMissNextProbe.value ? false : await actual.canConnect(...args); mockMissNextProbe.value = false; probeAnswers.push(reachable); + probedPorts.push(args[0].port); return reachable; }, }; @@ -69,33 +90,6 @@ function resolveCurrentDaemonCodeSignature(): string { return computeDaemonCodeSignature(entryPath, root); } -/** Blocks the event loop right after the first socket arms its timeout, as a loaded client does. */ -function stallAfterFirstProbeArms(): { stalled: () => boolean; restore: () => void } { - const originalCreateConnection = net.createConnection; - let stalled = false; - (net as unknown as { createConnection: typeof net.createConnection }).createConnection = (( - ...args: Parameters - ) => { - const socket = originalCreateConnection(...args); - if (stalled) return socket; - stalled = true; - const armTimeout = socket.setTimeout.bind(socket); - socket.setTimeout = ((timeoutMs: number) => { - armTimeout(timeoutMs); - Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 700); - return socket; - }) as typeof socket.setTimeout; - return socket; - }) as typeof net.createConnection; - return { - stalled: () => stalled, - restore: () => { - (net as unknown as { createConnection: typeof net.createConnection }).createConnection = - originalCreateConnection; - }, - }; -} - type LiveStandIn = { stateDir: string; pid: number }; /** @@ -165,6 +159,8 @@ async function withLiveStandIn( 'utf8', ); probeAnswers.length = 0; + probedPorts.length = 0; + mockEmitDiagnostic.mockClear(); await body({ stateDir, pid }); } finally { mockMissNextProbe.value = false; @@ -199,23 +195,82 @@ test('sendToDaemon keeps a live daemon whose first probe missed', async (t) => { assert.deepEqual(probeAnswers.slice(0, 2), [false, true]); assert.deepEqual(response, { ok: true, data: { via: 'live-daemon' } }); assert.equal(isProcessAlive(pid), true); + assert.ok( + mockEmitDiagnostic.mock.calls.some( + ([event]) => event.phase === 'daemon_probe_recovered' && event.data?.pid === pid, + ), + 'a recovered probe names the daemon it kept', + ); }); }); -test('sendToDaemon keeps a live socket daemon whose probe the client stalled past', async (t) => { - await withLiveStandIn(t, async ({ stateDir, pid }) => { - const stall = stallAfterFirstProbeArms(); - try { - const response = await sendSmoke(stateDir); - if (probeAnswers[0] !== false) { - t.skip('this host completed the connect before the stalled timer fired'); - return; - } - assert.equal(stall.stalled(), true); - assert.deepEqual(response, { ok: true, data: { via: 'live-daemon' } }); - assert.equal(isProcessAlive(pid), true); - } finally { - stall.restore(); - } +test('sendToDaemon replaces a daemon whose process is gone after a single probe', async (t) => { + if (!(await supportsLoopbackBind())) { + t.skip('loopback listeners are not permitted in this environment'); + return; + } + const stateDir = mkdtempForTestSync('agent-device-dead-probe-'); + const gone = runCmdBackground(process.execPath, ['-e', ''], { + stdio: 'ignore', + allowFailure: true, + captureOutput: false, + }); + await gone.wait.catch(() => {}); + const deadPid = gone.child.pid; + assert.ok(deadPid, 'spawned child should have a pid'); + const unused = net.createServer(); + const deadPort = await listenOnLoopback(unused); + await closeLoopbackServer(unused); + const fresh = net.createServer((socket) => { + socket.setEncoding('utf8'); + socket.on('data', () => { + socket.end(`${JSON.stringify({ ok: true, data: { via: 'fresh-daemon' } })}\n`); + }); }); + const writeInfo = (port: number, pid: number) => { + const paths = resolveDaemonPaths(stateDir); + fs.mkdirSync(paths.baseDir, { recursive: true }); + fs.writeFileSync( + paths.infoPath, + `${JSON.stringify({ + port, + transport: 'socket', + token: 'local-secret', + pid, + version: readVersion(), + codeSignature: resolveCurrentDaemonCodeSignature(), + processStartTime: readProcessStartTime(process.pid) ?? undefined, + })}\n`, + 'utf8', + ); + }; + + try { + const freshPort = await listenOnLoopback(fresh); + writeInfo(deadPort, deadPid); + mockSpawnDaemon.mockImplementation(() => { + writeInfo(freshPort, process.pid); + return { pid: process.pid, exited: new Promise(() => {}) }; + }); + probeAnswers.length = 0; + probedPorts.length = 0; + mockEmitDiagnostic.mockClear(); + + const response = await sendSmoke(stateDir); + + assert.deepEqual(response, { ok: true, data: { via: 'fresh-daemon' } }); + assert.equal( + probedPorts.filter((port) => port === deadPort).length, + 1, + 'a daemon whose pid is gone gets no patient retry', + ); + assert.equal( + mockEmitDiagnostic.mock.calls.some(([event]) => event.phase === 'daemon_probe_recovered'), + false, + ); + } finally { + mockSpawnDaemon.mockReset(); + await closeLoopbackServer(fresh); + fs.rmSync(stateDir, { recursive: true, force: true }); + } }); diff --git a/src/daemon-client/__tests__/daemon-launch-spec.test.ts b/src/daemon-client/__tests__/daemon-launch-spec.test.ts index 37af703064..d69d0aa727 100644 --- a/src/daemon-client/__tests__/daemon-launch-spec.test.ts +++ b/src/daemon-client/__tests__/daemon-launch-spec.test.ts @@ -122,6 +122,20 @@ test('an unreachable newer daemon is replaced like any version mismatch', async ); }); +test('a version mismatch is decided without asking the client transport', async () => { + let asked = false; + const decision = await resolveDaemonTakeover(runningDaemon({ version: '0.0.1' }), { + onClientTransport: async () => { + asked = true; + return false; + }, + onAnyAdvertisedTransport: async () => false, + }); + + assert.equal(decision.kind, 'replace'); + assert.equal(asked, false); +}); + test('a newer daemon alive only on a transport the client does not prefer is still refused', async () => { assert.deepEqual( await resolveDaemonTakeover(runningDaemon({ version: '999.0.0' }), onlyOnAnotherTransport()), @@ -159,15 +173,15 @@ function useClientTree(sourceCheckout: boolean): void { } function reachable(): DaemonReachability { - return { viaClientTransport: true, onAnyAdvertisedTransport: async () => true }; + return { onClientTransport: async () => true, onAnyAdvertisedTransport: async () => true }; } function unreachable(): DaemonReachability { - return { viaClientTransport: false, onAnyAdvertisedTransport: async () => false }; + return { onClientTransport: async () => false, onAnyAdvertisedTransport: async () => false }; } function onlyOnAnotherTransport(): DaemonReachability { - return { viaClientTransport: false, onAnyAdvertisedTransport: async () => true }; + return { onClientTransport: async () => false, onAnyAdvertisedTransport: async () => true }; } function runningDaemon(info: { diff --git a/src/daemon-client/daemon-client-lifecycle.ts b/src/daemon-client/daemon-client-lifecycle.ts index f44c3278d4..2e2a727ab1 100644 --- a/src/daemon-client/daemon-client-lifecycle.ts +++ b/src/daemon-client/daemon-client-lifecycle.ts @@ -194,11 +194,9 @@ async function readReusableLocalDaemon(settings: DaemonClientSettings): Promise< const existing = readDaemonInfo(settings.paths.infoPath); if (!existing) return null; - const viaClientTransport = await canReachReusableDaemon(existing, settings.transportPreference); const decision = await resolveDaemonTakeover(existing, { - viaClientTransport, - onAnyAdvertisedTransport: async () => - viaClientTransport || (await canConnectReusableDaemon(existing, 'auto')), + onClientTransport: () => canReachReusableDaemon(existing, settings.transportPreference), + onAnyAdvertisedTransport: () => canConnectReusableDaemon(existing, 'auto'), }); if (decision.kind === 'reuse') return existing; if (decision.kind === 'refuseNewer') { diff --git a/src/daemon-client/daemon-launch-spec.ts b/src/daemon-client/daemon-launch-spec.ts index d0cd527bb3..102cdf3332 100644 --- a/src/daemon-client/daemon-launch-spec.ts +++ b/src/daemon-client/daemon-launch-spec.ts @@ -119,7 +119,8 @@ export type DaemonTakeoverDecision = * transport preference the daemon does not serve must still see a live newer daemon. */ export type DaemonReachability = { - viaClientTransport: boolean; + /** Asked last, only once version and code identity leave the decision to it. */ + onClientTransport: () => Promise; onAnyAdvertisedTransport: () => Promise; }; @@ -152,7 +153,7 @@ export async function resolveDaemonTakeover( const localIdentity = await resolveLocalDaemonCodeIdentity(); const codeMismatch = resolveCodeIdentityMismatch(localIdentity, info); if (codeMismatch) return { kind: 'replace', reason: codeMismatch }; - if (!reachability.viaClientTransport) return { kind: 'replace', reason: 'unreachable' }; + if (!(await reachability.onClientTransport())) return { kind: 'replace', reason: 'unreachable' }; return { kind: 'reuse' }; } From ec33fb2b1ca13a16c4ab396980eb3749ff054fce Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Oskar=20Kwas=CC=81niewski?= Date: Tue, 29 Sep 2026 11:54:03 +0200 Subject: [PATCH 5/6] test(daemon): pin the dead-pid replacement on a port the fresh daemon cannot reuse The fresh stand-in binds before the dead port is picked, and the test asserts the replacement spawn ran and the one dead probe failed; it skips when the host recycled the exited pid. --- .../daemon-client-stalled-probe.test.ts | 20 +++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts index de200604a8..f7c06991e5 100644 --- a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts +++ b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts @@ -218,9 +218,6 @@ test('sendToDaemon replaces a daemon whose process is gone after a single probe' await gone.wait.catch(() => {}); const deadPid = gone.child.pid; assert.ok(deadPid, 'spawned child should have a pid'); - const unused = net.createServer(); - const deadPort = await listenOnLoopback(unused); - await closeLoopbackServer(unused); const fresh = net.createServer((socket) => { socket.setEncoding('utf8'); socket.on('data', () => { @@ -246,7 +243,12 @@ test('sendToDaemon replaces a daemon whose process is gone after a single probe' }; try { + // Bound before the dead port is picked, so the port the dead daemon recorded cannot be + // handed back to the fresh one. const freshPort = await listenOnLoopback(fresh); + const unused = net.createServer(); + const deadPort = await listenOnLoopback(unused); + await closeLoopbackServer(unused); writeInfo(deadPort, deadPid); mockSpawnDaemon.mockImplementation(() => { writeInfo(freshPort, process.pid); @@ -255,15 +257,17 @@ test('sendToDaemon replaces a daemon whose process is gone after a single probe' probeAnswers.length = 0; probedPorts.length = 0; mockEmitDiagnostic.mockClear(); + if (isProcessAlive(deadPid)) { + t.skip('the host recycled the exited stand-in pid before the probe'); + return; + } const response = await sendSmoke(stateDir); assert.deepEqual(response, { ok: true, data: { via: 'fresh-daemon' } }); - assert.equal( - probedPorts.filter((port) => port === deadPort).length, - 1, - 'a daemon whose pid is gone gets no patient retry', - ); + assert.equal(mockSpawnDaemon.mock.calls.length, 1, 'the dead daemon was replaced'); + const deadProbes = probeAnswers.filter((_, index) => probedPorts[index] === deadPort); + assert.deepEqual(deadProbes, [false], 'a daemon whose pid is gone gets no patient retry'); assert.equal( mockEmitDiagnostic.mock.calls.some(([event]) => event.phase === 'daemon_probe_recovered'), false, From 74b15cdf774aab0eea14f410fdc554cb84478dac Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Wed, 30 Sep 2026 12:49:39 +0200 Subject: [PATCH 6/6] fix(daemon): re-probe a live newer daemon before replacing it as unreachable The newer-daemon refusal decided on one bare probe, so a client that stalled through it replaced a live newer daemon as a version mismatch and killed its sessions. Both reachability answers now go through canReachReusableDaemon. --- .../daemon-client-stalled-probe.test.ts | 49 ++++++++++++++++--- src/daemon-client/daemon-client-lifecycle.ts | 2 +- 2 files changed, 43 insertions(+), 8 deletions(-) diff --git a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts index f7c06991e5..6005159793 100644 --- a/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts +++ b/src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts @@ -20,6 +20,9 @@ import { supportsLoopbackBind, } from '../../__tests__/test-utils/loopback.ts'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; +import { AppError } from '@agent-device/kernel/errors'; + +const NEWER_DAEMON_VERSION = '999.0.0'; // The spawned stand-in's identity is read once and pinned: a second real `ps` can miss its // deadline under suite load and misclassify the live process as gone. @@ -98,6 +101,7 @@ type LiveStandIn = { stateDir: string; pid: number }; */ async function withLiveStandIn( t: { skip: (reason: string) => void }, + version: string, body: (standIn: LiveStandIn) => Promise, ): Promise { if (!(await supportsLoopbackBind())) { @@ -152,7 +156,7 @@ async function withLiveStandIn( transport: 'socket', token: 'local-secret', pid, - version: readVersion(), + version, codeSignature: resolveCurrentDaemonCodeSignature(), processStartTime, })}\n`, @@ -164,6 +168,7 @@ async function withLiveStandIn( await body({ stateDir, pid }); } finally { mockMissNextProbe.value = false; + mockSpawnDaemon.mockReset(); mockReadProcessStartTime.mockReset(); mockReadProcessCommand.mockReset(); await closeLoopbackServer(server); @@ -187,7 +192,7 @@ function sendSmoke(stateDir: string) { } test('sendToDaemon keeps a live daemon whose first probe missed', async (t) => { - await withLiveStandIn(t, async ({ stateDir, pid }) => { + await withLiveStandIn(t, readVersion(), async ({ stateDir, pid }) => { mockMissNextProbe.value = true; const response = await sendSmoke(stateDir); @@ -204,7 +209,29 @@ test('sendToDaemon keeps a live daemon whose first probe missed', async (t) => { }); }); -test('sendToDaemon replaces a daemon whose process is gone after a single probe', async (t) => { +test('sendToDaemon refuses a live newer daemon whose first probe missed', async (t) => { + await withLiveStandIn(t, NEWER_DAEMON_VERSION, async ({ stateDir, pid }) => { + mockMissNextProbe.value = true; + + await assert.rejects( + () => sendSmoke(stateDir), + (error: unknown) => { + assert.ok(error instanceof AppError); + assert.equal(error.details?.daemonVersion, NEWER_DAEMON_VERSION); + return true; + }, + ); + + assert.deepEqual(probeAnswers.slice(0, 2), [false, true]); + assert.equal(isProcessAlive(pid), true); + assert.equal(mockSpawnDaemon.mock.calls.length, 0, 'no replacement daemon is spawned'); + }); +}); + +async function expectDeadDaemonReplaced( + t: { skip: (reason: string) => void }, + deadVersion: string, +): Promise { if (!(await supportsLoopbackBind())) { t.skip('loopback listeners are not permitted in this environment'); return; @@ -224,7 +251,7 @@ test('sendToDaemon replaces a daemon whose process is gone after a single probe' socket.end(`${JSON.stringify({ ok: true, data: { via: 'fresh-daemon' } })}\n`); }); }); - const writeInfo = (port: number, pid: number) => { + const writeInfo = (port: number, pid: number, version: string) => { const paths = resolveDaemonPaths(stateDir); fs.mkdirSync(paths.baseDir, { recursive: true }); fs.writeFileSync( @@ -234,7 +261,7 @@ test('sendToDaemon replaces a daemon whose process is gone after a single probe' transport: 'socket', token: 'local-secret', pid, - version: readVersion(), + version, codeSignature: resolveCurrentDaemonCodeSignature(), processStartTime: readProcessStartTime(process.pid) ?? undefined, })}\n`, @@ -249,9 +276,9 @@ test('sendToDaemon replaces a daemon whose process is gone after a single probe' const unused = net.createServer(); const deadPort = await listenOnLoopback(unused); await closeLoopbackServer(unused); - writeInfo(deadPort, deadPid); + writeInfo(deadPort, deadPid, deadVersion); mockSpawnDaemon.mockImplementation(() => { - writeInfo(freshPort, process.pid); + writeInfo(freshPort, process.pid, readVersion()); return { pid: process.pid, exited: new Promise(() => {}) }; }); probeAnswers.length = 0; @@ -277,4 +304,12 @@ test('sendToDaemon replaces a daemon whose process is gone after a single probe' await closeLoopbackServer(fresh); fs.rmSync(stateDir, { recursive: true, force: true }); } +} + +test('sendToDaemon replaces a daemon whose process is gone after a single probe', async (t) => { + await expectDeadDaemonReplaced(t, readVersion()); +}); + +test('sendToDaemon replaces a newer daemon whose process is gone after a single probe', async (t) => { + await expectDeadDaemonReplaced(t, NEWER_DAEMON_VERSION); }); diff --git a/src/daemon-client/daemon-client-lifecycle.ts b/src/daemon-client/daemon-client-lifecycle.ts index 98a8761cf9..f8bd690557 100644 --- a/src/daemon-client/daemon-client-lifecycle.ts +++ b/src/daemon-client/daemon-client-lifecycle.ts @@ -203,7 +203,7 @@ async function readReusableLocalDaemon(settings: DaemonClientSettings): Promise< const decision = await resolveDaemonTakeover(existing, { onClientTransport: () => canReachReusableDaemon(existing, settings.transportPreference), - onAnyAdvertisedTransport: () => canConnectReusableDaemon(existing, 'auto'), + onAnyAdvertisedTransport: () => canReachReusableDaemon(existing, 'auto'), }); if (decision.kind === 'reuse') return existing; if (decision.kind === 'refuseNewer') {