diff --git a/packages/platform-apple/src/runner/__tests__/runner-session-close.test.ts b/packages/platform-apple/src/runner/__tests__/runner-session-close.test.ts index ab1a2292a4..67b43e3133 100644 --- a/packages/platform-apple/src/runner/__tests__/runner-session-close.test.ts +++ b/packages/platform-apple/src/runner/__tests__/runner-session-close.test.ts @@ -1,5 +1,7 @@ import assert from 'node:assert/strict'; +import { EventEmitter } from 'node:events'; import { beforeEach, test, vi } from 'vitest'; +import type { ExecBackgroundResult } from '@agent-device/host-kit/command'; import { IOS_SIMULATOR } from './device-fixtures.ts'; import { appleRunnerTestHost } from '../test-host.ts'; import { @@ -105,6 +107,8 @@ import { readRunnerSessionLiveness, releaseIosRunnerOnClose, } from '../runner-session.ts'; +import { registerRunnerPrepProcess } from '../runner-artifact.ts'; +import { runnerPrepProcessChildren } from '../runner-xctestrun.ts'; // Test-only stand-in for the daemon's runtime lease-owner-state-dir setter (root-only; the package // cannot import it). Backs the host.leaseOwnerStateDir() getter the package reads instead. @@ -317,3 +321,77 @@ test('releaseIosRunnerOnClose retains an idle runner, disposes a busy one, and t await releaseIosRunnerOnClose(device.id, { retain: false }); assert.equal(readRunnerSessionLiveness(device.id), null); }); + +/** + * A close that stops the device stops its in-flight runner build too (#3177). The build behind a + * cold start is the resource a session teardown promises to leave behind, not an orphan for the + * next `open` to race on the shared runner derived-data root. + */ +test('releaseIosRunnerOnClose stops the device build still in flight (#3177)', async () => { + const device = { ...IOS_SIMULATOR, id: 'runner-session-close-build-sim' }; + await ensureRunnerSession(device, {}); + const build = Object.assign(new EventEmitter(), { + pid: 4747, + exitCode: null, + }) as ExecBackgroundResult['child']; + registerRunnerPrepProcess(device.id, build); + + await releaseIosRunnerOnClose(device.id, { retain: false }); + + const signaledPids = mockSignalProcessGroupBestEffort.mock.calls.map(([pid]) => pid); + assert.ok( + signaledPids.includes(4747), + 'the non-retained close tree-killed the build still in flight', + ); + assert.equal(runnerPrepProcessChildren(device.id).length, 0, 'the build left the prep ledger'); +}); + +/** + * A close arriving during `build-for-testing` stops the build BEFORE waiting for the session + * lock (#3177 review). The start holds that lock for its whole cold build, so a close that stopped + * the session first would only reach the build after it finished — exactly the orphan close is + * supposed to prevent. Here the start is parked mid-build (the lock held), and the build must be + * signaled without the session stop completing first. + */ +test('releaseIosRunnerOnClose stops the build while the start still holds the session lock (#3177)', async () => { + const device = { ...IOS_SIMULATOR, id: 'runner-session-close-ordering-sim' }; + let releaseBuild: () => void = () => {}; + mockEnsureXctestrunArtifact.mockImplementationOnce( + () => + new Promise>>((resolve) => { + releaseBuild = () => + resolve({ + xctestrunPath: '/tmp/base-runner.xctestrun', + derived: '/tmp/derived', + cacheKey: RUNNER_CACHE_KEY_FIXTURE, + cache: 'miss', + artifact: 'rebuilt', + buildMs: 12, + xctestrunPathSource: 'build', + }); + }), + ); + const starting = ensureRunnerSession(device, {}).catch(() => undefined); + await vi.waitFor(() => { + assert.equal(mockEnsureXctestrunArtifact.mock.calls.length, 1); + }); + + const build = Object.assign(new EventEmitter(), { + pid: 4748, + exitCode: null, + }) as ExecBackgroundResult['child']; + registerRunnerPrepProcess(device.id, build); + const closing = releaseIosRunnerOnClose(device.id, { retain: false }); + + await vi.waitFor(() => { + const signaledPids = mockSignalProcessGroupBestEffort.mock.calls.map(([pid]) => pid); + assert.ok( + signaledPids.includes(4748), + 'the close killed the build while the start still held the session lock', + ); + }); + + releaseBuild(); + await closing; + await starting; +}); diff --git a/packages/platform-apple/src/runner/__tests__/runner-start-budget.test.ts b/packages/platform-apple/src/runner/__tests__/runner-start-budget.test.ts new file mode 100644 index 0000000000..dbe8e555f2 --- /dev/null +++ b/packages/platform-apple/src/runner/__tests__/runner-start-budget.test.ts @@ -0,0 +1,198 @@ +import assert from 'node:assert/strict'; +import { EventEmitter } from 'node:events'; +import { beforeEach, test, vi } from 'vitest'; +import type { ExecBackgroundResult } from '@agent-device/host-kit/command'; +import { + AppError, + createRequestCanceledError, + isRequestCanceledError, +} from '@agent-device/kernel/errors'; +import { IOS_SIMULATOR } from './device-fixtures.ts'; +import { appleRunnerTestHost } from '../test-host.ts'; +import { raceRunnerStartAgainstCaller } from '../runner-start-budget.ts'; +import { registerRunnerPrepProcess } from '../runner-artifact.ts'; +import { runnerPrepProcessChildren } from '../runner-xctestrun.ts'; + +const mockSignalPidsBestEffort = vi.fn(); +const mockSignalProcessGroupBestEffort = vi.fn(); +const mockRunAppleToolCommand = vi.fn(); +const mockGetRequestSignal = vi.fn(); + +beforeEach(() => { + vi.resetAllMocks(); + mockRunAppleToolCommand.mockResolvedValue({ exitCode: 0, stdout: '', stderr: '' }); + mockGetRequestSignal.mockReturnValue(undefined); + appleRunnerTestHost.update({ + signalPidsBestEffort: mockSignalPidsBestEffort, + signalProcessGroupBestEffort: mockSignalProcessGroupBestEffort, + runAppleToolCommand: mockRunAppleToolCommand, + getRequestSignal: mockGetRequestSignal, + }); +}); + +function callerDeadline(): DOMException { + return new DOMException('Wait deadline exceeded', 'TimeoutError'); +} + +/** A detached start nobody has finished — the shape of a cold build still under the lock. */ +function hangingStart(): Promise { + return new Promise(() => {}); +} + +function makePrepChild(pid: number): ExecBackgroundResult['child'] { + return Object.assign(new EventEmitter(), { + pid, + exitCode: null, + }) as ExecBackgroundResult['child']; +} + +function signaledPids(): number[] { + return mockSignalProcessGroupBestEffort.mock.calls.map(([pid]) => pid as number); +} + +function canceled(error: unknown): boolean { + return isRequestCanceledError(error) && error instanceof AppError; +} + +async function settleAsyncWork(): Promise { + for (let i = 0; i < 4; i += 1) { + await new Promise((resolve) => setTimeout(resolve, 5)); + } +} + +/** + * The waiter's cancel owns the build it waited on (#3177). A request canceled while queued behind + * a detached cold build must stop that build: its spawn carried only the STARTING request's + * cancellation signal, so without the waiter's device-scoped prep kill the build would keep + * compiling under the daemon on the shared runner derived-data root, and a retried `open` would + * race it. The kill is the same tree-kill escalation a session stop uses. Here the build's owner + * has no live request signal (its start is detached), so the waiter is allowed to stop it. + */ +test('a request canceled while waiting on the detached start stops that device build', async () => { + const device = { ...IOS_SIMULATOR, id: 'runner-waiter-cancel-sim' }; + const build = makePrepChild(4848); + registerRunnerPrepProcess(device.id, build, 'owner-request-gone'); + mockGetRequestSignal.mockImplementation((requestId?: string) => + requestId === 'owner-request-gone' ? AbortSignal.abort() : undefined, + ); + + const controller = new AbortController(); + controller.abort(createRequestCanceledError()); + await assert.rejects( + raceRunnerStartAgainstCaller(hangingStart(), controller.signal, device.id), + canceled, + ); + await settleAsyncWork(); + + assert.ok( + signaledPids().includes(4848), + 'the waiter cancel reached the build through the prep tree-kill path', + ); + assert.equal( + runnerPrepProcessChildren(device.id).length, + 0, + 'the killed build left the prep ledger', + ); +}); + +/** + * A canceled waiter must not reach a build still owned by an in-flight request (#3177 review). + * The owner cancels its own build through its live request signal at the exec layer; a waiter + * tearing it down would SIGTERM another active request's work out from under it. Here the build's + * owner still has a registered, un-aborted signal, so the canceled waiter leaves it running. + */ +test('a canceled waiter leaves a build owned by a still-active request', async () => { + const device = { ...IOS_SIMULATOR, id: 'runner-waiter-owner-active-sim' }; + const build = makePrepChild(4850); + registerRunnerPrepProcess(device.id, build, 'owner-request-live'); + mockGetRequestSignal.mockImplementation((requestId?: string) => + requestId === 'owner-request-live' ? new AbortController().signal : undefined, + ); + + const controller = new AbortController(); + controller.abort(createRequestCanceledError()); + await assert.rejects( + raceRunnerStartAgainstCaller(hangingStart(), controller.signal, device.id), + canceled, + ); + await settleAsyncWork(); + + assert.equal( + mockSignalProcessGroupBestEffort.mock.calls.length, + 0, + 'a canceled waiter never signals a build another active request owns', + ); + assert.deepEqual( + runnerPrepProcessChildren(device.id).map((child) => child.pid), + [4850], + 'the owned build stayed registered', + ); +}); + +/** + * #2894 still holds on this seam: the caller's own deadline (a bounded poll) is NOT a + * cancellation. The start it interrupts is the one the retry joins, so a deadline must leave the + * device's build running and un-signaled. + */ +test('a caller deadline on the same waiter leaves the build running', async () => { + const device = { ...IOS_SIMULATOR, id: 'runner-waiter-deadline-sim' }; + const build = makePrepChild(4949); + registerRunnerPrepProcess(device.id, build); + + const controller = new AbortController(); + const waiting = raceRunnerStartAgainstCaller(hangingStart(), controller.signal, device.id); + controller.abort(callerDeadline()); + await assert.rejects(waiting, canceled); + await settleAsyncWork(); + + assert.equal( + mockSignalProcessGroupBestEffort.mock.calls.length, + 0, + 'a deadline never signals the build a retry needs (#2894)', + ); + assert.deepEqual( + runnerPrepProcessChildren(device.id).map((child) => child.pid), + [4949], + 'the deadline left the build registered', + ); +}); + +/** + * The kill is scoped to the waiting device. A canceled waiter for device A must not signal the + * build device B is still paying for — the same request-scoping rule this PR applies to the + * daemon-side timeout recovery. + */ +test('a canceled waiter signals only its own device build', async () => { + const device = { ...IOS_SIMULATOR, id: 'runner-waiter-scope-sim' }; + const ownBuild = makePrepChild(5050); + const siblingBuild = makePrepChild(5151); + registerRunnerPrepProcess(device.id, ownBuild); + registerRunnerPrepProcess('other-device', siblingBuild); + + const controller = new AbortController(); + controller.abort(createRequestCanceledError()); + await assert.rejects( + raceRunnerStartAgainstCaller(hangingStart(), controller.signal, device.id), + canceled, + ); + await settleAsyncWork(); + + assert.ok(signaledPids().includes(5050), 'the waiting device build was signaled'); + assert.ok( + !signaledPids().includes(5151), + 'a build for another device was left running (#3177 sibling protection)', + ); +}); + +/** A prep child that exits on its own leaves the ledger through its close event. */ +test('a closed build leaves the prep ledger on its own', () => { + const device = { ...IOS_SIMULATOR, id: 'runner-prep-close-sim' }; + const build = new EventEmitter(); + registerRunnerPrepProcess( + device.id, + Object.assign(build, { pid: 5252 }) as ExecBackgroundResult['child'], + ); + assert.equal(runnerPrepProcessChildren(device.id).length, 1); + build.emit('close'); + assert.equal(runnerPrepProcessChildren(device.id).length, 0); +}); diff --git a/packages/platform-apple/src/runner/runner-artifact.ts b/packages/platform-apple/src/runner/runner-artifact.ts index bb5ecaa504..defe03ce6c 100644 --- a/packages/platform-apple/src/runner/runner-artifact.ts +++ b/packages/platform-apple/src/runner/runner-artifact.ts @@ -9,6 +9,7 @@ import { withProcessLock, emitRequestProgress, findProjectRoot, + getRequestSignal, isCommandTimeoutError, } from './host.ts'; import type { ExecBackgroundResult } from '@agent-device/host-kit/command'; @@ -60,7 +61,78 @@ import { resolveAppleRunnerProjectPath } from './runner-source.ts'; export { prepareXctestrunWithEnv } from './runner-artifact-env.ts'; const runnerXctestrunBuildLocks = new Map>(); -export const runnerPrepProcesses = new Set(); + +type RunnerPrepProcess = Readonly<{ + deviceId: string; + requestId: string | undefined; + child: ExecBackgroundResult['child']; +}>; + +const runnerPrepProcessLedger = new Set(); + +/** + * Records a prep subprocess (`xcodebuild build-for-testing`) against the device it builds for and + * the request whose start spawned it, so a request canceled while waiting on that build can stop + * it (#3177). The build child keeps its owning start's signal as its first cancel path; this + * ledger is the device-scoped second one, for the waiters whose cancellation the spawn never saw. + * The owner is what keeps the second path from reaching a build another, still-active request + * launched: a canceled waiter stops only builds whose owner can no longer cancel them. + */ +export function registerRunnerPrepProcess( + deviceId: string, + child: ExecBackgroundResult['child'], + requestId?: string, +): void { + const entry: RunnerPrepProcess = { deviceId, requestId, child }; + runnerPrepProcessLedger.add(entry); + child.on('close', () => { + runnerPrepProcessLedger.delete(entry); + }); +} + +/** The prep subprocesses still running, for one device or for every device when none is named. */ +export function runnerPrepProcessChildren( + deviceId?: string, +): readonly ExecBackgroundResult['child'][] { + return prepProcessEntries(deviceId).map((entry) => entry.child); +} + +/** + * The prep subprocesses whose owning start is detached: the request that spawned it has finished, + * was canceled, or never existed (an in-process build with no request), so no live cancellation + * can reach the build anymore. A canceled waiter stops exactly these, never a build still owned by + * an in-flight request — that one belongs to its owner and dies through the owner's own signal. + */ +export function runnerPrepProcessChildrenWithoutActiveOwner( + deviceId?: string, +): readonly ExecBackgroundResult['child'][] { + return prepProcessEntries(deviceId) + .filter((entry) => !isRunnerPrepOwnerActive(entry.requestId)) + .map((entry) => entry.child); +} + +function prepProcessEntries(deviceId?: string): readonly RunnerPrepProcess[] { + return [...runnerPrepProcessLedger].filter( + (entry) => deviceId === undefined || entry.deviceId === deviceId, + ); +} + +/** + * Whether the request owning a prep build can still cancel it: a registered, un-aborted request + * signal means the owner is in flight and its cancellation reaches the build through the exec + * layer directly. + */ +function isRunnerPrepOwnerActive(requestId: string | undefined): boolean { + if (!requestId) return false; + const ownerSignal = getRequestSignal(requestId); + return ownerSignal !== undefined && !ownerSignal.aborted; +} + +export function forgetRunnerPrepProcess(child: ExecBackgroundResult['child']): void { + for (const entry of runnerPrepProcessLedger) { + if (entry.child === child) runnerPrepProcessLedger.delete(entry); + } +} export type RunnerXctestrunArtifactState = 'valid' | 'rebuilt'; @@ -87,6 +159,8 @@ type RunnerXctestrunBuildOptions = { verbose?: boolean; logPath?: string; traceLogPath?: string; + /** The request whose start owns this build; recorded so cancel rules respect the owner. */ + requestId?: string; /** * The build phase's one budget, opened by whoever owns the build: the cache decision's * blocking toolchain probes and `xcodebuild` spend the same clock, and the owning @@ -484,10 +558,7 @@ async function buildRunnerXctestrun( timeoutMs: buildTimeoutMs, signal: options.budget?.signal, onSpawn: (child) => { - runnerPrepProcesses.add(child); - child.on('close', () => { - runnerPrepProcesses.delete(child); - }); + registerRunnerPrepProcess(device.id, child, options.requestId); }, onStdoutChunk: (chunk) => { logChunk(chunk, options.logPath, options.traceLogPath, options.verbose); diff --git a/packages/platform-apple/src/runner/runner-disposal.ts b/packages/platform-apple/src/runner/runner-disposal.ts index 1dc18c2d5e..0aa61ec94d 100644 --- a/packages/platform-apple/src/runner/runner-disposal.ts +++ b/packages/platform-apple/src/runner/runner-disposal.ts @@ -25,7 +25,12 @@ import { type RunnerLeaseCleanupAdapter, type RunnerXcodebuildCleanupTarget, } from './runner-lease.ts'; -import { IOS_RUNNER_CONTAINER_BUNDLE_IDS, runnerPrepProcesses } from './runner-xctestrun.ts'; +import { + forgetRunnerPrepProcess, + IOS_RUNNER_CONTAINER_BUNDLE_IDS, + runnerPrepProcessChildren, + runnerPrepProcessChildrenWithoutActiveOwner, +} from './runner-xctestrun.ts'; import { advanceRunnerSessionState, type RunnerSession } from './runner-session-types.ts'; export const RUNNER_INVALIDATE_WAIT_TIMEOUT_MS = 1_000; @@ -88,7 +93,7 @@ export async function cleanupOwnedIosRunnerLease(deviceId: string): Promise { - const prepProcesses = Array.from(runnerPrepProcesses); + const prepProcesses = runnerPrepProcessChildren(); const macOsSessions = activeSessions.filter((session) => isMacOs(session.device)); const otherSessions = activeSessions.filter((session) => !isMacOs(session.device)); for (const session of activeSessions) { @@ -108,15 +113,35 @@ export async function abortRunnerSessionsAndPrepProcesses( ); } -export async function stopRunnerPrepProcesses(): Promise { - const prepProcesses = Array.from(runnerPrepProcesses); +/** + * Stops the prep subprocesses (the `xcodebuild build-for-testing` behind a cold runner start) + * with the tree-kill escalation the sessions get. A device stops only its own builds: the caller + * that stops device A's session must not sweep device B's in-flight build (#3177). + */ +export async function stopRunnerPrepProcesses(deviceId?: string): Promise { + await stopPrepProcessList(runnerPrepProcessChildren(deviceId)); +} + +/** + * Stops the device builds no active request owns anymore (#3177). A canceled waiter may stop the + * build it waited on only once that build's owning start is detached; a build still owned by an + * in-flight request belongs to its owner, which cancels it through its own signal, and a waiter + * must not SIGTERM another request's work out from under it. + */ +export async function stopRunnerPrepProcessesWithoutActiveOwner(deviceId?: string): Promise { + await stopPrepProcessList(runnerPrepProcessChildrenWithoutActiveOwner(deviceId)); +} + +async function stopPrepProcessList( + prepProcesses: readonly ExecBackgroundResult['child'][], +): Promise { await Promise.allSettled( prepProcesses.map(async (child) => { try { await killRunnerProcessTree(child.pid, 'SIGTERM'); await killRunnerProcessTree(child.pid, 'SIGKILL'); } finally { - runnerPrepProcesses.delete(child); + forgetRunnerPrepProcess(child); } }), ); @@ -311,7 +336,7 @@ async function signalRunnerPrepProcesses( prepProcesses.map(async (child) => { await killRunnerProcessTree(child.pid, signal); if (signal === 'SIGKILL') { - runnerPrepProcesses.delete(child); + forgetRunnerPrepProcess(child); } }), ); diff --git a/packages/platform-apple/src/runner/runner-session.ts b/packages/platform-apple/src/runner/runner-session.ts index bd922c0631..ce4fac520e 100644 --- a/packages/platform-apple/src/runner/runner-session.ts +++ b/packages/platform-apple/src/runner/runner-session.ts @@ -120,7 +120,7 @@ export async function ensureRunnerSession( } }); const { raceRunnerStartAgainstCaller } = await import('./runner-start-budget.ts'); - return await raceRunnerStartAgainstCaller(start, options.signal); + return await raceRunnerStartAgainstCaller(start, options.signal, device.id); } /** How long the device-readiness probe may take, bounded by the startup budget it runs inside. */ @@ -675,6 +675,8 @@ export async function stopIosRunnerSession(deviceId: string): Promise { * or wedges, so pooling it back hands the same stalled process to the next `open` (#2552). An idle * retained runner keeps warm reuse via the idle-stop timer. The decision is owned here because the * occupancy fact lives on the session, and awaited so `close` returns only once the lease is gone. + * Non-retained close stops the device's current prep processes before taking the session lock, + * which an in-flight cold start holds through its build. Later prep spawns are not fenced here. */ export async function releaseIosRunnerOnClose( deviceId: string, @@ -692,6 +694,7 @@ export async function releaseIosRunnerOnClose( data: { deviceId }, }); } + await stopRunnerPrepProcesses(deviceId); await stopIosRunnerSession(deviceId); } diff --git a/packages/platform-apple/src/runner/runner-start-budget.ts b/packages/platform-apple/src/runner/runner-start-budget.ts index 63288ec44a..50d087d917 100644 --- a/packages/platform-apple/src/runner/runner-start-budget.ts +++ b/packages/platform-apple/src/runner/runner-start-budget.ts @@ -4,7 +4,8 @@ import { isRequestCanceledError, } from '@agent-device/kernel/errors'; import { emitDiagnostic } from './host.ts'; -import { resolveRunnerStartupSignal } from './runner-contract.ts'; +import { isCallerDeadlineAbortReason, resolveRunnerStartupSignal } from './runner-contract.ts'; +import { stopRunnerPrepProcessesWithoutActiveOwner } from './runner-disposal.ts'; import { createRunnerPhaseBudget, type RunnerPhaseBudget } from './runner-xctestrun.ts'; import { normalizeRunnerStartupTimeoutMs, type RunnerSession } from './runner-session-types.ts'; import type { AppleRunnerLifecycleOptions } from './runner-provider.ts'; @@ -85,18 +86,34 @@ function runnerStartBudgetExhaustedError(timeoutMs: number, explicit: boolean): * signal allows. A caller whose deadline lands during a cold xctestrun build leaves on time, the * build keeps going under the lock, and the next request for the device queues behind it and joins * the session it registers (#2894). Whatever the abort reason, the caller sees the same cancelled - * request it would have seen from any later step; a cancelled request's abort also reaches the - * start through its own startup signal, so nothing here decides whether the start survives. A - * start that fails after its caller left has nobody to report to, so its failure is logged here. + * request it would have seen from any later step. A start that fails after its caller left has + * nobody to report to, so its failure is logged here. + * + * The two abort reasons get opposite treatment of the detached start, and the difference is the + * whole point (#2894 vs #3177). A caller's own deadline (a bounded poll) must leave the start + * running: it is the start the retry joins. A cancelled request means the client is gone, and the + * start it was waiting for belongs to nobody — its spawn carried the *waiting request's* + * cancellation signal only when that request opened the start, so a request that merely joined a + * start another request spawned has no path to it otherwise. On cancel the waiter stops the + * device's prep subprocesses through the same tree-kill path a session stop uses, so a timed-out + * `open` cannot orphan a `build-for-testing` on the shared runner derived-data root where a + * retried `open` would race it. Only builds whose owning start is detached are stopped: a build + * still owned by an in-flight request belongs to its owner and dies through the owner's own + * signal, never under a canceled waiter. The start itself keeps running under the lock (bounded + * by its own budget); its build is left running only while its owner is still there to cancel it. */ export async function raceRunnerStartAgainstCaller( start: Promise, signal: AbortSignal | undefined, + deviceId: string, ): Promise { if (!signal) return await start; return await new Promise((resolve, reject) => { const abort = () => { reject(createRequestCanceledError(undefined, signal.reason)); + if (!isCallerDeadlineAbortReason(signal.reason)) { + void stopRunnerPrepProcessesWithoutActiveOwner(deviceId); + } start.catch(emitDetachedRunnerStartFailed); }; if (signal.aborted) { diff --git a/packages/platform-apple/src/runner/runner-xctestrun.ts b/packages/platform-apple/src/runner/runner-xctestrun.ts index 944fb4652d..2be0a78a3a 100644 --- a/packages/platform-apple/src/runner/runner-xctestrun.ts +++ b/packages/platform-apple/src/runner/runner-xctestrun.ts @@ -1,9 +1,10 @@ export { ensureXctestrunArtifact, + forgetRunnerPrepProcess, hasCachedAppleRunnerArtifact, prepareXctestrunWithEnv, - runnerPrepProcesses, - type ExternalXctestRunnerOptions, + runnerPrepProcessChildren, + runnerPrepProcessChildrenWithoutActiveOwner, type RunnerXctestrunArtifact, type RunnerXctestrunArtifactState, } from './runner-artifact.ts'; diff --git a/scripts/layering/daemon-client-entry.ts b/scripts/layering/daemon-client-entry.ts index 1c7ced3664..16b687df4f 100644 --- a/scripts/layering/daemon-client-entry.ts +++ b/scripts/layering/daemon-client-entry.ts @@ -48,6 +48,11 @@ export const DAEMON_CLIENT_ENTRY_EDGES: readonly DaemonClientEntryEdge[] = [ target: 'src/daemon/daemon-request.ts', rationale: REQUEST_RATIONALE, }, + { + file: 'src/daemon-client/daemon-client-liveness-probe.ts', + target: 'src/daemon/daemon-request.ts', + rationale: REQUEST_RATIONALE, + }, { file: 'src/daemon-client/daemon-client-progress.ts', target: 'src/daemon/daemon-request.ts', diff --git a/src/__tests__/test-utils/loopback.ts b/src/__tests__/test-utils/loopback.ts index 0f8d1cf5aa..e119d9805a 100644 --- a/src/__tests__/test-utils/loopback.ts +++ b/src/__tests__/test-utils/loopback.ts @@ -44,7 +44,18 @@ export async function skipWhenLoopbackUnavailable( return true; } +// Connections the test server accepted after it started listening. A `net.Server` refuses to close +// while a connection is live, and a hung-request stand-in keeps one open past the client's RST, so +// teardown destroys them — the job `DaemonServer.destroyConnections` does for the real daemon. +const acceptedConnections = new WeakMap>(); + export async function listenOnLoopback(server: LoopbackServer): Promise { + const sockets = new Set(); + acceptedConnections.set(server, sockets); + server.on('connection', (socket) => { + sockets.add(socket); + socket.on('close', () => sockets.delete(socket)); + }); await new Promise((resolve, reject) => { server.once('error', reject); server.listen(0, '127.0.0.1', () => { @@ -60,6 +71,10 @@ export async function listenOnLoopback(server: LoopbackServer): Promise } export async function closeLoopbackServer(server: LoopbackServer): Promise { + for (const socket of acceptedConnections.get(server) ?? []) { + socket.destroy(); + } + acceptedConnections.delete(server); if (!server.listening) return; closeHttpConnections(server); await new Promise((resolve, reject) => { diff --git a/src/daemon-client/__tests__/boundary-fault-transport.test.ts b/src/daemon-client/__tests__/boundary-fault-transport.test.ts index 83242e6363..cff373a906 100644 --- a/src/daemon-client/__tests__/boundary-fault-transport.test.ts +++ b/src/daemon-client/__tests__/boundary-fault-transport.test.ts @@ -70,7 +70,10 @@ for (const [commandClass, row] of Object.entries(COMMAND_ROWS)) { return true; }, ); - assert.equal(mockRunCmdSync.mock.calls.filter(([command]) => command === 'pkill').length, 3); + // #3177: a timed-out request never spawns the host-wide runner `pkill` sweep again. The + // timeout must stay bounded not only in wall-clock but in blast radius: the daemon cancels + // this request when the connection dies, and processes belong to the daemon's leases. + assert.equal(mockRunCmdSync.mock.calls.filter(([command]) => command === 'pkill').length, 0); }); for (const responseShape of ['partial', 'malformed'] as const) { diff --git a/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts b/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts index 294eff629b..9c017ef986 100644 --- a/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts +++ b/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts @@ -34,7 +34,7 @@ const { handleRequestTimeoutCalls } = vi.hoisted(() => ({ })); vi.mock('../daemon-client-timeout.ts', () => ({ - handleRequestTimeout: (params: unknown) => { + handleRequestTimeout: async (params: unknown) => { handleRequestTimeoutCalls.push(params); return new AppError('COMMAND_FAILED', 'Daemon request timed out', { reason: 'daemon_transport_timeout', diff --git a/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts b/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts index 39e80a6123..0d16a79520 100644 --- a/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts +++ b/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts @@ -8,7 +8,6 @@ import { afterEach, test, vi } from 'vitest'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; import { spawnRegisteredDaemonFixture, - waitForRegisteredDaemonFixture, finishRegisteredDaemonFixture, finishRegisteredDaemonFixtures, } from '../../__tests__/test-utils/registered-daemon-fixture.ts'; @@ -25,9 +24,7 @@ vi.mock('@agent-device/host-kit/retry', async (importOriginal) => ({ })); import { resolveDaemonPaths, type DaemonPaths } from '../../daemon-resolution.ts'; -import { sendToDaemon, type DaemonRequest } from '../daemon-client.ts'; -import { sendRequest } from '../daemon-client-transport.ts'; -import type { DaemonRetirementResult } from '../../daemon-registration-owner.ts'; +import { sendToDaemon } from '../daemon-client.ts'; import { closeLoopbackServer, listenOnLoopback, @@ -95,18 +92,6 @@ function writeDaemonInfo(paths: DaemonPaths, info: DaemonInfoFixture): void { ); } -function writeDaemonLock( - paths: DaemonPaths, - lock: { pid: number; processStartTime?: string; startedAt?: number }, -): void { - fs.mkdirSync(paths.baseDir, { recursive: true }); - fs.writeFileSync( - paths.lockPath, - `${JSON.stringify({ startedAt: Date.now(), ...lock })}\n`, - 'utf8', - ); -} - /** Like `startHttpDaemonFixture`, but every RPC call returns `errorResult` as an `{ok:false}` result. */ async function startHttpDaemonErrorFixture( errorResult: Record, @@ -177,37 +162,6 @@ function installSpawnedHttpDaemonAtOwnedStateDir( }); } -async function startHangingHttpDaemonFixture(): Promise { - const seenPaths: string[] = []; - const rpcRequests: Record[] = []; - const server = http.createServer((req, res) => { - const url = new URL(req.url || '/', 'http://127.0.0.1'); - seenPaths.push(`${req.method ?? 'GET'} ${url.pathname}`); - - if (req.method === 'GET' && url.pathname === '/health') { - res.writeHead(200); - res.end('ok'); - return; - } - - if (req.method === 'POST' && url.pathname === '/rpc') { - const chunks: Buffer[] = []; - req.on('data', (chunk) => { - chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); - }); - req.on('end', () => { - rpcRequests.push(JSON.parse(Buffer.concat(chunks).toString('utf8')) as Record); - }); - return; - } - - res.writeHead(404); - res.end('not found'); - }); - const port = await listenOnLoopback(server); - return { server, port, seenPaths, rpcRequests }; -} - function installSpawnedHttpDaemon(paths: DaemonPaths, httpPort: number): void { mockSleep.mockImplementation(actualRetry.sleep); mockRunCmdDetached.mockImplementation((_command, _args, options) => { @@ -604,73 +558,6 @@ test('sendToDaemon replaces socket-only daemon metadata when HTTP transport is r } }); -test('sendRequest timeout cleanup uses resolved daemon paths instead of request flags', async (t) => { - if (!(await supportsLoopbackBind())) { - t.skip('loopback listeners are not permitted in this environment'); - return; - } - - const daemonStateDir = makeTempStateDir('agent-device-daemon-timeout-active-'); - const requestFlagStateDir = makeTempStateDir('agent-device-daemon-timeout-request-'); - const daemonPaths = resolveDaemonPaths(daemonStateDir); - const requestFlagPaths = resolveDaemonPaths(requestFlagStateDir); - const daemon = await startHangingHttpDaemonFixture(); - 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', - pid: 999_998, - }); - writeDaemonLock(requestFlagPaths, { pid: 999_998 }); - - const request: DaemonRequest = { - session: 'default', - command: 'replay', - positionals: [], - flags: { stateDir: requestFlagStateDir, daemonTransport: 'http' }, - token: 'local-secret', - meta: { requestId: 'req-timeout-paths' }, - }; - - try { - const info = await waitForRegisteredDaemonFixture(daemonPaths, child); - let thrown: unknown; - try { - 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); - assert.equal(fs.existsSync(requestFlagPaths.infoPath), true); - assert.equal(fs.existsSync(requestFlagPaths.lockPath), true); - } finally { - await closeLoopbackServer(daemon.server); - await finishRegisteredDaemonFixture(daemonStateDir); - fs.rmSync(requestFlagStateDir, { recursive: true, force: true }); - } -}); - test('sendToDaemon falls back from failed socket transport to HTTP using daemon metadata ports', async (t) => { if (!(await supportsLoopbackBind())) { t.skip('loopback listeners are not permitted in this environment'); diff --git a/src/daemon-client/__tests__/daemon-client-liveness-probe.test.ts b/src/daemon-client/__tests__/daemon-client-liveness-probe.test.ts new file mode 100644 index 0000000000..161af40649 --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-liveness-probe.test.ts @@ -0,0 +1,347 @@ +// Source-mirroring coverage for the post-timeout liveness probe (#3177) — the question that +// authorizes (or refuses) a daemon reset. `daemon-client-timeout-route.test.ts` drives the route +// through `sendRequest`; this file pins the probe's own verdicts, budget, and wire request, which +// the route tests exercise only one endpoint at a time. Stand-ins answer or refuse immediately; +// only the tests that need a SILENT endpoint (the regressions the concurrent-legs design and the +// absolute deadline exist for) spend wall-clock, and each stays inside the unit budget. + +import net from 'node:net'; +import http from 'node:http'; +import assert from 'node:assert/strict'; +import { test } from 'vitest'; + +import { PUBLIC_COMMANDS } from '@agent-device/command-registry/catalog'; +import { resolveCommandTimeoutPolicy } from '@agent-device/command-registry/registry'; +import { loadNodeHttpRequester } from '@agent-device/host-kit/transport'; +import { + LIVENESS_PROBE_BUDGET_MS, + probeDaemonResponsive, +} from '../daemon-client-liveness-probe.ts'; +import { + closeLoopbackServer, + listenOnLoopback, + skipWhenLoopbackUnavailable, + type LoopbackServer, +} from '../../__tests__/test-utils/loopback.ts'; + +function daemonInfo(over: { port?: number; httpPort?: number }): { + port?: number; + httpPort?: number; + token: string; + pid: number; +} { + return { ...over, token: 'test-token', pid: process.pid }; +} + +async function withLoopback( + server: LoopbackServer, + run: (port: number) => Promise, +): Promise { + const port = await listenOnLoopback(server); + try { + return await run(port); + } finally { + await closeLoopbackServer(server); + } +} + +test('the probe budget stays an order of magnitude under the narrowest reset-eligible envelope', () => { + // A timed-out reset-eligible request may pay one extra probe window before its error lands. + // The narrowest envelope that class carries is the registry default (90s): keep the probe a + // rounding error against it, so the probe never dominates the operation it follows (AGENTS.md + // probe rule). If someone widens the probe toward the envelope, this fails at the owning number. + const narrowestResetEligibleEnvelopeMs = resolveCommandTimeoutPolicy( + PUBLIC_COMMANDS.open, + ).envelopeMs; + assert.equal(typeof narrowestResetEligibleEnvelopeMs, 'number'); + assert.ok( + LIVENESS_PROBE_BUDGET_MS * 10 <= (narrowestResetEligibleEnvelopeMs as number), + `probe budget ${LIVENESS_PROBE_BUDGET_MS}ms must stay 10x under the ${narrowestResetEligibleEnvelopeMs}ms envelope`, + ); +}); + +test('a daemon answering /health is responsive', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // The url is asserted, not just "an answer": a probe that drifted to any other route would + // still be answered by a server that answers everything, and the finding would then describe + // some other endpoint's liveness rather than the health route the transport also asks. + const requestedUrls: string[] = []; + const server = http.createServer((req, res) => { + requestedUrls.push(String(req.url)); + res.statusCode = 200; + res.end('{}'); + }); + await withLoopback(server, async (port) => { + assert.equal(await probeDaemonResponsive(daemonInfo({ httpPort: port })), true); + assert.deepEqual(requestedUrls, ['/health']); + }); +}); + +test('a 5xx answer is still an answer: the finding is liveness, not reachability policy', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // `readDaemonHttpHealth` (the reachability reader the transport uses before each command) reports + // a 5xx as UNREACHABLE — correct for its own policy, wrong for this one. The probe asks only + // whether the endpoint answered, so a daemon serving a 500 is alive and must not be killed: + // reading a reachability flag here would reproduce the bug #3177 is about. + const server = http.createServer((_req, res) => { + res.statusCode = 503; + res.end('overloaded'); + }); + await withLoopback(server, async (port) => { + assert.equal(await probeDaemonResponsive(daemonInfo({ httpPort: port })), true); + }); +}); + +test('headers alone answer: a daemon that stalls mid-body is still served', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // The finding is "the event loop served a fresh request", and that is already proven by a status + // line. A reader that waits for the whole BODY would call a daemon which answers and then stalls + // mid-response silent — the negative verdict would reset a live shared daemon, which is the bug + // #3177 is about. (A body-reading health helper does exactly that for its own purposes; the probe + // must not borrow it.) So the affirmative must arrive by the deadline's FIRST fraction, not at it. + const stallingBodyServer = http.createServer((_req, res) => { + res.writeHead(200, { 'content-length': '1000' }); + // Node buffers headers with the first chunk, so this is what makes the status line actually + // reach the client. The declared body is never sent: headers delivered, body pending forever. + res.flushHeaders(); + res.on('error', () => {}); + }); + await withLoopback(stallingBodyServer, async (port) => { + const startedAt = Date.now(); + assert.equal( + await probeDaemonResponsive(daemonInfo({ httpPort: port })), + true, + 'a stalled body must not be read as an unresponsive daemon', + ); + assert.ok( + Date.now() - startedAt < LIVENESS_PROBE_BUDGET_MS / 2, + 'the answer is headers-arrival, not the deadline running out', + ); + }); +}); + +test('the health leg rides a fresh connection, never the keep-alive pool', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // Since Node 19 `http.globalAgent` has `keepAlive: true`, so a default request can REUSE an idle + // socket from an earlier command's health read. Reuse would make "did this probe connect?" a + // question about a socket this probe never opened — and a pooled socket the daemon had half torn + // down would answer silent, resetting a live daemon. Discriminator: on a reused connection the + // server sees the SAME socket object for both requests, so `clientSockets.size` is 1 for reuse + // and 2 for a fresh connection. The prime must use the SAME module object the probe loads + // (`loadNodeHttpRequester('http:')` resolves the real `node:http`) or the pool holds no candidate + // and the test proves nothing. + // `keepAlive` is set at runtime (Node >=19) but not on the base `Agent` type. + if (!(http.globalAgent as { keepAlive?: boolean }).keepAlive) return; // pool cannot prime without keep-alive + const clientSockets = new Set(); + let probeRequests = 0; + let primed = false; + const server = http.createServer((req, res) => { + clientSockets.add(req.socket); + req.socket.on('error', () => {}); + if (primed) probeRequests += 1; + else primed = true; + res.end('{}'); + }); + const httpRequester = await loadNodeHttpRequester('http:'); + await withLoopback(server, async (port) => { + // Prime: a full keep-alive request/response whose socket returns to the global pool. + await new Promise((resolve, reject) => { + const request = httpRequester.request( + { host: '127.0.0.1', port, path: '/health', method: 'GET' }, + (res) => { + res.resume(); + res.on('end', resolve); + }, + ); + request.on('error', reject); + request.end(); + }); + assert.equal(clientSockets.size, 1, 'the prime must have connected'); + assert.equal(await probeDaemonResponsive(daemonInfo({ httpPort: port })), true); + assert.equal(probeRequests, 1, 'the health leg must actually have been asked'); + assert.equal( + clientSockets.size, + 2, + 'the probe rode the primed keep-alive socket instead of opening a fresh connection', + ); + }); +}); + +test('a refused endpoint is negative, not a hang: the verdict arrives well inside the budget', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // The wedged-daemon shape: connections are accepted and destroyed without an answer. The + // negative verdict must come from the refusal itself (fast), not from waiting out the window — + // a timeout error's latency must not depend on the probe budget when the host already answered. + const server = net.createServer((socket) => { + socket.on('error', () => {}); + socket.destroy(); + }); + await withLoopback(server, async (port) => { + const startedAt = Date.now(); + assert.equal(await probeDaemonResponsive(daemonInfo({ port })), false); + assert.ok( + Date.now() - startedAt < LIVENESS_PROBE_BUDGET_MS / 2, + 'refusal must settle the probe immediately, not by spending the window', + ); + }); +}); + +test('a silent HTTP leg cannot veto a live socket leg', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // The regression the concurrent-legs design exists for: sequential probing lets the first + // endpoint spend the whole budget (a transport that accepts but never answers), leaving the + // second endpoint no window — and the all-negative verdict SIGKILLs a daemon the other + // transport would have proven alive. The wrong verdict here kills every session on the host, + // which is the bug #3177 is about. + // The HTTP leg must be OBSERVED, not assumed: a verdict reached with no /health request ever + // sent would also pass if the probe simply skipped the leg. So the socket waits for the receipt + // before it answers — if the HTTP leg never asks, the socket never answers, the probe spends + // its window, and the `true` assertion below fails for the right reason. + const healthRequested: { value: boolean } = { value: false }; + const rpcHangingServer = http.createServer((req, res) => { + if (req.url === '/health') healthRequested.value = true; + // Accepts and never answers: this leg will spend the full budget. + res.on('error', () => {}); + }); + const liveSocketServer = net.createServer((socket) => { + socket.on('error', () => {}); + socket.on('data', () => { + const answer = () => { + socket.write(`${JSON.stringify({ jsonrpc: '2.0', id: 'probe', result: { ok: true } })}\n`); + }; + const waitForHttpLeg = (): void => { + if (healthRequested.value) answer(); + else setTimeout(waitForHttpLeg, 5); + }; + waitForHttpLeg(); + }); + }); + const httpPort = await listenOnLoopback(rpcHangingServer); + try { + await withLoopback(liveSocketServer, async (port) => { + const startedAt = Date.now(); + assert.equal( + await probeDaemonResponsive(daemonInfo({ port, httpPort })), + true, + 'the socket answer alone is the finding; the silent HTTP leg must not outweigh it', + ); + assert.ok(healthRequested.value, 'a verdict with no request on the HTTP leg skipped the leg'); + assert.ok( + Date.now() - startedAt < LIVENESS_PROBE_BUDGET_MS, + 'the affirmative short-circuits instead of waiting the silent leg out', + ); + }); + } finally { + await closeLoopbackServer(rpcHangingServer); + } +}); + +// `readDaemonInfo` accepts any positive integer port, and BOTH transports throw synchronously on +// one out of range (`ERR_SOCKET_BAD_PORT`). The HTTP leg builds detached, so its throw would +// surface as an unhandled rejection; the socket leg builds synchronously, so its throw would +// escape `probeDaemonResponsive` altogether — replacing the caller's timeout error with a crash +// from the recovery path and discarding the other leg's answer. +async function expectMalformedLegAnswersNotThisOne( + info: Parameters[0], + expectResponsive: boolean, + name: string, +): Promise { + const probe = probeDaemonResponsive(info); + const rejection = probe.then( + () => null, + (error: unknown) => error, + ); + assert.equal(await probe, expectResponsive, `${name}: the malformed leg answers 'not this one'`); + assert.equal(await rejection, null, `${name}: the probe never rejects on a malformed record`); +} + +// These two rows bind no listener at all: a loopback guard here would let an environment that +// cannot bind silently skip the socket-leg escape regression this file exists to hold. +test('a malformed port is an endpoint that did not answer, on either leg and never a crash', async () => { + process.on('unhandledRejection', failFastOnUnhandledRejection); + try { + await expectMalformedLegAnswersNotThisOne( + daemonInfo({ httpPort: 70_000 }), + false, + 'http leg malformed', + ); + await expectMalformedLegAnswersNotThisOne( + daemonInfo({ port: 70_000 }), + false, + 'socket leg malformed', + ); + } finally { + process.off('unhandledRejection', failFastOnUnhandledRejection); + } +}); + +test('a malformed socket port does not discard the answer of a live HTTP peer', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + const liveHttp = http.createServer((_req, res) => { + res.statusCode = 200; + res.end('{}'); + }); + const httpPort = await listenOnLoopback(liveHttp); + process.on('unhandledRejection', failFastOnUnhandledRejection); + try { + await expectMalformedLegAnswersNotThisOne( + daemonInfo({ port: 70_000, httpPort }), + true, + 'socket malformed alongside a live http peer', + ); + } finally { + process.off('unhandledRejection', failFastOnUnhandledRejection); + await closeLoopbackServer(liveHttp); + } +}); +function failFastOnUnhandledRejection(error: unknown): never { + throw error; +} + +test('a slow-trickle endpoint is cut off by the absolute deadline, not extended by its dribble', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // `socket.setTimeout`/`http.request({timeout})` are IDLE timeouts: a daemon that dribbles bytes + // forever would keep the probe (and the caller's error) open indefinitely. The probe deadline + // must be absolute, so a wedged-but-trickling endpoint still gets the negative verdict by + // `LIVENESS_PROBE_BUDGET_MS`. + // NB: a server-side socket never emits `connect`, so the dribble starts immediately on accept. + // It must actually keep the client's IDLE timeout reset, or this test proves nothing about the + // absolute deadline: an idle-only seam would pass this same assertion for a socket that simply + // got nothing. + const server = net.createServer((socket) => { + socket.on('error', () => {}); + socket.write('x'); + const dribble = setInterval(() => socket.write('x'), 20); + socket.on('close', () => clearInterval(dribble)); + }); + await withLoopback(server, async (port) => { + const startedAt = Date.now(); + assert.equal(await probeDaemonResponsive(daemonInfo({ port })), false); + const elapsedMs = Date.now() - startedAt; + assert.ok( + elapsedMs >= LIVENESS_PROBE_BUDGET_MS / 2 && elapsedMs <= LIVENESS_PROBE_BUDGET_MS * 1.5, + `trickle must end at the deadline, measured ${elapsedMs}ms (budget ${LIVENESS_PROBE_BUDGET_MS}ms)`, + ); + }); +}); + +test('the socket probe asks with the session-lock-exempt inventory command', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // The probe must never queue behind the session/device work that made the original request time + // out. That is a registry claim (`sessionExecutionLockExempt`), verified here end-to-end at the + // wire: the command on the probe's request line is the inventory one. + let requestLine = ''; + const server = net.createServer((socket) => { + socket.on('error', () => {}); + socket.once('data', (chunk) => { + requestLine = String(chunk); + socket.write(`${JSON.stringify({ jsonrpc: '2.0', id: 'p', result: { ok: true } })}\n`); + }); + }); + await withLoopback(server, async (port) => { + assert.equal(await probeDaemonResponsive(daemonInfo({ port }), { session: 'worker-3' }), true); + const asked = JSON.parse(requestLine) as { command: string; session?: string }; + assert.equal(asked.command, 'session_list'); + assert.equal(asked.session, 'worker-3'); + }); +}); 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 4399bf2b3a..69a5ddd8ac 100644 --- a/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts +++ b/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts @@ -1,37 +1,32 @@ -// Production-seam coverage for the real request-timeout route. -// -// src/daemon-client/__tests__/daemon-client-timeout.test.ts covers -// `resolveRequestTimeoutHint` as a pure formatter, but a pure-formatter test cannot catch a bug in -// CLEANUP ELIGIBILITY: whether `cleanupTimedOutIosRunnerBuilds` (the Apple -// xcodebuild pkill sweep) actually runs. This file spies on the real -// process-execution seam (`runCmdSync`, @agent-device/host-kit/command) and drives an -// actual socket/HTTP timeout through `sendRequest` so the assertions exercise -// the same code path a real client does. -// -// Why cleanup eligibility must stay unconditional for local timeouts: the -// client's declared --platform is not authoritative for session-bound -// execution. `applyStripLockPolicy` (src/daemon/request-lock-policy.ts) lets -// an existing session's real device platform silently override a conflicting -// declared selector under --session-lock strip, and the common session-bound -// request omits --platform entirely. So a request declaring `platform: -// 'android'` can still legitimately execute against an Apple-bound session -// (the "rebound-session" case below), and a request with no platform at all -// (the "unknown-session" case) is the common route the original bug misled. -// A design that skips the pkill sweep based on the declared flag alone would -// skip real cleanup in the rebound case — the dangerous direction. This test -// proves the sweep fires for every local timeout except `record` (excluded by -// command, never by platform), and that the HINT text (not the cleanup) is -// what carries the platform-evidence gating. +// Timeout recovery through the real socket/HTTP route. Connection counts distinguish the timed-out +// RPC from the fresh liveness probe. Confirmed reset controls use owned registration children; +// missing or superseded identity must retain state. No control signals an unrelated process. import net from 'node:net'; import http from 'node:http'; -import path from 'node:path'; import fs from 'node:fs'; +import path from 'node:path'; import assert from 'node:assert/strict'; import { beforeEach, afterEach, test, vi } from 'vitest'; -const { mockRunCmdSync } = vi.hoisted(() => ({ mockRunCmdSync: vi.fn() })); +const { mockRunCmdSync, mockEmitDiagnostic } = vi.hoisted(() => ({ + mockRunCmdSync: vi.fn(), + mockEmitDiagnostic: vi.fn(), +})); + +// Records what the route reported while still emitting for real: the timeout's own diagnostic is +// expected, a transport-failure diagnostic for the same request is not. +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 () => { const actual = await vi.importActual( @@ -40,40 +35,55 @@ vi.mock('@agent-device/host-kit/command', async () => { return { ...actual, runCmdSync: mockRunCmdSync }; }); -import { AppError, normalizeError } from '@agent-device/kernel/errors'; -import { sleep } from '@agent-device/host-kit/retry'; -import { withDiagnosticsScope } from '@agent-device/host-kit/diagnostics'; +import { AppError } from '@agent-device/kernel/errors'; +import { PUBLIC_COMMANDS } from '@agent-device/command-registry/catalog'; import { sendRequest } from '../daemon-client-transport.ts'; import type { DaemonRequest } from '../../daemon/daemon-request.ts'; import type { DaemonInfo } from '../daemon-client-metadata.ts'; -import { resolveDaemonPaths, type DaemonPaths } from '../../daemon-resolution.ts'; -import type { DaemonRetirementResult } from '../../daemon-registration-owner.ts'; +import type { DaemonPaths } from '../../daemon-resolution.ts'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; +import { + closeLoopbackServer, + listenOnLoopback, + skipWhenLoopbackUnavailable, + type LoopbackServer, +} from '../../__tests__/test-utils/loopback.ts'; + import { spawnRegisteredDaemonFixture, waitForRegisteredDaemonFixture, - finishRegisteredDaemonFixture, finishRegisteredDaemonFixtures, } from '../../__tests__/test-utils/registered-daemon-fixture.ts'; -import { - closeLoopbackServer, - skipWhenLoopbackUnavailable, -} from '../../__tests__/test-utils/loopback.ts'; -const TIMEOUT_MS = 120; +async function registeredOwner(paths: DaemonPaths, port: number, transport: 'http' | 'socket') { + const child = spawnRegisteredDaemonFixture( + paths, + { + ...(transport === 'http' ? { httpPort: port } : { socketPort: port }), + token: 'test-token', + version: 'test', + codeOrigin: 'checkout', + codeSignature: 'test', + }, + undefined, + ); + return await waitForRegisteredDaemonFixture(paths, child); +} + +const TIMEOUT_MS = 60; -// `snapshot`'s timeout policy preserves the daemon (onTimeout !== -// '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'; +// A reset-ELIGIBLE policy command — exercised only if the probe finds the daemon unresponsive; +// `open` is the command from the motivating report (#3177). A preserve-daemon command (`snapshot`) +// makes no reset reachable, so the probe must never run for it. +const RESET_POLICY_COMMAND = PUBLIC_COMMANDS.open; +const PRESERVE_POLICY_COMMAND = PUBLIC_COMMANDS.snapshot; function dummyStatePaths(): DaemonPaths { const baseDir = path.join( mkdtempForTestSync('agent-device-timeout-route-test'), 'agent-device-timeout-route-test', ); - return { + const paths: DaemonPaths = { baseDir, infoPath: path.join(baseDir, 'daemon.json'), lockPath: path.join(baseDir, 'daemon.lock'), @@ -81,466 +91,399 @@ function dummyStatePaths(): DaemonPaths { allocationsDir: path.join(baseDir, 'allocations'), sessionsDir: path.join(baseDir, 'sessions'), }; + return paths; +} + +// The start time the seeded registration and the request's DaemonInfo agree on. The ownership +// fence only deletes a record that MATCHES the timed-out daemon on pid and start time, so a row +// that expects removal has to name both. The pid is the test process, which fails the identity +// gate (`isAgentDeviceDaemonProcess`) on its real start time: these tests prove WHICH branch the +// route takes and never signal a process. +const TEST_DAEMON_START_TIME = 'test-start-time'; +// A registration whose pid and start time BOTH differ from the timed-out daemon: proof of a +// replacement (#3125), which the ownership fence must refuse to delete. +const REPLACEMENT_DAEMON_PID = 99_999; + +function seedRegistration(paths: DaemonPaths, owner: { pid: number; startTime: string }): void { + fs.mkdirSync(paths.baseDir, { recursive: true }); + fs.writeFileSync( + paths.infoPath, + JSON.stringify({ pid: owner.pid, processStartTime: owner.startTime }), + ); } -function buildRequest(platform: 'android' | 'ios' | undefined): DaemonRequest { +function seedProtocolLockDir(paths: DaemonPaths, owner: { pid: number; startTime: string }): void { + // `daemon.lock` is the ADR 0030 directory, not a file: seeding it shaped like production means + // a reset that still deleted it out-of-band would be RECLAIMING someone else's lock, and the + // survival assertion below catches that. (The pre-fix `unlinkSync` even failed on this shape.) + fs.mkdirSync(path.join(paths.lockPath), { recursive: true }); + fs.writeFileSync( + path.join(paths.lockPath, 'owner.json'), + JSON.stringify({ pid: owner.pid, startTime: owner.startTime, acquiredAtMs: Date.now() }), + ); +} + +function buildRequest( + command: string, + platform: 'android' | 'ios' | undefined, + positionals: readonly string[] = [], +): DaemonRequest { return { token: 'test-token', session: 'default', - command: SNAPSHOT_COMMAND, - positionals: [], + command, + positionals: [...positionals], flags: platform ? { platform } : {}, meta: { requestId: 'req-timeout-route' }, }; } -function startHangingSocketServer(): Promise<{ server: net.Server; port: number }> { - return new Promise((resolve, reject) => { - const server = net.createServer((socket) => { - // Accept the connection but never write a response — forces the - // client's own request-timeout envelope to fire. - socket.on('error', () => {}); - socket.resume(); +/** + * A daemon stand-in whose FIRST connection (the RPC) is accepted and never answered, so the + * client's own envelope cuts the round trip off, and whose LATER connections — the timeout + * handler's fresh probe — either answer like a live daemon (`answer`) or are destroyed + * unanswered (`refuse`, the wedged daemon the reset path is built for: fails the probe fast + * instead of spending its window). + * + * `connections` counts accepted TCP connections, NOT requests: the probe's whole claim is that it + * asks on a FRESH connection, so an HTTP request counter — which an RPC kept alive by keep-alive + * would also advance — would let the probe pass by reusing the timed-out socket. The stand-in only + * answers `/health` on a connection after the first, so a row that reaches its kept-alive hint is + * simultaneously proving a fresh connection carried a health request. + */ +async function startStandIn( + transport: 'http' | 'socket', + afterFirst: 'answer' | 'refuse', +): Promise<{ server: LoopbackServer; port: number; connections: () => number }> { + let connections = 0; + // The TCP ordinal lives on the socket the request arrived on, so two requests sharing one + // socket (keep-alive) count as one connection: the probe's claim is a FRESH connection, and a + // request counter would advance for an RPC kept alive on the timed-out socket too. + type OrdinalSocket = { __connectionOrdinal?: number }; + const server: LoopbackServer = + transport === 'http' + ? http.createServer((req, res) => { + const connection = (req.socket as OrdinalSocket).__connectionOrdinal ?? 0; + if (connection > 1 && afterFirst === 'answer' && req.url === '/health') { + res.statusCode = 200; + res.end('{}'); + return; + } + if (connection > 1) { + res.destroy(); + return; + } + res.on('error', () => {}); + }) + : net.createServer((socket) => { + const connection = ++connections; + socket.on('error', () => {}); + if (connection === 1) return; + if (afterFirst === 'refuse') { + socket.destroy(); + return; + } + socket.on('data', () => { + socket.write( + `${JSON.stringify({ jsonrpc: '2.0', id: 'probe', result: { ok: true } })}\n`, + ); + }); + }); + if (transport === 'http') { + (server as http.Server).on('clientError', (_err, socket) => socket.destroy()); + (server as http.Server).on('connection', (socket) => { + (socket as OrdinalSocket).__connectionOrdinal = ++connections; }); - server.on('error', reject); - server.listen(0, '127.0.0.1', () => { - const address = server.address(); - if (address && typeof address === 'object') { - resolve({ server, port: address.port }); - } else { - reject(new Error('failed to bind hanging socket test server')); - } - }); - }); + } + const port = await listenOnLoopback(server); + return { server, port, connections: () => connections }; } -function startHangingHttpServer(): Promise<{ server: http.Server; port: number }> { - return new Promise((resolve, reject) => { - const server = http.createServer((_req, res) => { - // Never call res.end() — forces the client's own request-timeout - // envelope to fire instead of a real response. - res.on('error', () => {}); - }); - server.on('clientError', (_err, socket) => socket.destroy()); - server.on('error', reject); - server.listen(0, '127.0.0.1', () => { - const address = server.address(); - if (address && typeof address === 'object') { - resolve({ server, port: address.port }); - } else { - reject(new Error('failed to bind hanging http test server')); - } - }); - }); +async function expectRouteError(run: Promise, hintPattern: RegExp): Promise { + let thrown: unknown; + try { + await run; + } catch (error) { + thrown = error; + } + assert.ok(thrown instanceof AppError, 'the route must reject with a typed AppError'); + assert.match(String(thrown.details?.hint), hintPattern); + // The regression this suite exists to catch: no host-wide process sweep, whatever the request + // declared. + assert.equal(mockRunCmdSync.mock.calls.length, 0); + // The timeout settles the request and then destroys it, so the transport's `error` event lands + // AFTER the rejection. A timed-out request must not also be diagnosed as a transport FAILURE — + // that describes a canceled request as a broken host and buries the timeout's own reason. The + // destroy surfaces on a later tick, so give it one before reading the record. + await new Promise((resolve) => setImmediate(() => setImmediate(resolve))); + assert.deepEqual( + mockEmitDiagnostic.mock.calls + .map(([entry]) => (entry as { phase?: string })?.phase) + .filter((phase) => phase === 'daemon_request_socket_error'), + [], + 'a timeout must not also diagnose a transport failure', + ); } beforeEach(() => { mockRunCmdSync.mockReset(); + mockEmitDiagnostic.mockReset(); }); 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 - // existing Apple session: the client declared `platform: 'android'`, but - // real Apple xcodebuild work was in flight and the pkill sweep kills it. - mockRunCmdSync.mockImplementation((cmd: string) => - cmd === 'pkill' - ? { exitCode: 0, stdout: '', stderr: '' } - : { exitCode: 1, stdout: '', stderr: '' }, - ); - - const { server, port } = await startHangingSocketServer(); - try { - const info: DaemonInfo = { port, token: 'test-token', pid: process.pid }; - const req = buildRequest('android'); +type RouteRow = Readonly<{ + name: string; + transport: 'http' | 'socket' | 'remote'; + command: string; + platform: 'android' | 'ios' | undefined; + positionals?: readonly string[]; + afterFirst: 'answer' | 'refuse'; + hintPattern: RegExp; + // The connections the stand-in must observe: 1 = the RPC only (route never probed), + // 2 = RPC + probe. + connections: 1 | 2; + /** Whether this row's verdict is the reset branch, which clears the daemon's metadata. */ + resets: boolean; +}>; + +// The `ios` declarations are deliberate: eligibility never keyed off the declared platform, and +// neither may recovery. The prepare follow-up on the snapshot row is keyed on a declared Apple +// platform — the only Apple evidence this route can back up now that the sweep is gone. +const ROUTE_ROWS: readonly RouteRow[] = [ + { + name: 'keeps a responsive daemon alive (the motivating #3177 case)', + transport: 'http', + command: RESET_POLICY_COMMAND, + platform: 'ios', + afterFirst: 'answer', + hintPattern: /The timed-out open request was canceled; the daemon was kept alive/, + connections: 2, + resets: false, + }, + { + name: 'is proven responsive over the socket transport', + transport: 'socket', + command: RESET_POLICY_COMMAND, + platform: undefined, + afterFirst: 'answer', + hintPattern: /the daemon was kept alive so the session can still be closed or inspected/, + connections: 2, + resets: false, + }, + { + name: 'resets a daemon that answers no probe endpoint', + transport: 'socket', + command: RESET_POLICY_COMMAND, + platform: undefined, + afterFirst: 'refuse', + hintPattern: /The daemon did not answer the liveness probe and was reset after the timeout/, + connections: 2, + resets: true, + }, + { + // `snapshot` declares preserve-daemon: no reset is reachable, so the probe would be pure + // latency and the route must skip it entirely. + name: 'skips the probe for a preserve-policy command', + transport: 'http', + command: PRESERVE_POLICY_COMMAND, + platform: 'ios', + afterFirst: 'answer', + hintPattern: + /The timed-out snapshot request was canceled; the daemon was kept alive.*prepare ios-runner/s, + connections: 1, + resets: false, + }, + { + // `record` declares preserve-daemon (#3199), so a LOCAL timed-out `record stop` reaches the + // retry hint without the probe ever running: the export the surviving daemon may still be + // finishing must not lose the runner to a sweep either (#3177 removed the sweep for all + // commands; this row pins that it is gone for the recorder too, on the real action positional). + name: 'names the record stop retry without probing a preserve-policy recorder', + transport: 'socket', + command: PUBLIC_COMMANDS.record, + platform: 'ios', + positionals: ['stop'], + afterFirst: 'answer', + hintPattern: + /^The daemon may still be exporting the recording\. Run agent-device record stop --session default again/, + connections: 1, + resets: false, + }, + { + // A remote client cannot reset anything on the daemon's host: no probe window, no sweep. + name: 'keeps a remote timeout declarative', + transport: 'remote', + command: RESET_POLICY_COMMAND, + platform: 'android', + afterFirst: 'answer', + hintPattern: /verify the remote daemon URL, auth token, and remote host logs/, + connections: 1, + resets: false, + }, +]; + +function standInInfo(row: RouteRow, port: number): DaemonInfo { + const identity = { + token: 'test-token', + pid: process.pid, + processStartTime: TEST_DAEMON_START_TIME, + }; + if (row.transport === 'remote') return { ...identity, baseUrl: `http://127.0.0.1:${port}` }; + return row.transport === 'http' ? { ...identity, httpPort: port } : { ...identity, port }; +} - await assert.rejects( - sendRequest(info, req, 'socket', dummyStatePaths(), TIMEOUT_MS), - (error: unknown) => { - assert.ok(error instanceof AppError); - assert.match(error.details?.hint as string, /Apple runner work was aborted when detected/); - return true; - }, +for (const row of ROUTE_ROWS) { + test(`request-timeout route: ${row.name}`, async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + const daemon = await startStandIn( + row.transport === 'remote' ? 'http' : row.transport, + row.afterFirst, ); - } finally { - server.close(); - } - - // The eligibility assertion: cleanup ran (all three kill patterns - // attempted) even though the request declared a non-Apple platform. A - // design that skips cleanup based on the declared flag would fail this. - const pkillCalls = mockRunCmdSync.mock.calls.filter(([cmd]) => cmd === 'pkill'); - assert.equal(pkillCalls.length, 3); - - // The session-xctestrun pattern is pinned by bytes, not derived from the runner's writer module: - // a client version in the field already pkills this exact string, and it must keep selecting - // launches that older writers named, since it cannot know which version started a timed-out - // launch. Deriving it would move this sweep off those names on any rename. - const sessionPattern = pkillCalls - .map(([, args]) => String(args?.[1])) - .find((pattern) => pattern.includes('session')); - assert.equal(sessionPattern, String.raw`xcodebuild .*AgentDeviceRunner\.env\.session-`); - assert.equal( - new RegExp(sessionPattern).test( - 'xcodebuild test-without-building -xctestrun /d/AgentDeviceRunner.env.session-SIM-1-owner-1-ff-8123.xctestrun', - ), - true, - ); - assert.equal( - new RegExp(sessionPattern).test( - 'xcodebuild test-without-building -xctestrun /d/AgentDeviceRunner.env.session-SIM-1-8123.xctestrun', - ), - true, - ); -}); - -test('http timeout: pkill cleanup still runs for an undeclared platform (unknown-session case) that terminates nothing, and the hint stays platform-neutral', async () => { - // Simulates the common session-bound request that never repeats - // --platform, on a real Android/web/Harmony session: no processes match - // the Apple-specific kill patterns. - mockRunCmdSync.mockImplementation(() => ({ exitCode: 1, stdout: '', stderr: '' })); + // Preserved rows retain seeded metadata. Reset rows acquire a real child-owned registration. + const statePaths = dummyStatePaths(); + const owned = { pid: process.pid, startTime: TEST_DAEMON_START_TIME }; + if (!row.resets) { + seedRegistration(statePaths, owned); + seedProtocolLockDir(statePaths, owned); + } + const registered = row.resets + ? await registeredOwner(statePaths, daemon.port, row.transport === 'http' ? 'http' : 'socket') + : undefined; + try { + const info = registered ?? standInInfo(row, daemon.port); + await expectRouteError( + sendRequest( + info, + buildRequest(row.command, row.platform, row.positionals), + row.transport === 'remote' ? 'http' : row.transport, + statePaths, + TIMEOUT_MS, + ), + row.hintPattern, + ); + assert.equal(daemon.connections(), row.connections, 'probe-vs-skip is a route decision'); + assert.equal( + fs.existsSync(statePaths.infoPath), + !row.resets, + row.resets ? 'the reset clears the owned registration' : 'no reset touches metadata', + ); + assert.equal( + fs.existsSync(statePaths.lockPath), + !row.resets, + 'only confirmed retirement releases the owned lock', + ); + } finally { + await closeLoopbackServer(daemon.server); + } + }); +} - const { server, port } = await startHangingHttpServer(); +test('request-timeout route: a reset never deletes a registration a replacement daemon published', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + // The probe window is exactly when another client can start a replacement daemon and publish + // ITS record (#3177 review, on top of #3125's fence). The timed-out request's info still names + // the wedged pid, so the reset must read the record before clearing it: here it names a + // different owner, and deleting it would orphan the live replacement — the same class of + // host-scoped damage this PR removes everywhere else. + const daemon = await startStandIn('socket', 'refuse'); + const statePaths = dummyStatePaths(); + seedRegistration(statePaths, { pid: REPLACEMENT_DAEMON_PID, startTime: 'replacement-start' }); + seedProtocolLockDir(statePaths, { pid: process.pid, startTime: TEST_DAEMON_START_TIME }); try { - const info: DaemonInfo = { httpPort: port, token: 'test-token', pid: process.pid }; - const req = buildRequest(undefined); - - await assert.rejects( - sendRequest(info, req, 'http', dummyStatePaths(), TIMEOUT_MS), - (error: unknown) => { - assert.ok(error instanceof AppError); - const hint = error.details?.hint as string; - assert.doesNotMatch(hint, /Apple/); - assert.match( - hint, - /The timed-out snapshot request was canceled; the daemon was kept alive/, - ); - return true; - }, + await expectRouteError( + sendRequest( + { + port: daemon.port, + token: 'test-token', + pid: process.pid, + processStartTime: TEST_DAEMON_START_TIME, + }, + buildRequest(RESET_POLICY_COMMAND, undefined), + 'socket', + statePaths, + TIMEOUT_MS, + ), + /The daemon could not be safely retired/, ); - } finally { - server.close(); - } - - // Cleanup still ran — this is the regression this suite exists to catch: - // an eligibility design keyed off the (here, absent) declared platform - // would either skip cleanup entirely or — under the original unconditional - // hint — falsely claim Apple involvement anyway. Neither happens here. - const pkillCalls = mockRunCmdSync.mock.calls.filter(([cmd]) => cmd === 'pkill'); - assert.equal(pkillCalls.length, 3); -}); - -test('http timeout: an explicitly declared Apple platform keeps the Apple hint even when the sweep terminates nothing', async () => { - mockRunCmdSync.mockImplementation(() => ({ exitCode: 1, stdout: '', stderr: '' })); - - const { server, port } = await startHangingHttpServer(); - try { - const info: DaemonInfo = { httpPort: port, token: 'test-token', pid: process.pid }; - const req = buildRequest('ios'); - - await assert.rejects( - sendRequest(info, req, 'http', dummyStatePaths(), TIMEOUT_MS), - (error: unknown) => { - assert.ok(error instanceof AppError); - assert.match(error.details?.hint as string, /Apple runner work was aborted when detected/); - return true; - }, + assert.ok( + fs.existsSync(statePaths.infoPath), + 'the replacement keeps the registration it published', ); - } finally { - server.close(); - } - - const pkillCalls = mockRunCmdSync.mock.calls.filter(([cmd]) => cmd === 'pkill'); - assert.equal(pkillCalls.length, 3); -}); - -test('remote HTTP timeout never runs the Apple pkill cleanup and uses the remote-specific hint', async () => { - mockRunCmdSync.mockImplementation(() => ({ exitCode: 0, stdout: '', stderr: '' })); - - const { server, port } = await startHangingHttpServer(); - try { - const info: DaemonInfo = { - baseUrl: `http://127.0.0.1:${port}`, - token: 'test-token', - pid: process.pid, - }; - const req = buildRequest('android'); - - await assert.rejects( - sendRequest(info, req, 'http', dummyStatePaths(), TIMEOUT_MS), - (error: unknown) => { - assert.ok(error instanceof AppError); - assert.match( - error.details?.hint as string, - /verify the remote daemon URL, auth token, and remote host logs/, - ); - return true; - }, + assert.equal( + JSON.parse(fs.readFileSync(statePaths.infoPath, 'utf8')).pid, + REPLACEMENT_DAEMON_PID, ); } finally { - server.close(); + await closeLoopbackServer(daemon.server); } - - assert.equal(mockRunCmdSync.mock.calls.length, 0); }); -async function waitForForceStop(requested: Promise): Promise { - let timer: ReturnType | undefined; - try { - await Promise.race([ - requested, - new Promise((_resolve, reject) => { - timer = setTimeout(() => reject(new Error('force stop was not requested')), 1_500); - }), - ]); - } finally { - clearTimeout(timer); - } -} - -function timeoutDiagnostic(paths: DaemonPaths): Record { - const events = fs - .readFileSync(path.join(paths.baseDir, 'timeout-diagnostics.ndjson'), 'utf8') - .trim() - .split('\n') - .map((line) => JSON.parse(line)); - const event = events.find((entry) => entry.phase === 'daemon_request_timeout'); - assert.ok(event); - return event.data; -} - -function assertForcedRetirement( - retirement: DaemonRetirementResult | undefined, - removalFails: boolean, -): void { - if (removalFails) { - assert.ok(retirement?.status === 'retained'); - assert.equal(retirement.reason, 'retirement-unconfirmed'); - assert.ok(retirement.termination?.status === 'exited'); - assert.equal(retirement.termination.mode, 'forced'); - } else { - assert.ok(retirement?.status === 'retired'); - assert.equal(retirement.termination.mode, 'forced'); - } -} - -for (const transport of ['socket', 'http'] as const) { - for (const removalFails of [false, true]) { - test(`${transport} timeout awaits force exit when metadata removal ${removalFails ? 'fails' : 'succeeds'}`, 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); - }); - const actualUnlink = fs.unlinkSync.bind(fs); - const remove = vi.spyOn(fs, 'unlinkSync').mockImplementation((file) => { - if (removalFails && file === paths.infoPath) - throw Object.assign(new Error('retained registration control'), { code: 'EACCES' }); - actualUnlink(file); - }); - try { - const info = await waitForRegisteredDaemonFixture(paths, child); - outcome = withDiagnosticsScope( - { debug: true, logPath: path.join(paths.baseDir, 'timeout-diagnostics.ndjson') }, - () => - sendRequest( - info, - { ...buildRequest(undefined), command: 'open' }, - transport, - paths, - TIMEOUT_MS, - ), - ).then( - () => assert.fail('hanging request unexpectedly succeeded'), - (error: unknown) => { - settled = true; - return error; - }, - ); - await waitForForceStop(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'); - assertForcedRetirement( - error.details?.retirement as DaemonRetirementResult | undefined, - removalFails, - ); - assert.equal(fs.existsSync(paths.infoPath), removalFails); - assert.equal(fs.existsSync(paths.lockPath), false); - assert.equal(fs.existsSync(paths.baseDir), true); - assert.equal(timeoutDiagnostic(paths).daemonPreservedAfterTimeout, false); - assert.equal(timeoutDiagnostic(paths).daemonPidForceKilled, true); - assert.equal(mockRunCmdSync.mock.calls.filter(([cmd]) => cmd === 'pkill').length, 3); - } finally { - remove.mockRestore(); - kill.mockRestore(); - if (!exited) actualKill(child.pid, 'SIGKILL'); - await child.exited; - await outcome; - await finishRegisteredDaemonFixture(paths.baseDir); - await closeLoopbackServer(endpoint.server); - } - }); - } -} - -test('timeout retains a live registration without captured birth proof and reports that outcome', async (t) => { +test('a refused timeout fallback preserves the timeout without an unhandled rejection', 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'); + // The reset branch's fallback: a SIGKILL the kernel refuses goes through the confirmed-retirement + // path (#3126), and a retirement that cannot confirm the daemon's exit must surface as the + // timeout the caller already has — not as an unhandled rejection and not as a different error. + // The stand-in refuses the probe's fresh connection, so this row really does reach the reset. + const daemon = await startStandIn('socket', 'refuse'); + const paths = dummyStatePaths(); + const info = await registeredOwner(paths, daemon.port, 'socket'); + const actualKill = process.kill.bind(process); + const killSpy = vi.spyOn(process, 'kill').mockImplementation((pid, signal) => { + if (pid === info.pid && signal !== 0) + throw Object.assign(new Error('refused'), { code: 'EPERM' }); + return actualKill(pid, signal); + }); try { - const info = await waitForRegisteredDaemonFixture(paths, child); - 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( - withDiagnosticsScope( - { debug: true, logPath: path.join(paths.baseDir, 'timeout-diagnostics.ndjson') }, - () => - sendRequest( - { ...info, processStartTime: undefined }, - { ...buildRequest(undefined), command: 'open' }, - 'http', - paths, - TIMEOUT_MS, - ), - ), + sendRequest(info, buildRequest(RESET_POLICY_COMMAND, undefined), 'socket', paths, TIMEOUT_MS), (error: unknown) => { 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 === '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/); + assert.equal(error.details?.reason, 'daemon_transport_timeout'); return true; }, ); - assert.equal(timeoutDiagnostic(paths).daemonPreservedAfterTimeout, true); - 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, - ); + await new Promise((resolve) => setImmediate(resolve)); + assert.ok(killSpy.mock.calls.some(([pid, signal]) => pid === info.pid && signal === 'SIGKILL')); + assert.ok(fs.existsSync(paths.infoPath)); + assert.ok(fs.existsSync(paths.lockPath)); } finally { - kill.mockRestore(); - await finishRegisteredDaemonFixture(paths.baseDir); - await closeLoopbackServer(endpoint.server); + await closeLoopbackServer(daemon.server); } }); -test('a local record stop timeout leaves the exporting daemon and runner alive and names the retry', async (t) => { +test('request-timeout route: the reset acts on the resolved state paths, not the request flags', async (t) => { if (await skipWhenLoopbackUnavailable(t)) return; - mockRunCmdSync.mockReturnValue({ exitCode: 0, stdout: '', stderr: '' }); - const endpoint = await startHangingSocketServer(); - const paths = resolveDaemonPaths(mkdtempForTestSync('agent-device-record-timeout-owner-')); - const child = spawnRegisteredDaemonFixture( - paths, - { - socketPort: endpoint.port, - token: 'test-token', - version: 'test', - codeOrigin: 'checkout', - codeSignature: 'test', - }, - undefined, - ); - const kill = vi.spyOn(process, 'kill'); + // Path resolution happens upstream of `sendRequest`, which receives the state paths the caller + // resolved. The reset must clear THAT registration and leave every other state dir alone — even + // one a request flag names, which this route never reads. (Moved from daemon-client-lifecycle + // coverage, where the timeout route lived before #3177 scoped it.) + const daemon = await startStandIn('http', 'refuse'); + const resolvedPaths = dummyStatePaths(); + const requestFlagPaths = dummyStatePaths(); + const info = await registeredOwner(resolvedPaths, daemon.port, 'http'); + seedRegistration(requestFlagPaths, { pid: REPLACEMENT_DAEMON_PID, startTime: 'other-owner' }); + seedProtocolLockDir(requestFlagPaths, { pid: REPLACEMENT_DAEMON_PID, startTime: 'other-owner' }); try { - const info = await waitForRegisteredDaemonFixture(paths, child); - const registration = fs.readFileSync(paths.infoPath, 'utf8'); - await assert.rejects( + await expectRouteError( sendRequest( info, { - ...buildRequest('ios'), - session: 'e2e-ios-0', - command: 'record', - positionals: ['stop'], + ...buildRequest(RESET_POLICY_COMMAND, undefined), + flags: { stateDir: requestFlagPaths.baseDir }, }, - 'socket', - paths, + 'http', + resolvedPaths, TIMEOUT_MS, ), - (error: unknown) => { - assert.ok(error instanceof AppError); - assert.equal(error.details?.reason, 'daemon_transport_timeout'); - assert.match( - error.details?.hint as string, - /^The daemon may still be exporting the recording\. Run agent-device record stop --session e2e-ios-0 again/, - ); - return true; - }, + /The daemon did not answer the liveness probe and was reset after the timeout/, ); - assert.equal(mockRunCmdSync.mock.calls.length, 0); - assert.equal(kill.mock.calls.length, 0); - assert.equal(process.kill(child.pid, 0), true); - assert.equal(fs.readFileSync(paths.infoPath, 'utf8'), registration); - assert.equal(fs.existsSync(paths.lockPath), true); + assert.ok(!fs.existsSync(resolvedPaths.infoPath), 'the resolved registration is cleared'); + assert.ok(fs.existsSync(requestFlagPaths.infoPath), 'the flag-named state dir is untouched'); + assert.ok(fs.existsSync(requestFlagPaths.lockPath), 'the flag-named lock is untouched'); } finally { - kill.mockRestore(); - await finishRegisteredDaemonFixture(paths.baseDir); - await closeLoopbackServer(endpoint.server); + await closeLoopbackServer(daemon.server); } }); diff --git a/src/daemon-client/__tests__/daemon-client-timeout.test.ts b/src/daemon-client/__tests__/daemon-client-timeout.test.ts index 46c330c2b6..f1559d7f5b 100644 --- a/src/daemon-client/__tests__/daemon-client-timeout.test.ts +++ b/src/daemon-client/__tests__/daemon-client-timeout.test.ts @@ -1,109 +1,81 @@ // The pure hint formatter in src/daemon-client/daemon-client-timeout.ts: what a timed-out request // tells the caller to do next. daemon-client-timeout-route.test.ts covers the same route at its -// production seam, where cleanup eligibility is decided; these assertions only fix the wording. +// production seam, where the liveness-gated recovery is decided; these assertions only fix wording. import assert from 'node:assert/strict'; import { test } from 'vitest'; import { resolveRequestTimeoutHint } from '../daemon-client-timeout.ts'; -test('request timeout hint only names Apple runner cleanup on actual evidence', () => { - // Before this change, handleRequestTimeout emitted Apple-specific hint - // wording for EVERY local timeout, regardless of the request's - // --platform. That was misleading for Android/web/Harmony sessions, which - // never had any Apple runner work to abort. - // - // The fix is evidence-based, not platform-guess-based: - // `appleCleanupEvidence` is true only when the request declared an - // AFFIRMATIVELY Apple platform (apple/ios/macos) or the pkill cleanup - // itself terminated a matching process — never from an undeclared or - // declared-non-Apple platform alone. (Why not trust the declared platform - // directly: it is not authoritative for session-bound execution — see - // `handleRequestTimeout`'s comment and the production-seam coverage in - // daemon-client-timeout-route.test.ts for the cleanup-eligibility half of - // this contract that this pure formatter test cannot prove.) - - // appleCleanupEvidence: true keeps the exact historical wording — nothing - // regresses for the true-Apple case. +test('request timeout hint names an Apple runner only on a declared Apple platform', () => { + // The hint used to derive its Apple claim partly from a host-wide `pkill` sweep that counted + // whatever it terminated (#1751). That sweep ended every session on the host and is gone + // (#3177), so the only Apple evidence this route has left is a selector that says so outright. assert.equal( resolveRequestTimeoutHint({ remote: false, resetDaemon: false, command: 'press', - appleCleanupEvidence: true, - }), - 'Retry with --debug and check daemon diagnostics logs. The timed-out press request was canceled and Apple runner work was aborted when detected; the daemon was kept alive so the session can still be closed or inspected.', - ); - assert.equal( - resolveRequestTimeoutHint({ - remote: false, - resetDaemon: true, - command: 'open', - appleCleanupEvidence: true, + applePlatformDeclared: true, }), - 'Retry with --debug and check daemon diagnostics logs. Timed-out Apple runner xcodebuild processes were terminated when detected.', + 'Retry with --debug and check daemon diagnostics logs. The timed-out press request was canceled; the daemon was kept alive so the session can still be closed or inspected.', ); assert.equal( resolveRequestTimeoutHint({ remote: false, resetDaemon: false, command: 'snapshot', - appleCleanupEvidence: true, + applePlatformDeclared: true, }), - 'Retry with --debug and check daemon diagnostics logs. The timed-out snapshot request was canceled and Apple runner work was aborted when detected; the daemon was kept alive so the session can still be closed or inspected. If this was the first Apple-platform snapshot on the device, run agent-device prepare ios-runner with the same --platform before snapshot/test so runner startup is handled explicitly.', + 'Retry with --debug and check daemon diagnostics logs. The timed-out snapshot request was canceled; the daemon was kept alive so the session can still be closed or inspected. If this was the first Apple-platform snapshot on the device, run agent-device prepare ios-runner with the same --platform before snapshot/test so runner startup is handled explicitly.', ); - // appleCleanupEvidence: false — no Apple-runner claim in any branch, and - // the Apple-only iOS-prepare follow-up drops entirely. This is the - // motivating fix: it fires equally whether the platform was declared - // non-Apple OR left undeclared (the common session-bound case), because - // neither is Apple evidence on its own. - assert.equal( - resolveRequestTimeoutHint({ - remote: false, - resetDaemon: false, - command: 'press', - appleCleanupEvidence: false, - }), - 'Retry with --debug and check daemon diagnostics logs. The timed-out press request was canceled; the daemon was kept alive so the session can still be closed or inspected.', - ); - assert.equal( - resolveRequestTimeoutHint({ - remote: false, - resetDaemon: true, - command: 'open', - appleCleanupEvidence: false, - }), - 'Retry with --debug and check daemon diagnostics logs. The daemon was reset after the timeout.', - ); + // An undeclared or declared non-Apple platform names no Apple runner in any branch, and the + // Apple-only prepare follow-up drops entirely. assert.equal( resolveRequestTimeoutHint({ remote: false, resetDaemon: false, command: 'snapshot', - appleCleanupEvidence: false, + applePlatformDeclared: false, }), 'Retry with --debug and check daemon diagnostics logs. The timed-out snapshot request was canceled; the daemon was kept alive so the session can still be closed or inspected.', ); + // A reset is now reported as what it was decided from: an unanswered liveness probe. The reset + // SIGKILLs the daemon pid only, so the hint claims nothing about runner children it cannot + // prove stopped — and a declared Apple platform changes the reset wording not at all. + for (const applePlatformDeclared of [true, false]) { + assert.equal( + resolveRequestTimeoutHint({ + remote: false, + resetDaemon: true, + command: 'open', + applePlatformDeclared, + }), + 'Retry with --debug and check daemon diagnostics logs. The daemon did not answer the liveness probe and was reset after the timeout.', + ); + } + // Remote requests were never Apple-specific and stay evidence-independent. assert.equal( resolveRequestTimeoutHint({ remote: true, resetDaemon: false, command: 'press', - appleCleanupEvidence: false, + applePlatformDeclared: false, }), 'Retry with --debug and verify the remote daemon URL, auth token, and remote host logs.', ); }); test('a timed-out record stop on a surviving daemon names the retry that returns the export', () => { + // A remote client never touches the daemon's host, so its daemon survives by construction. assert.equal( resolveRequestTimeoutHint({ remote: true, resetDaemon: false, command: 'record', - appleCleanupEvidence: false, + applePlatformDeclared: false, action: 'stop', session: 'recording', }), @@ -114,34 +86,36 @@ test('a timed-out record stop on a surviving daemon names the retry that returns remote: true, resetDaemon: false, command: 'record', - appleCleanupEvidence: false, + applePlatformDeclared: false, action: 'stop', }), 'The remote daemon may still be exporting the recording. Run agent-device record stop again to wait for that export and receive the completed recording.', ); - // A local daemon preserved across the timeout may still be exporting too. + // A LOCAL daemon preserved across the timeout (#3199 declares `record` preserve-daemon) may + // still be exporting too, so it gets the same retry instead of the generic kept-alive wording. assert.equal( resolveRequestTimeoutHint({ remote: false, resetDaemon: false, command: 'record', - appleCleanupEvidence: true, + applePlatformDeclared: false, action: 'stop', session: 'recording', }), 'The daemon may still be exporting the recording. Run agent-device record stop --session recording again to wait for that export and receive the completed recording.', ); - // A reset daemon is no longer exporting, so no keep-exporting promise is made. + // A daemon the probe proved unresponsive was reset, and a reset daemon is no longer exporting: + // no keep-exporting promise is made, and the wording says what happened to it instead. assert.equal( resolveRequestTimeoutHint({ remote: false, resetDaemon: true, command: 'record', - appleCleanupEvidence: false, + applePlatformDeclared: false, action: 'stop', session: 'recording', }), - 'Retry with --debug and check daemon diagnostics logs. The daemon was reset after the timeout.', + 'Retry with --debug and check daemon diagnostics logs. The daemon did not answer the liveness probe and was reset after the timeout.', ); // `record start` runs no export, so it keeps the generic remote wording. assert.equal( @@ -149,7 +123,7 @@ test('a timed-out record stop on a surviving daemon names the retry that returns remote: true, resetDaemon: false, command: 'record', - appleCleanupEvidence: false, + applePlatformDeclared: false, action: 'start', session: 'recording', }), diff --git a/src/daemon-client/daemon-client-liveness-probe.ts b/src/daemon-client/daemon-client-liveness-probe.ts new file mode 100644 index 0000000000..66a3605ece --- /dev/null +++ b/src/daemon-client/daemon-client-liveness-probe.ts @@ -0,0 +1,170 @@ +import net from 'node:net'; +import { INTERNAL_COMMANDS } from '@agent-device/command-registry/catalog'; +import { createRequestId } from '@agent-device/host-kit/diagnostics'; +import { loadNodeHttpRequester, consumeTextLines } from '@agent-device/host-kit/transport'; +import type { DaemonRequest } from '../daemon/daemon-request.ts'; +import type { DaemonInfo } from './daemon-client-metadata.ts'; + +// The post-timeout liveness probe (#3177): the question a daemon reset may no longer assume. The +// daemon owns request-scoped recovery, so a reset is only justified for a daemon that answers +// nothing, and it must be asked on FRESH connections — the timed-out connection is the one the +// transport already tore down to cancel that request. + +// The extra latency a caller of a reset-eligible timed-out request can pay before its error lands. +// An order of magnitude under the narrowest reset-eligible envelope (pinned in +// `__tests__/daemon-client-liveness-probe.test.ts`): a wrong "unresponsive" verdict kills a live +// shared daemon. +export const LIVENESS_PROBE_BUDGET_MS = 1_000; + +/** + * Whether this local daemon answers anything within the probe budget: its HTTP health route, or + * one sessionless RPC on its socket. An answer on either leg proves the daemon's event loop serves + * fresh requests; a refused or silent leg answers "not this one", not "not the daemon". + * + * Both legs run CONCURRENTLY against one shared deadline: the finding is per-daemon, not per-leg, + * and a leg that hangs must not spend the window the other leg needs. The first AFFIRMATIVE is + * unrevocable and settles at once. The deadline is ABSOLUTE — a leg whose endpoint trickles bytes + * would outlive an idle timeout, and the negative verdict must land inside the window. A leg never + * rejects: a probe that cannot ask is an endpoint that did not answer. `session` rides along so an + * isolation-scoped daemon routes the probe like the request it follows. + */ +export async function probeDaemonResponsive( + info: DaemonInfo, + params: Readonly<{ session?: string }> = {}, +): Promise { + const deadlineAtMs = Date.now() + LIVENESS_PROBE_BUDGET_MS; + const legs: ProbeLeg[] = []; + if (info.httpPort) legs.push(probeHttpHealth(info.httpPort, deadlineAtMs)); + if (info.port) legs.push(probeSocketRpc(info.port, info.token, params.session, deadlineAtMs)); + if (legs.length === 0) return false; + return await new Promise((resolve) => { + let outstanding = legs.length; + const settle = (answered: boolean): void => { + if (answered) { + for (const leg of legs) leg.cancel(); + resolve(true); + return; + } + if (--outstanding === 0) resolve(false); + }; + for (const leg of legs) leg.answer.then(settle, () => settle(false)); + }); +} + +type ProbeLeg = Readonly<{ answer: Promise; cancel: () => void }>; + +function createProbeLeg(deadlineAtMs: number): { + leg: ProbeLeg; + settle: (answered: boolean) => void; + attach: (destroyTransport: () => void) => void; +} { + let settled = false; + let destroy: (() => void) | undefined; + let resolveAnswer: (answered: boolean) => void = () => {}; + const answer = new Promise((resolve) => { + resolveAnswer = resolve; + }); + const deadlineHandle = setTimeout(() => settle(false), Math.max(1, deadlineAtMs - Date.now())); + deadlineHandle.unref?.(); + function settle(answered: boolean): void { + if (settled) return; + settled = true; + clearTimeout(deadlineHandle); + destroy?.(); + resolveAnswer(answered); + } + return { + leg: { answer, cancel: () => settle(false) }, + settle, + attach: (destroyTransport) => { + destroy = destroyTransport; + // The deadline can close before a leg builds its transport, so a late arrival is torn down. + if (settled) destroyTransport(); + }, + }; +} + +function probeHttpHealth(httpPort: number, deadlineAtMs: number): ProbeLeg { + const probe = createProbeLeg(deadlineAtMs); + // `readDaemonInfo` accepts any positive integer port, and `transport.request` throws on one out + // of range; this leg builds detached, so the throw would reach the caller as an unhandled + // rejection. + void (async () => { + try { + const transport = await loadNodeHttpRequester('http:'); + const request = transport.request( + // `agent: false` is load-bearing, not hygiene: since Node 19 the global agent keeps sockets + // alive, so a pooled request would ride an idle socket from an earlier command's health read + // instead of the FRESH connection this probe exists to make — and a half-torn-down pooled + // socket would answer "silent", resetting a live daemon. + { host: '127.0.0.1', port: String(httpPort), path: '/health', method: 'GET', agent: false }, + (res) => { + // Any status is an answer: the health route is served before any request handling, so + // reaching it proves the event loop serves fresh requests. Drain so the response cannot + // hold the socket open past the finding. + res.resume(); + probe.settle(true); + }, + ); + request.on('error', () => probe.settle(false)); + probe.attach(() => request.destroy()); + request.end(); + } catch { + probe.settle(false); + } + })(); + return probe.leg; +} + +function probeSocketRpc( + port: number, + token: string, + session: string | undefined, + deadlineAtMs: number, +): ProbeLeg { + const probe = createProbeLeg(deadlineAtMs); + // The same malformed-record shape as the HTTP leg, and this one builds synchronously: a throw + // here would escape `probeDaemonResponsive` altogether, replacing the caller's timeout error + // with a crash from the recovery path and discarding the other leg's answer. + let socket: net.Socket; + try { + socket = net.createConnection({ host: '127.0.0.1', port }); + } catch { + probe.settle(false); + return probe.leg; + } + probe.attach(() => socket.destroy()); + const request = buildLivenessProbeRequest(token, session); + let buffer = ''; + socket.setEncoding('utf8'); + socket.on('connect', () => { + socket.write(`${JSON.stringify(request)}\n`); + }); + socket.on('data', (chunk) => { + const parsed = consumeTextLines(buffer, chunk); + buffer = parsed.buffer; + // Any answer line counts, including an error for a probe the daemon declined: the finding is + // that it reads and answers, not that it accepted this command. + if (parsed.lines.length > 0) probe.settle(true); + }); + socket.on('error', () => probe.settle(false)); + socket.on('close', () => probe.settle(false)); + return probe.leg; +} + +/** + * The probe's request: `session_list`, which derives `sessionExecutionLockExempt: true` (plus + * lease-admission and selector-validation exemptions), so it never queues behind the session or + * device locks the timed-out request's work may still hold. The `session` field is required on + * the wire; an omitted one falls back to the daemon's own routing default. + */ +function buildLivenessProbeRequest(token: string, session: string | undefined): DaemonRequest { + return { + token, + session: session || 'default', + command: INTERNAL_COMMANDS.sessionList, + positionals: [], + flags: {}, + meta: { requestId: createRequestId() }, + }; +} diff --git a/src/daemon-client/daemon-client-timeout.ts b/src/daemon-client/daemon-client-timeout.ts index 547c5f2799..5ae3336658 100644 --- a/src/daemon-client/daemon-client-timeout.ts +++ b/src/daemon-client/daemon-client-timeout.ts @@ -1,7 +1,5 @@ import { AppError } from '@agent-device/kernel/errors'; -import { runCmdSync } from '@agent-device/host-kit/command'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; - 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'; @@ -9,31 +7,16 @@ import type { DaemonPaths } from '../daemon-resolution.ts'; import type { PlatformSelector } from '@agent-device/kernel/device'; import type { DaemonInfo } from './daemon-client-metadata.ts'; -const IOS_RUNNER_XCODEBUILD_KILL_PATTERNS = [ - 'xcodebuild .*AgentDeviceRunnerUITests/RunnerTests/testCommand', - // A client in the field already pkills these exact bytes. This sweep ships separately from the - // runner and cannot know which version wrote a timed-out launch, so it never follows a rename: - // it must keep matching the names older writers used. The literal stays pinned here rather than - // derived from `runner-artifact-env.ts`, which builds only today's session name. - String.raw`xcodebuild .*AgentDeviceRunner\.env\.session-`, - String.raw`xcodebuild build-for-testing .*apple/runner/AgentDeviceRunner/AgentDeviceRunner\.xcodeproj`, -]; - -// `--platform` selectors that AFFIRMATIVELY name (or alias) an Apple device. -// This is deliberately narrower than "not proven non-Apple": the client's -// declared platform is not authoritative for session-bound execution (see -// the eligibility note on `handleRequestTimeout` below), so an undeclared or -// declared-non-Apple platform is not evidence of anything — it only counts -// as Apple evidence when it says so outright. -const AFFIRMATIVE_APPLE_PLATFORM_SELECTORS: ReadonlySet = new Set([ - 'apple', - 'ios', - 'macos', -]); - -function isAffirmativelyApplePlatform(platform: PlatformSelector | undefined): boolean { - return platform !== undefined && AFFIRMATIVE_APPLE_PLATFORM_SELECTORS.has(platform); -} +// Recovery for a timed-out request is scoped to that request or to nothing. The daemon owns +// request-scoped recovery (#3177): destroying the timed-out connection cancels exactly that +// request, and the cancel retires the runner work the request owned. The client used to recover +// host-scoped on both ends — a `pkill -f` sweep matching every agent-device runner xcodebuild, +// which sabotaged that cancel path (a swept kill reads as a host failure, not a canceled request) +// and duplicated cleanup the daemon scopes to its own leases — plus a daemon SIGKILL that ended +// every sibling session. The SIGKILL survives only behind the liveness probe below: the question +// the reset used to assume. That probe and the registration fence load on demand: a timed-out +// request is their only caller, and the CLI's eager closure (the ADR 0019 loading-shape probe) +// must not pay for a recovery path a healthy invocation never runs. export async function handleRequestTimeout( params: Readonly<{ @@ -51,24 +34,25 @@ export async function handleRequestTimeout( ): Promise { const { info, statePaths, remote, timeoutMs, requestId, command, platform, session, action } = params; - // Cleanup eligibility never depends on the declared platform, on purpose: - // the request's declared --platform is not - // authoritative for session-bound execution. An existing session's real - // device platform can silently override a conflicting declared selector - // (`applyStripLockPolicy` in request-lock-policy.ts, reached via - // --session-lock strip), and the common session-bound request omits - // --platform entirely — so there is no client-visible signal that proves - // a request cannot touch an Apple runner. The pkill patterns are - // Apple-process-name-specific, so sweeping them on a non-Apple host or - // session matches nothing and costs a few no-op subprocess spawns, never - // a wrong skip. - // `record` is excluded by command, which is authoritative: on a physical iOS device or macOS the - // runner is the recorder, and the sweep would kill the export the preserved daemon is finishing. - const sweepRunnerBuilds = !remote && command !== PUBLIC_COMMANDS.record; - const cleanup = sweepRunnerBuilds ? cleanupTimedOutIosRunnerBuilds() : { terminated: 0 }; - const resetDaemon = !remote && shouldResetDaemonAfterRequestTimeout(command); + // The command's declared policy (`timeoutPolicy.onTimeout`, ADR 0008) decides whether a local + // timed-out request is RESET-ELIGIBLE; the liveness probe decides whether the eligibility runs. + // A remote client cannot reach the host's process table or the daemon's state dir, so a remote + // timeout stays purely declarative here. + const resetEligible = !remote && shouldResetDaemonAfterRequestTimeout(command); + // A daemon that answers the probe is busy or slow, not gone: the request-scoped cancel its + // transport performed when this client's connection died is the recovery that request needed, + // and sibling sessions survive. Only a daemon answering neither endpoint in the probe window is + // the hung daemon a reset was designed for. + const probeAnswered = resetEligible + ? await ( + await import('./daemon-client-liveness-probe.ts') + ).probeDaemonResponsive(info, { + session, + }) + : undefined; + const unresponsive = probeAnswered === false; let retirement: DaemonRetirementResult | undefined; - if (resetDaemon) { + if (unresponsive) { const { stopAndRetireDaemon } = await import('../daemon-registration-owner.ts'); retirement = await stopAndRetireDaemon({ paths: statePaths, @@ -76,16 +60,9 @@ export async function handleRequestTimeout( mode: 'force', }); } + const resetDaemon = retirement?.status === 'retired'; const preserved = retirement?.status === 'retained' && retirement.termination?.status !== 'exited'; - // 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 - // process (proof positive regardless of what --platform claimed). Any - // other combination — undeclared platform, declared non-Apple platform, - // zero processes terminated — gets platform-neutral wording instead of - // asserting Apple specifics the client cannot back up. - const appleCleanupEvidence = isAffirmativelyApplePlatform(platform) || cleanup.terminated > 0; emitDiagnostic({ level: 'error', phase: 'daemon_request_timeout', @@ -93,12 +70,11 @@ export async function handleRequestTimeout( timeoutMs, requestId, command, - timedOutRunnerPidsTerminated: cleanup.terminated, - timedOutRunnerCleanupError: cleanup.error, - daemonPidReset: retirement?.status === 'retired' ? info.pid : undefined, + daemonPidReset: resetDaemon ? info.pid : undefined, daemonPidForceKilled: resetDaemon ? daemonWasForceKilled(retirement) : undefined, daemonRetirement: retirement, - daemonPreservedAfterTimeout: preserved || (!remote && !resetDaemon), + daemonPreservedAfterTimeout: preserved || (!remote && !unresponsive), + daemonLivenessProbeAnswered: probeAnswered, daemonBaseUrl: info.baseUrl, }, }); @@ -112,50 +88,62 @@ export async function handleRequestTimeout( ? `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, + resetDaemon: resetDaemon, command, - appleCleanupEvidence, + applePlatformDeclared: isAffirmativelyApplePlatform(platform), session, action, }), }); } -function daemonWasForceKilled(retirement: DaemonRetirementResult | undefined): boolean { - if (!retirement || retirement.status === 'absent') return false; - const termination = retirement.termination; - return termination?.status === 'exited' && termination.mode === 'forced'; -} - -// Whether a timed-out request tears down the local daemon is declared on the -// command's descriptor (ADR 0008, `timeoutPolicy.onTimeout`): read-only -// capture/polling commands preserve the daemon so sessions survive and evidence -// commands still work; everything else resets it. Unknown/undefined commands -// fall back to the default reset-daemon policy. +// Whether a timed-out request is eligible to tear down the local daemon is declared on the +// command's descriptor (ADR 0008, `timeoutPolicy.onTimeout`): read-only capture/polling commands +// preserve the daemon so sessions survive and evidence commands still work; everything else is +// eligible. Execution of the eligibility is gated by the liveness probe — the descriptor says a +// reset is allowed where the daemon is unreachable, never that a slow-but-alive daemon may lose +// every session it owns. Unknown/undefined commands fall back to the default reset-eligible +// policy, which matches the old hand lists: not listed meant default envelope + reset. function shouldResetDaemonAfterRequestTimeout(command: string | undefined): boolean { return resolveCommandTimeoutPolicy(command).onTimeout === 'reset-daemon'; } -// Exported for direct hint-matrix testing: handleRequestTimeout also triggers -// real pkill/process-kill side effects, so its wording is verified through -// this pure sub-function rather than the full timeout path (see also the -// production-seam route tests in -// src/daemon-client/__tests__/daemon-client-timeout-route.test.ts, which -// prove the cleanup-eligibility side of this contract that a pure formatter -// test cannot). +// `--platform` selectors that AFFIRMATIVELY name (or alias) an Apple device. This is the +// hint-wording gate: the hint names Apple-runner involvement only where this call site has +// evidence for it — an explicitly declared Apple platform selector. The sweep that used to offer +// "terminated something" as evidence is gone (#3177), and an undeclared or declared non-Apple +// platform is not evidence of anything. +const AFFIRMATIVE_APPLE_PLATFORM_SELECTORS: ReadonlySet = new Set([ + 'apple', + 'ios', + 'macos', +]); + +function isAffirmativelyApplePlatform(platform: PlatformSelector | undefined): boolean { + return platform !== undefined && AFFIRMATIVE_APPLE_PLATFORM_SELECTORS.has(platform); +} + +// Exported for direct hint-matrix testing: handleRequestTimeout also runs the real liveness probe +// and process-kill side effects, so its wording is verified through this pure sub-function +// (see also the production-seam route tests in +// src/daemon-client/__tests__/daemon-client-timeout-route.test.ts, which prove the liveness-gated +// recovery this route performs — the side of the contract a pure formatter test cannot reach). export function resolveRequestTimeoutHint(params: { remote: boolean; resetDaemon: boolean; command: string | undefined; - appleCleanupEvidence: boolean; + /** Only ever consulted on the kept-alive branch, where an Apple hint has evidence to stand on. */ + applePlatformDeclared: boolean; /** The request's first positional, for commands whose recovery depends on which action ran. */ action?: string; session?: string; }): string { - const { remote, resetDaemon, command, appleCleanupEvidence, session, action } = params; + const { remote, resetDaemon, command, applePlatformDeclared, session, action } = params; // A daemon that survives this client window may still be exporting a `record stop` that ran out - // of time (a stop still queued for the device lock is dropped before any export starts), and a - // finished file stays retrievable by asking again. A reset daemon makes no such promise. + // of time — a stop still queued for the device lock is dropped before any export starts, hence + // "may" — and a finished file stays retrievable by asking again. This holds whether the daemon + // survived because the request was remote or because `record` declares preserve-daemon (#3199). + // A reset daemon makes no such promise. if (!resetDaemon && command === PUBLIC_COMMANDS.record && action === 'stop') { return `The ${remote ? 'remote ' : ''}daemon may still be exporting the recording. Run agent-device record stop${ session ? ` --session ${session}` : '' @@ -164,33 +152,25 @@ export function resolveRequestTimeoutHint(params: { if (remote) { return 'Retry with --debug and verify the remote daemon URL, auth token, and remote host logs.'; } - if (!resetDaemon) { - const iosPrepareHint = - appleCleanupEvidence && command === PUBLIC_COMMANDS.snapshot - ? ' If this was the first Apple-platform snapshot on the device, run agent-device prepare ios-runner with the same --platform before snapshot/test so runner startup is handled explicitly.' - : ''; - const appleCleanupNote = appleCleanupEvidence - ? ' and Apple runner work was aborted when detected' - : ''; - return `Retry with --debug and check daemon diagnostics logs. The timed-out ${command ?? 'request'} request was canceled${appleCleanupNote}; the daemon was kept alive so the session can still be closed or inspected.${iosPrepareHint}`; + if (resetDaemon) { + // The reset SIGKILLs the daemon pid only; child processes it owned (Apple runners included) + // are not proven stopped from here, so the hint claims nothing about them. + return 'Retry with --debug and check daemon diagnostics logs. The daemon did not answer the liveness probe and was reset after the timeout.'; } - return appleCleanupEvidence - ? 'Retry with --debug and check daemon diagnostics logs. Timed-out Apple runner xcodebuild processes were terminated when detected.' - : 'Retry with --debug and check daemon diagnostics logs. The daemon was reset after the timeout.'; + const iosPrepareHint = + applePlatformDeclared && command === PUBLIC_COMMANDS.snapshot + ? ' If this was the first Apple-platform snapshot on the device, run agent-device prepare ios-runner with the same --platform before snapshot/test so runner startup is handled explicitly.' + : ''; + // The daemon canceled the request when this client's connection was destroyed, and stayed + // reachable (the probe answered, or the command's policy forbids the reset), so the session + // survives and can still be closed or inspected. + return `Retry with --debug and check daemon diagnostics logs. The timed-out ${ + command ?? 'request' + } request was canceled; the daemon was kept alive so the session can still be closed or inspected.${iosPrepareHint}`; } -function cleanupTimedOutIosRunnerBuilds(): { terminated: number; error?: string } { - let terminated = 0; - try { - for (const pattern of IOS_RUNNER_XCODEBUILD_KILL_PATTERNS) { - const result = runCmdSync('pkill', ['-f', pattern], { allowFailure: true }); - if (result.exitCode === 0) terminated += 1; - } - return { terminated }; - } catch (error) { - return { - terminated, - error: error instanceof Error ? error.message : String(error), - }; - } +function daemonWasForceKilled(retirement: DaemonRetirementResult | undefined): boolean { + if (!retirement || retirement.status === 'absent') return false; + const termination = retirement.termination; + return termination?.status === 'exited' && termination.mode === 'forced'; }