diff --git a/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts b/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts index 88d57bbc14..487df065d2 100644 --- a/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts +++ b/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts @@ -27,6 +27,8 @@ import { resolveDaemonPaths, type DaemonPaths } from '../../daemon-resolution.ts import { sendToDaemon, type DaemonRequest, type DaemonResponse } from '../daemon-client.ts'; import { attachActiveSessionAddressHint } from '../daemon-client-lifecycle.ts'; import { sendRequest } from '../daemon-client-transport.ts'; +import { readDaemonInfo } from '../daemon-client-metadata.ts'; +import type { DaemonRetirementResult } from '../../daemon-registration-owner.ts'; import { closeLoopbackServer, listenOnLoopback, @@ -615,12 +617,18 @@ test('sendRequest timeout cleanup uses resolved daemon paths instead of request const daemonPaths = resolveDaemonPaths(daemonStateDir); const requestFlagPaths = resolveDaemonPaths(requestFlagStateDir); const daemon = await startHangingHttpDaemonFixture(); - writeDaemonInfo(daemonPaths, { - httpPort: daemon.port, - transport: 'http', - pid: 999_999, - }); - writeDaemonLock(daemonPaths, { pid: 999_999 }); + mockSleep.mockImplementation(actualRetry.sleep); + const child = spawnRegisteredDaemonFixture( + daemonPaths, + { + httpPort: daemon.port, + token: 'local-secret', + version: readVersion(), + codeOrigin: 'checkout', + codeSignature: currentDaemonCodeSignature(), + }, + undefined, + ); writeDaemonInfo(requestFlagPaths, { httpPort: daemon.port, transport: 'http', @@ -638,26 +646,26 @@ test('sendRequest timeout cleanup uses resolved daemon paths instead of request }; try { + let info = readDaemonInfo(daemonPaths.infoPath); + for (let attempt = 0; !info && attempt < 200; attempt += 1) { + await actualRetry.sleep(10); + info = readDaemonInfo(daemonPaths.infoPath); + } + assert.ok(info); let thrown: unknown; try { - await sendRequest( - { - token: 'local-secret', - pid: 999_999, - httpPort: daemon.port, - transport: 'http', - }, - request, - 'http', - daemonPaths, - 50, - ); + await sendRequest(info, request, 'http', daemonPaths, 50); } catch (error) { thrown = error; } assert.ok(thrown instanceof AppError); assert.equal(thrown.message, 'Daemon request timed out'); + assert.equal( + (thrown.details?.retirement as DaemonRetirementResult | undefined)?.status, + 'retired', + ); + await child.exited; assert.deepEqual(daemon.seenPaths, ['POST /rpc']); assert.equal(fs.existsSync(daemonPaths.infoPath), false); assert.equal(fs.existsSync(daemonPaths.lockPath), false); @@ -665,7 +673,7 @@ test('sendRequest timeout cleanup uses resolved daemon paths instead of request assert.equal(fs.existsSync(requestFlagPaths.lockPath), true); } finally { await closeLoopbackServer(daemon.server); - fs.rmSync(daemonStateDir, { recursive: true, force: true }); + await finishRegisteredDaemonFixture(daemonStateDir); fs.rmSync(requestFlagStateDir, { recursive: true, force: true }); } }); diff --git a/src/daemon-client/__tests__/daemon-client-metadata.test.ts b/src/daemon-client/__tests__/daemon-client-metadata.test.ts index 6bfec46b6e..bae4084103 100644 --- a/src/daemon-client/__tests__/daemon-client-metadata.test.ts +++ b/src/daemon-client/__tests__/daemon-client-metadata.test.ts @@ -1,29 +1,13 @@ import assert from 'node:assert/strict'; -import { AppError, normalizeError } from '@agent-device/kernel/errors'; import fs from 'node:fs'; import path from 'node:path'; -import { afterEach, test, vi } from 'vitest'; +import { test } from 'vitest'; import type { DaemonCodeOrigin } from '@agent-device/host-kit/code-signature'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; -import { - stopAndRetireDaemon, - tryAcquireDaemonRegistration, -} from '../../daemon-registration-owner.ts'; -import { - readDaemonInfo, - stopDaemonProcessForTakeover, - type DaemonInfo, -} from '../daemon-client-metadata.ts'; -import { isAgentDeviceDaemonProcess, stopDaemonProcess } from '../../daemon-process.ts'; +import { tryAcquireDaemonRegistration } from '../../daemon-registration-owner.ts'; +import { readDaemonInfo, type DaemonInfo } from '../daemon-client-metadata.ts'; import { resolveDaemonPaths } from '../../daemon-resolution.ts'; -vi.mock('../../daemon-process.ts', async (importOriginal) => ({ - ...(await importOriginal()), - isAgentDeviceDaemonProcess: vi.fn(), - stopDaemonProcess: vi.fn(), -})); -afterEach(() => vi.resetAllMocks()); - // The reuse decision is only as good as the identity that survives the round trip // through `daemon.json`: a client cannot compare what the file lost (#2458). @@ -68,41 +52,3 @@ test('a registration this version did not write reads back unreported', () => { assert.equal(readDaemonInfo(infoPath)?.codeOrigin, undefined); } }); - -for (const artifact of ['daemon.json', 'daemon.lock']) { - test(`unconfirmed startup stop retains ${artifact} without claiming cleanup`, async () => { - const [stateDir] = scratchStateDir(); - const paths = resolveDaemonPaths(stateDir); - const file = path.join(stateDir, artifact); - const contents = JSON.stringify({ - pid: 7, - processStartTime: 'start', - port: 1234, - token: 'secret', - }); - fs.writeFileSync(file, contents); - vi.mocked(isAgentDeviceDaemonProcess).mockReturnValue(true); - vi.mocked(stopDaemonProcess).mockResolvedValue({ status: 'retained', reason: 'exit-timeout' }); - const result = await stopAndRetireDaemon({ - paths, - observed: { pid: 7, startTime: 'start' }, - mode: 'graceful', - }); - assert.equal(fs.readFileSync(file, 'utf8'), contents); - assert.equal(result.removedInfo, false); - assert.equal(result.status, 'retained'); - if (result.status === 'retained') assert.equal(result.reason, 'exit-unconfirmed'); - }); -} - -test('a retained takeover keeps its reason at the normalized error boundary', async () => { - vi.mocked(stopDaemonProcess).mockResolvedValue({ status: 'retained', reason: 'exit-timeout' }); - await assert.rejects( - stopDaemonProcessForTakeover({ pid: 7, token: 'secret', processStartTime: 'start' }), - (error: unknown) => { - assert.ok(error instanceof AppError); - assert.equal(normalizeError(error).details?.reason, 'daemon_exit_unconfirmed'); - return true; - }, - ); -}); diff --git a/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts b/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts index f02582c40e..04d4154403 100644 --- a/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts +++ b/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts @@ -26,19 +26,11 @@ import net from 'node:net'; import http from 'node:http'; import path from 'node:path'; +import fs from 'node:fs'; import assert from 'node:assert/strict'; import { beforeEach, afterEach, test, vi } from 'vitest'; -const { mockRunCmdSync, mockIsDaemon, mockStop } = vi.hoisted(() => ({ - mockRunCmdSync: vi.fn(), - mockIsDaemon: vi.fn(), - mockStop: vi.fn(), -})); -vi.mock('../../daemon-process.ts', async (importOriginal) => ({ - ...(await importOriginal()), - isAgentDeviceDaemonProcess: mockIsDaemon, - stopDaemonProcess: mockStop, -})); +const { mockRunCmdSync } = vi.hoisted(() => ({ mockRunCmdSync: vi.fn() })); vi.mock('@agent-device/host-kit/command', async () => { const actual = await vi.importActual( @@ -47,18 +39,29 @@ vi.mock('@agent-device/host-kit/command', async () => { return { ...actual, runCmdSync: mockRunCmdSync }; }); -import { AppError } from '@agent-device/kernel/errors'; +import { AppError, normalizeError } from '@agent-device/kernel/errors'; +import { sleep } from '@agent-device/host-kit/retry'; import { sendRequest } from '../daemon-client-transport.ts'; import type { DaemonRequest } from '../../daemon/daemon-request.ts'; -import type { DaemonInfo } from '../daemon-client-metadata.ts'; -import type { DaemonPaths } from '../../daemon-resolution.ts'; +import { readDaemonInfo, type DaemonInfo } from '../daemon-client-metadata.ts'; +import { resolveDaemonPaths, type DaemonPaths } from '../../daemon-resolution.ts'; +import type { DaemonRetirementResult } from '../../daemon-registration-owner.ts'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; +import { + spawnRegisteredDaemonFixture, + finishRegisteredDaemonFixture, + finishRegisteredDaemonFixtures, +} from '../../__tests__/test-utils/registered-daemon-fixture.ts'; +import { + closeLoopbackServer, + skipWhenLoopbackUnavailable, +} from '../../__tests__/test-utils/loopback.ts'; const TIMEOUT_MS = 120; // `snapshot`'s timeout policy preserves the daemon (onTimeout !== -// 'reset-daemon'), so `handleRequestTimeout` never reaches -// `resetDaemonAfterTimeout` (`process.kill`) here — keeping this suite +// 'reset-daemon'), so `handleRequestTimeout` never signals the daemon +// in the hint controls — keeping them // side-effect-free outside the mocked pkill sweep. const SNAPSHOT_COMMAND = 'snapshot'; @@ -94,6 +97,7 @@ function startHangingSocketServer(): Promise<{ server: net.Server; port: number // Accept the connection but never write a response — forces the // client's own request-timeout envelope to fire. socket.on('error', () => {}); + socket.resume(); }); server.on('error', reject); server.listen(0, '127.0.0.1', () => { @@ -129,10 +133,11 @@ function startHangingHttpServer(): Promise<{ server: http.Server; port: number } beforeEach(() => { mockRunCmdSync.mockReset(); - mockIsDaemon.mockReset(); - mockStop.mockReset(); }); -afterEach(() => vi.restoreAllMocks()); +afterEach(async () => { + vi.restoreAllMocks(); + await finishRegisteredDaemonFixtures(); +}); test('socket timeout: pkill cleanup still runs for a declared non-Apple platform that actually terminates a runner (rebound-session case), and the hint claims Apple on that evidence', async () => { // Simulates --session-lock strip silently rebinding this request onto an @@ -279,32 +284,156 @@ test('remote HTTP timeout never runs the Apple pkill cleanup and uses the remote assert.equal(mockRunCmdSync.mock.calls.length, 0); }); -test('a refused timeout fallback preserves the timeout without an unhandled rejection', async () => { - mockRunCmdSync.mockReturnValue({ exitCode: 1, stdout: '', stderr: '' }); - mockIsDaemon.mockReturnValue(true); - mockStop.mockResolvedValue({ status: 'retained', reason: 'exit-timeout' }); - vi.spyOn(process, 'kill').mockImplementation(() => { - throw Object.assign(new Error('refused'), { code: 'EPERM' }); +async function publishedInfo(paths: DaemonPaths): Promise { + for (let attempt = 0; attempt < 200; attempt += 1) { + const info = readDaemonInfo(paths.infoPath); + if (info) return info; + await sleep(10); + } + throw new Error('registered child did not publish'); +} + +for (const transport of ['socket', 'http'] as const) { + test(`${transport} timeout waits for force retirement before reporting reset`, async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + mockRunCmdSync.mockReturnValue({ exitCode: 1, stdout: '', stderr: '' }); + const endpoint = await (transport === 'socket' + ? startHangingSocketServer() + : startHangingHttpServer()); + const paths = resolveDaemonPaths(mkdtempForTestSync('agent-device-timeout-owner-')); + const child = spawnRegisteredDaemonFixture( + paths, + { + ...(transport === 'socket' ? { socketPort: endpoint.port } : { httpPort: endpoint.port }), + token: 'test-token', + version: 'test', + codeOrigin: 'checkout', + codeSignature: 'test', + }, + undefined, + ); + let exited = false; + void child.exited.then(() => { + exited = true; + }); + let killRequested!: () => void; + const requested = new Promise((resolve) => { + killRequested = resolve; + }); + const actualKill = process.kill.bind(process); + let settled = false; + let outcome: Promise | undefined; + const kill = vi.spyOn(process, 'kill').mockImplementation((pid, signal) => { + if (pid === child.pid && signal === 'SIGKILL') { + killRequested(); + return true; + } + return actualKill(pid, signal); + }); + try { + const info = await publishedInfo(paths); + outcome = sendRequest( + info, + { ...buildRequest(undefined), command: 'open' }, + transport, + paths, + TIMEOUT_MS, + ).then( + () => assert.fail('hanging request unexpectedly succeeded'), + (error: unknown) => { + settled = true; + return error; + }, + ); + await Promise.race([ + requested, + sleep(1_500).then(() => assert.fail('force stop was not requested')), + ]); + await sleep(30); + assert.equal(actualKill(child.pid, 0), true); + assert.equal(settled, false, 'request must remain pending while the daemon is alive'); + kill.mockRestore(); + actualKill(child.pid, 'SIGKILL'); + await child.exited; + const error = await outcome; + assert.ok(error instanceof AppError); + assert.equal(normalizeError(error).details?.reason, 'daemon_transport_timeout'); + const retirement = error.details?.retirement as DaemonRetirementResult | undefined; + assert.ok(retirement?.status === 'retired'); + assert.equal(retirement.termination.mode, 'forced'); + assert.equal(fs.existsSync(paths.infoPath), false); + assert.equal(fs.existsSync(paths.lockPath), false); + assert.equal(fs.existsSync(paths.baseDir), true); + assert.equal(mockRunCmdSync.mock.calls.filter(([cmd]) => cmd === 'pkill').length, 3); + } finally { + kill.mockRestore(); + if (!exited) actualKill(child.pid, 'SIGKILL'); + await child.exited; + await outcome; + await finishRegisteredDaemonFixture(paths.baseDir); + await closeLoopbackServer(endpoint.server); + } }); - const { server, port } = await startHangingSocketServer(); +} + +test('timeout retains a live registration without captured birth proof and reports that outcome', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + mockRunCmdSync.mockReturnValue({ exitCode: 1, stdout: '', stderr: '' }); + const endpoint = await startHangingHttpServer(); + const paths = resolveDaemonPaths(mkdtempForTestSync('agent-device-timeout-retained-')); + const child = spawnRegisteredDaemonFixture( + paths, + { + httpPort: endpoint.port, + token: 'test-token', + version: 'test', + codeOrigin: 'checkout', + codeSignature: 'test', + }, + undefined, + ); + const kill = vi.spyOn(process, 'kill'); try { + const info = await publishedInfo(paths); + const before = fs.readFileSync(paths.infoPath, 'utf8'); + const lockBefore = fs + .readdirSync(paths.lockPath) + .map((name) => [name, fs.readFileSync(path.join(paths.lockPath, name), 'utf8')]); await assert.rejects( sendRequest( - { port, pid: 7, token: 'test-token', processStartTime: 'start' }, + { ...info, processStartTime: undefined }, { ...buildRequest(undefined), command: 'open' }, - 'socket', - dummyStatePaths(), + 'http', + paths, TIMEOUT_MS, ), (error: unknown) => { assert.ok(error instanceof AppError); - assert.equal(error.details?.reason, 'daemon_transport_timeout'); + assert.equal(normalizeError(error).details?.reason, 'daemon_transport_timeout'); + const retirement = error.details?.retirement as DaemonRetirementResult | undefined; + assert.ok(retirement?.status === 'retained'); + assert.ok(retirement.termination?.status === 'retained'); + assert.equal(retirement.termination.reason, 'missing-start-time'); + assert.match(normalizeError(error).hint ?? '', /State was retained/); + assert.doesNotMatch(normalizeError(error).hint ?? '', /daemon was reset/); return true; }, ); - await new Promise((resolve) => setImmediate(resolve)); - assert.equal(mockStop.mock.calls.length, 1); + assert.equal(process.kill(child.pid, 0), true); + assert.equal(fs.readFileSync(paths.infoPath, 'utf8'), before); + assert.deepEqual( + fs + .readdirSync(paths.lockPath) + .map((name) => [name, fs.readFileSync(path.join(paths.lockPath, name), 'utf8')]), + lockBefore, + ); + assert.equal( + kill.mock.calls.some(([pid, signal]) => pid === child.pid && signal !== 0), + false, + ); } finally { - server.close(); + kill.mockRestore(); + await finishRegisteredDaemonFixture(paths.baseDir); + await closeLoopbackServer(endpoint.server); } }); diff --git a/src/daemon-client/__tests__/daemon-client-transport.test.ts b/src/daemon-client/__tests__/daemon-client-transport.test.ts index d9b60a8bd8..d9860c241a 100644 --- a/src/daemon-client/__tests__/daemon-client-transport.test.ts +++ b/src/daemon-client/__tests__/daemon-client-transport.test.ts @@ -1,6 +1,9 @@ import assert from 'node:assert/strict'; import http from 'node:http'; +import net from 'node:net'; import { test, vi } from 'vitest'; +import * as hostTransport from '@agent-device/host-kit/transport'; +import { sleep } from '@agent-device/host-kit/retry'; import { AppError } from '@agent-device/kernel/errors'; import { DAEMON_HTTP_INSTANCE_HEADER, @@ -36,6 +39,64 @@ function sendWithStaleInstance(port: number, timeoutMs: number) { ); } +test('auto health probing reserves time for a healthy fallback when HTTP hangs', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + const httpServer = http.createServer(() => {}); + const socketServer = net.createServer((socket) => socket.on('error', () => {})); + try { + const httpPort = await listenOnLoopback(httpServer); + const port = await listenOnLoopback(socketServer); + assert.equal( + await canConnect({ token: 'secret', pid: 1, transport: 'http', httpPort, port }, 'auto', 120), + true, + ); + } finally { + await closeLoopbackServer(httpServer); + await closeLoopbackServer(socketServer); + } +}); + +test('the health deadline includes requester loading and forbids a late request', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + let requests = 0; + const server = http.createServer((_req, res) => { + requests += 1; + res.end('{}'); + }); + let release!: () => void; + const blocked = new Promise((resolve) => { + release = resolve; + }); + const actualLoad = hostTransport.loadNodeHttpRequester; + const load = vi + .spyOn(hostTransport, 'loadNodeHttpRequester') + .mockImplementation(async (protocol) => { + await blocked; + return actualLoad(protocol); + }); + let probing: Promise | undefined; + try { + const httpPort = await listenOnLoopback(server); + let settled = false; + probing = canConnect({ token: 'secret', pid: 1, httpPort }, 'http', 40).then((reachable) => { + settled = true; + return reachable; + }); + await sleep(90); + assert.equal(settled, true, 'loading must not extend the probe deadline'); + assert.equal(await probing, false); + release(); + await blocked; + await sleep(10); + assert.equal(requests, 0, 'a timed-out loader must not open a request later'); + } finally { + release(); + await probing; + load.mockRestore(); + await closeLoopbackServer(server); + } +}); + test('persistent remote client caches health and retries a refused stale instance before dispatch', async (t) => { if (await skipWhenLoopbackUnavailable(t)) return; const paths: string[] = []; diff --git a/src/daemon-client/daemon-client-metadata.ts b/src/daemon-client/daemon-client-metadata.ts index 5e9d060daa..d517da18e4 100644 --- a/src/daemon-client/daemon-client-metadata.ts +++ b/src/daemon-client/daemon-client-metadata.ts @@ -1,6 +1,4 @@ import fs from 'node:fs'; -import { AppError } from '@agent-device/kernel/errors'; -import { stopDaemonProcess, type DaemonTerminationResult } from '../daemon-process.ts'; import type { DaemonCodeOrigin } from '@agent-device/host-kit/code-signature'; @@ -32,9 +30,6 @@ export type DaemonMetadataState = { hasLock: boolean; }; -const DAEMON_TAKEOVER_TERM_TIMEOUT_MS = 3000; -const DAEMON_TAKEOVER_KILL_TIMEOUT_MS = 1000; - export function readDaemonInfo(infoPath: string): DaemonInfo | null { const data = readJsonFile(infoPath); if (!data || typeof data !== 'object') return null; @@ -85,14 +80,6 @@ function readPositiveInteger(value: unknown): number | undefined { return Number.isInteger(value) && Number(value) > 0 ? Number(value) : undefined; } -export function removeDaemonInfo(infoPath: string): void { - removeFileIfExists(infoPath); -} - -export function removeDaemonLock(lockPath: string): void { - removeFileIfExists(lockPath); -} - export function getDaemonMetadataState(paths: DaemonPaths): DaemonMetadataState { return { hasInfo: fs.existsSync(paths.infoPath), @@ -100,29 +87,6 @@ export function getDaemonMetadataState(paths: DaemonPaths): DaemonMetadataState }; } -export async function stopDaemonProcessForTakeover( - info: DaemonInfo, -): Promise { - const termination = await stopDaemonProcess( - { pid: info.pid, startTime: info.processStartTime ?? null }, - { - mode: 'graceful', - termTimeoutMs: DAEMON_TAKEOVER_TERM_TIMEOUT_MS, - killTimeoutMs: DAEMON_TAKEOVER_KILL_TIMEOUT_MS, - }, - ); - requireDaemonExit(termination); - return termination; -} - -function requireDaemonExit(termination: DaemonTerminationResult): void { - if (termination.status !== 'retained') return; - throw new AppError('COMMAND_FAILED', 'Daemon exit could not be confirmed.', { - reason: 'daemon_exit_unconfirmed', - termination, - }); -} - export function isRemoteDaemon(info: DaemonInfo): boolean { return typeof info.baseUrl === 'string' && info.baseUrl.length > 0; } @@ -147,11 +111,3 @@ function readJsonFile(filePath: string): unknown | null { return null; } } - -function removeFileIfExists(filePath: string): void { - try { - if (fs.existsSync(filePath)) fs.unlinkSync(filePath); - } catch { - // Best-effort cleanup only. - } -} diff --git a/src/daemon-client/daemon-client-timeout.ts b/src/daemon-client/daemon-client-timeout.ts index d338f74d02..c27c95b1b7 100644 --- a/src/daemon-client/daemon-client-timeout.ts +++ b/src/daemon-client/daemon-client-timeout.ts @@ -1,18 +1,13 @@ -import { AppError, normalizeError } from '@agent-device/kernel/errors'; +import { AppError } from '@agent-device/kernel/errors'; import { runCmdSync } from '@agent-device/host-kit/command'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; -import { isAgentDeviceDaemonProcess } from '../daemon-process.ts'; +import type { DaemonRetirementResult } from '../daemon-registration-owner.ts'; import { PUBLIC_COMMANDS } from '@agent-device/command-registry/catalog'; import { resolveCommandTimeoutPolicy } from '@agent-device/command-registry/registry'; import type { DaemonPaths } from '../daemon-resolution.ts'; import type { PlatformSelector } from '@agent-device/kernel/device'; -import { - removeDaemonInfo, - removeDaemonLock, - stopDaemonProcessForTakeover, - type DaemonInfo, -} from './daemon-client-metadata.ts'; +import type { DaemonInfo } from './daemon-client-metadata.ts'; const IOS_RUNNER_XCODEBUILD_KILL_PATTERNS = [ 'xcodebuild .*AgentDeviceRunnerUITests/RunnerTests/testCommand', @@ -40,7 +35,7 @@ function isAffirmativelyApplePlatform(platform: PlatformSelector | undefined): b return platform !== undefined && AFFIRMATIVE_APPLE_PLATFORM_SELECTORS.has(platform); } -export function handleRequestTimeout( +export async function handleRequestTimeout( params: Readonly<{ info: DaemonInfo; statePaths: DaemonPaths; @@ -53,7 +48,7 @@ export function handleRequestTimeout( session?: string; action?: string; }>, -): AppError { +): Promise { const { info, statePaths, remote, timeoutMs, requestId, command, platform, session, action } = params; // Cleanup eligibility stays UNCONDITIONAL for every local (non-remote) @@ -69,9 +64,19 @@ export function handleRequestTimeout( // a wrong skip. const cleanup = remote ? { terminated: 0 } : cleanupTimedOutIosRunnerBuilds(); const resetDaemon = !remote && shouldResetDaemonAfterRequestTimeout(command); - const daemonReset = resetDaemon - ? resetDaemonAfterTimeout(info, statePaths) - : { forcedKill: false }; + let retirement: DaemonRetirementResult | undefined; + if (resetDaemon) { + const { stopAndRetireDaemon } = await import('../daemon-registration-owner.ts'); + retirement = await stopAndRetireDaemon({ + paths: statePaths, + observed: { pid: info.pid, startTime: info.processStartTime ?? null }, + mode: 'force', + }); + } + const forcedKill = + retirement?.status !== 'absent' && + retirement?.termination?.status === 'exited' && + retirement.termination.mode === 'forced'; // The HINT, unlike cleanup, may only name Apple-runner involvement on // evidence this call site actually has: an explicitly declared Apple // platform selector, or the cleanup itself having terminated a matching @@ -89,8 +94,9 @@ export function handleRequestTimeout( command, timedOutRunnerPidsTerminated: cleanup.terminated, timedOutRunnerCleanupError: cleanup.error, - daemonPidReset: resetDaemon ? info.pid : undefined, - daemonPidForceKilled: resetDaemon ? daemonReset.forcedKill : undefined, + daemonPidReset: retirement?.status === 'retired' ? info.pid : undefined, + daemonPidForceKilled: resetDaemon ? forcedKill : undefined, + daemonRetirement: retirement, daemonPreservedAfterTimeout: !remote && !resetDaemon, daemonBaseUrl: info.baseUrl, }, @@ -99,14 +105,18 @@ export function handleRequestTimeout( timeoutMs, requestId, reason: 'daemon_transport_timeout', - hint: resolveRequestTimeoutHint({ - remote, - resetDaemon, - command, - appleCleanupEvidence, - session, - action, - }), + ...(retirement ? { retirement, stateDir: statePaths.baseDir } : {}), + hint: + retirement?.status === 'retained' + ? `The daemon could not be safely retired. State was retained at ${statePaths.baseDir}. ${retirement.error?.hint ?? 'Retry with --debug and inspect daemon diagnostics before retrying.'}` + : resolveRequestTimeoutHint({ + remote, + resetDaemon, + command, + appleCleanupEvidence, + session, + action, + }), }); } @@ -177,25 +187,3 @@ function cleanupTimedOutIosRunnerBuilds(): { terminated: number; error?: string }; } } - -function resetDaemonAfterTimeout(info: DaemonInfo, paths: DaemonPaths): { forcedKill: boolean } { - let forcedKill = false; - try { - if (isAgentDeviceDaemonProcess(info.pid, info.processStartTime)) { - process.kill(info.pid, 'SIGKILL'); - forcedKill = true; - } - } catch { - void stopDaemonProcessForTakeover(info).catch((error: unknown) => { - emitDiagnostic({ - level: 'warn', - phase: 'daemon_timeout_stop_failed', - data: { error: normalizeError(error) }, - }); - }); - } finally { - removeDaemonInfo(paths.infoPath); - removeDaemonLock(paths.lockPath); - } - return { forcedKill }; -} diff --git a/src/daemon-client/daemon-client-transport.ts b/src/daemon-client/daemon-client-transport.ts index aa0f93820d..66d545b0a3 100644 --- a/src/daemon-client/daemon-client-transport.ts +++ b/src/daemon-client/daemon-client-transport.ts @@ -76,10 +76,13 @@ export async function canConnect( ): Promise { const deadline = Date.now() + (probeTimeoutMs ?? Number.POSITIVE_INFINITY); const transport = chooseTransport(info, preference); - if (await canConnectWithTransport(info, transport, deadline - Date.now())) return true; - const fallback = chooseAutoFallbackTransport(info, preference, transport); - return fallback ? await canConnectWithTransport(info, fallback, deadline - Date.now()) : false; + const firstBudget = (deadline - Date.now()) / (fallback ? 2 : 1); + if (await canConnectWithTransport(info, transport, firstBudget)) return Date.now() < deadline; + return fallback + ? (await canConnectWithTransport(info, fallback, deadline - Date.now())) && + Date.now() < deadline + : false; } async function canConnectWithTransport( @@ -160,17 +163,27 @@ async function readDaemonHttpHealth( : null; if (!endpoint) return { reachable: false }; const url = new URL(endpoint); - const transport = await loadNodeHttpRequester(url.protocol); const timeoutMs = Math.min( info.baseUrl ? REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS : LOCAL_DAEMON_HEALTHCHECK_TIMEOUT_MS, probeTimeoutMs ?? Number.POSITIVE_INFINITY, ); if (timeoutMs <= 0) return { reachable: false, timedOut: true }; + const deadline = performance.now() + timeoutMs; const signal = AbortSignal.timeout(Math.ceil(timeoutMs)); + const transport = await Promise.race([ + loadNodeHttpRequester(url.protocol), + new Promise((resolve) => { + signal.addEventListener('abort', () => resolve(null), { once: true }); + }), + ]); + if (!transport || healthProbeExpired(signal, deadline)) + return { reachable: false, timedOut: true }; return await new Promise((resolve) => { const headers = info.baseUrl ? buildDaemonHttpAuthHeaders(info.token) : {}; const unreachable = (): RemoteDaemonHealth => - signal.aborted ? { reachable: false, timedOut: true } : { reachable: false }; + healthProbeExpired(signal, deadline) + ? { reachable: false, timedOut: true } + : { reachable: false }; const req = transport.request( { protocol: url.protocol, @@ -190,11 +203,11 @@ async function readDaemonHttpHealth( }); res.on('end', () => { const statusCode = res.statusCode ?? 500; - resolve({ - reachable: statusCode < 500, - statusCode, - ...readHealthPayload(body), - }); + resolve( + healthProbeExpired(signal, deadline) + ? { reachable: false, timedOut: true } + : { reachable: statusCode < 500, statusCode, ...readHealthPayload(body) }, + ); }); res.on('error', () => resolve(unreachable())); res.on('aborted', () => resolve(unreachable())); @@ -211,6 +224,10 @@ async function readDaemonHttpHealth( }); } +function healthProbeExpired(signal: AbortSignal, deadline: number): boolean { + return signal.aborted || performance.now() >= deadline; +} + function readHealthPayload(body: string): Omit { try { const parsed = JSON.parse(body) as { upstream?: unknown }; @@ -245,6 +262,29 @@ export async function sendRequest( statePaths: DaemonPaths, timeoutMs: number | undefined, options: SendRequestOptions = {}, +): Promise { + try { + return await sendRequestWithFallback(info, req, preference, statePaths, timeoutMs, options); + } catch (error) { + if (!(error instanceof AppError) || error.details?.reason !== 'daemon_transport_timeout') + throw error; + const expiredBudget = error.details.timeoutMs; + if (typeof expiredBudget !== 'number') throw error; + throw await handleRequestTimeout({ + info, + statePaths, + ...timeoutRequestContext(req, isRemoteDaemon(info), expiredBudget), + }); + } +} + +async function sendRequestWithFallback( + info: DaemonInfo, + req: DaemonRequest, + preference: DaemonTransportPreference, + statePaths: DaemonPaths, + timeoutMs: number | undefined, + options: SendRequestOptions = {}, ): Promise { const transport = chooseTransport(info, preference); const deadline = typeof timeoutMs === 'number' ? performance.now() + timeoutMs : undefined; @@ -279,30 +319,14 @@ async function retryAfterRemoteInstanceMismatch( options: SendRequestOptions, ): Promise { invalidateRemoteDaemonHealth(info); - const probeTimeoutMs = remainingRemoteRequestTimeoutMs( - info, + const probeTimeoutMs = remainingRemoteRequestTimeoutMs(req, timeoutMs, deadline); + const health = await readRemoteDaemonHealth(info, probeTimeoutMs); + const remainingMs = remainingRemoteRequestTimeoutMs( req, - statePaths, timeoutMs, deadline, + health.timedOut ? probeTimeoutMs : undefined, ); - const health = await readRemoteDaemonHealth(info, probeTimeoutMs); - // The probe's timer starts from the event loop's cached clock, so it can expire while - // performance.now() is still short of the deadline: a probe the RPC deadline capped that ran out - // of time is the RPC timing out. - if ( - health.timedOut && - timeoutMs !== undefined && - probeTimeoutMs !== undefined && - probeTimeoutMs <= REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS - ) { - throw handleRequestTimeout({ - info, - statePaths, - ...timeoutRequestContext(req, true, timeoutMs), - }); - } - const remainingMs = remainingRemoteRequestTimeoutMs(info, req, statePaths, timeoutMs, deadline); if (!health.reachable) { throw new AppError('COMMAND_FAILED', 'Remote daemon is unavailable', { daemonBaseUrl: info.baseUrl, @@ -320,20 +344,20 @@ async function retryAfterRemoteInstanceMismatch( } function remainingRemoteRequestTimeoutMs( - info: DaemonInfo, req: DaemonRequest, - statePaths: DaemonPaths, timeoutMs: number | undefined, deadline: number | undefined, + timedOutProbeMs?: number, ): number | undefined { if (deadline === undefined || timeoutMs === undefined) return undefined; const remainingMs = deadline - performance.now(); - if (remainingMs > 0) return remainingMs; - throw handleRequestTimeout({ - info, - statePaths, - ...timeoutRequestContext(req, true, timeoutMs), - }); + // A probe capped by the request can expire before the monotonic clock catches up to its timer. + if ( + remainingMs > 0 && + (timedOutProbeMs === undefined || timedOutProbeMs > REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS) + ) + return remainingMs; + throw requestTimeoutError(timeoutMs, req.meta?.requestId); } function isRemoteInstanceMismatch(error: unknown): boolean { @@ -359,7 +383,7 @@ async function sendRequestWithTransport( ): Promise { return transport === 'http' ? await sendHttpRequest(info, req, statePaths, timeoutMs, options) - : await sendSocketRequest(info, req, statePaths, timeoutMs, options); + : await sendSocketRequest(info, req, timeoutMs, options); } function chooseTransport( @@ -454,7 +478,6 @@ function handleTransportError( async function sendSocketRequest( info: DaemonInfo, req: DaemonRequest, - statePaths: DaemonPaths, timeoutMs: number | undefined, options: SendRequestOptions, ): Promise { @@ -470,15 +493,10 @@ async function sendSocketRequest( const timeoutHandle = typeof timeoutMs === 'number' ? setTimeout(() => { + if (settled) return; settled = true; + reject(requestTimeoutError(timeoutMs, req.meta?.requestId)); socket.destroy(); - reject( - handleRequestTimeout({ - info, - statePaths, - ...timeoutRequestContext(req, false, timeoutMs), - }), - ); }, timeoutMs) : undefined; @@ -512,6 +530,14 @@ async function sendSocketRequest( }); } +function requestTimeoutError(timeoutMs: number, requestId: string | undefined): AppError { + return new AppError('COMMAND_FAILED', 'Daemon request timed out', { + reason: 'daemon_transport_timeout', + timeoutMs, + requestId, + }); +} + // The fields a timed-out request is described by, read once so a socket and an HTTP timeout cannot // describe the same request differently. type TimeoutRequestFields = Omit[0], 'info' | 'statePaths'>; @@ -644,14 +670,8 @@ async function sendHttpRequest( const timeoutHandle = typeof timeoutMs === 'number' ? setTimeout(() => { + reject(requestTimeoutError(timeoutMs, req.meta?.requestId)); request.destroy(); - reject( - handleRequestTimeout({ - info, - statePaths, - ...timeoutRequestContext(req, remote, timeoutMs), - }), - ); }, timeoutMs) : undefined; diff --git a/test/wire-compat/ledger.json b/test/wire-compat/ledger.json index 80e13dc3e9..20fc434604 100644 --- a/test/wire-compat/ledger.json +++ b/test/wire-compat/ledger.json @@ -67,7 +67,7 @@ "src/daemon-client/daemon-client-rpc.ts#toDaemonHttpRpcError": "sha256:888246763c48670e7da893054d025744654f8715c3b4906312617a2b5028316b", "src/daemon-client/daemon-client-transport.ts#RemoteDaemonHealth": "sha256:3bac36fa97090b273afe1128103e1bf476d05441fd9de41317f1b78775b88304", "src/daemon-client/daemon-client-transport.ts#RemoteDaemonHealthLink": "sha256:7702598468b4c82b4ffa93064cd6bdfec938c1c527f7e6007c0584f3884a6eea", - "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth": "sha256:eef0e153eb6f0bc02175e464dd6d98559b500b65ea8c74ba2203cfef09ec09a0", + "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth": "sha256:f40cc2f5bfb8add78f99b11db7842bc244735786aa390d10a38ef025e16dc53a", "src/daemon-client/daemon-client-transport.ts#readHealthLink": "sha256:f56404b94d73da57de7248719c9136b760d2be8b4a53cde5315a2311b088425b", "src/daemon-client/daemon-client-transport.ts#readHealthPayload": "sha256:4e85ffc3e35e02379c393e9312344757e003cf1f0ad9eb8d1f77d90c81c861f1", "src/daemon-client/daemon-client-transport.ts#readRemoteDaemonHealth": "sha256:bcefa89fbb7fcbd6fee1b5ecb217955b9eee8d6fd2ad175fb653011edbd199d9", @@ -431,8 +431,8 @@ }, { "declaration": "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth", - "digest": "sha256:eef0e153eb6f0bc02175e464dd6d98559b500b65ea8c74ba2203cfef09ec09a0", - "rationale": "A restart retry must tell a probe that ran out of time from one that failed, so the client can report the RPC deadline instead of 'Remote daemon is unavailable'. `timedOut` is client-local: the prober sets it on its own timeout and abort paths and never reads it from a /health payload. The health request and the accepted payload fields are unchanged, so a released daemon or proxy is probed and parsed exactly as before." + "digest": "sha256:f40cc2f5bfb8add78f99b11db7842bc244735786aa390d10a38ef025e16dc53a", + "rationale": "#3116 includes requester loading and response completion in the existing health-probe budget. `timedOut` remains client-local; it is never read from a /health payload. Request routes, authentication, accepted fields and protocol-version admission are unchanged, so released daemons and proxies still send payloads this client parses correctly. This is an implementation-only timing change under ADR 0006." } ] }