Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
280 changes: 280 additions & 0 deletions src/daemon-client/__tests__/daemon-client-stalled-probe.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,280 @@
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>(),
}));

// 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<typeof import('@agent-device/host-kit/diagnostics')>();
return {
...actual,
emitDiagnostic: (...args: Parameters<typeof actual.emitDiagnostic>) => {
mockEmitDiagnostic(...args);
actual.emitDiagnostic(...args);
},
};
});

vi.mock('@agent-device/host-kit/command', async (importOriginal) => {
const actual = await importOriginal<typeof import('@agent-device/host-kit/command')>();
return { ...actual, runCmdDetachedMonitored: mockSpawnDaemon };
});

vi.mock('../daemon-client-transport.ts', async (importOriginal) => {
const actual = await importOriginal<typeof import('../daemon-client-transport.ts')>();
return {
...actual,
canConnect: async (...args: Parameters<typeof actual.canConnect>) => {
const reachable = mockMissNextProbe.value ? false : await actual.canConnect(...args);
mockMissNextProbe.value = false;
probeAnswers.push(reachable);
probedPorts.push(args[0].port);
return reachable;
},
};
});

vi.mock('@agent-device/host-kit/process', async (importOriginal) => {
const actual = await importOriginal<typeof import('@agent-device/host-kit/process')>();
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);
}

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<void>,
): Promise<void> {
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`);
});
});

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',
);
probeAnswers.length = 0;
probedPorts.length = 0;
mockEmitDiagnostic.mockClear();
await body({ stateDir, pid });
} finally {
mockMissNextProbe.value = false;
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 });
}
}

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);
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 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', ''], {
Comment thread
okwasniewski marked this conversation as resolved.
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 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 {
// 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);
return { pid: process.pid, exited: new Promise(() => {}) };
});
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(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'),
Comment thread
okwasniewski marked this conversation as resolved.
false,
);
} finally {
mockSpawnDaemon.mockReset();
await closeLoopbackServer(fresh);
fs.rmSync(stateDir, { recursive: true, force: true });
}
});
20 changes: 17 additions & 3 deletions src/daemon-client/__tests__/daemon-launch-spec.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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()),
Expand Down Expand Up @@ -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: {
Expand Down
37 changes: 33 additions & 4 deletions src/daemon-client/daemon-client-lifecycle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ 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 { isProcessAlive } from '@agent-device/host-kit/process';
import { sleep } from '@agent-device/host-kit/retry';

import { findUnrecoveredRepairCommitFailure } from '../session-repair-tombstone.ts';
Expand Down Expand Up @@ -72,6 +73,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();
Expand Down Expand Up @@ -185,11 +188,9 @@ async function readReusableLocalDaemon(settings: DaemonClientSettings): Promise<
const existing = readDaemonInfo(settings.paths.infoPath);
if (!existing) return null;

const viaClientTransport = await canConnectReusableDaemon(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') {
Expand All @@ -202,6 +203,34 @@ async function readReusableLocalDaemon(settings: DaemonClientSettings): Promise<
return null;
}

/**
* 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,
preference: DaemonTransportPreference,
): Promise<boolean> {
if (await canConnectReusableDaemon(info, preference)) return true;
for (let retry = 1; retry <= LIVE_DAEMON_PROBE_RETRIES; retry += 1) {
if (!isProcessAlive(info.pid)) return false;
await sleep(LIVE_DAEMON_PROBE_RETRY_DELAY_MS);
Comment thread
okwasniewski marked this conversation as resolved.
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,
Expand Down
Loading
Loading