diff --git a/src/daemon-client/__tests__/daemon-client-transport.test.ts b/src/daemon-client/__tests__/daemon-client-transport.test.ts index 83eb72b4bd..9dde930eae 100644 --- a/src/daemon-client/__tests__/daemon-client-transport.test.ts +++ b/src/daemon-client/__tests__/daemon-client-transport.test.ts @@ -1,6 +1,6 @@ import assert from 'node:assert/strict'; import http from 'node:http'; -import { test } from 'vitest'; +import { test, vi } from 'vitest'; import { AppError } from '@agent-device/kernel/errors'; import { DAEMON_HTTP_INSTANCE_HEADER, @@ -222,6 +222,66 @@ test('a delayed restart health probe stops at the RPC deadline without retrying' } }); +test('a restart health probe that times out just before the RPC deadline still reports the deadline', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + let probing = false; + const realNow = performance.now.bind(performance); + const server = http.createServer((req, res) => { + if (req.url === '/health') { + probing = true; + return; + } + res.statusCode = 409; + res.setHeader(DAEMON_HTTP_INSTANCE_MISMATCH_HEADER, 'true'); + res.end(); + }); + const clock = vi + .spyOn(performance, 'now') + .mockImplementation(() => realNow() - (probing ? 50 : 0)); + try { + const port = await listenOnLoopback(server); + await assert.rejects( + sendWithStaleInstance(port, 150), + (error: unknown) => + error instanceof AppError && error.details?.reason === 'daemon_transport_timeout', + ); + assert.equal(probing, true); + } finally { + clock.mockRestore(); + await closeLoopbackServer(server); + } +}); + +test('a restart health probe refused near the RPC deadline reports the daemon as unavailable', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + let probed = false; + const server = http.createServer((req, res) => { + if (req.url === '/health') { + probed = true; + res.statusCode = 503; + res.end(); + return; + } + res.statusCode = 409; + res.setHeader(DAEMON_HTTP_INSTANCE_MISMATCH_HEADER, 'true'); + res.end(); + }); + try { + const port = await listenOnLoopback(server); + await assert.rejects( + sendWithStaleInstance(port, 150), + (error: unknown) => + error instanceof AppError && + error.message === 'Remote daemon is unavailable' && + error.details?.reason !== 'daemon_transport_timeout' && + error.details?.daemonBaseUrl === `http://127.0.0.1:${port}`, + ); + assert.equal(probed, true); + } finally { + await closeLoopbackServer(server); + } +}); + test('proxy forwards cached upstream identity and rejects a restarted upstream before dispatch', async (t) => { if (await skipWhenLoopbackUnavailable(t)) return; let upstreamInstance = 'upstream-one'; diff --git a/src/daemon-client/daemon-client-transport.ts b/src/daemon-client/daemon-client-transport.ts index d7ac9d2f43..3877df90d3 100644 --- a/src/daemon-client/daemon-client-transport.ts +++ b/src/daemon-client/daemon-client-transport.ts @@ -1,4 +1,5 @@ import type { RequestProgressSink } from '@agent-device/contracts/progress'; +import type http from 'node:http'; import net from 'node:net'; import { AppError } from '@agent-device/kernel/errors'; import { loadNodeHttpRequester, readNodeHttpResponseBody } from '@agent-device/host-kit/transport'; @@ -94,11 +95,20 @@ function canConnectHttp(info: DaemonInfo): Promise { return readDaemonHttpHealth(info).then((health) => health.reachable); } +type HealthProbeResult = { + health: RemoteDaemonHealth; + /** The probe's own budget timer ended it, not a refusal, reset, abort or 5xx. */ + timedOut: boolean; +}; + export async function readRemoteDaemonHealth( info: DaemonInfo, probeTimeoutMs?: number, + onProbeTimedOut?: () => void, ): Promise { - const health = await readDaemonHttpHealth(info, probeTimeoutMs); + const probe = await probeDaemonHttpHealth(info, probeTimeoutMs); + const health = probe.health; + if (probe.timedOut) onProbeTimedOut?.(); if (!info.baseUrl || !health.reachable) return health; // Every link a command RPC crosses has to speak the client's protocol: a proxy that reports a // skewed daemon behind it fails here, before the RPC, exactly like a skewed proxy does. @@ -125,21 +135,30 @@ async function readDaemonHttpHealth( info: DaemonInfo, probeTimeoutMs?: number, ): Promise { + return (await probeDaemonHttpHealth(info, probeTimeoutMs)).health; +} + +async function probeDaemonHttpHealth( + info: DaemonInfo, + probeTimeoutMs?: number, +): Promise { + const unreachable: HealthProbeResult = { health: { reachable: false }, timedOut: false }; const endpoint = info.baseUrl ? buildDaemonHttpUrl(info.baseUrl, 'health') : info.httpPort ? `http://127.0.0.1:${info.httpPort}/health` : null; - if (!endpoint) return { reachable: false }; + if (!endpoint) return unreachable; const url = new URL(endpoint); const transport = await loadNodeHttpRequester(url.protocol); const timeoutMs = Math.min( info.baseUrl ? REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS : LOCAL_DAEMON_HEALTHCHECK_TIMEOUT_MS, probeTimeoutMs ?? Number.POSITIVE_INFINITY, ); - if (timeoutMs <= 0) return { reachable: false }; + if (timeoutMs <= 0) return unreachable; return await new Promise((resolve) => { const headers = info.baseUrl ? buildDaemonHttpAuthHeaders(info.token) : {}; + const signal = AbortSignal.timeout(Math.ceil(timeoutMs)); const req = transport.request( { protocol: url.protocol, @@ -148,38 +167,44 @@ async function readDaemonHttpHealth( path: url.pathname + url.search, method: 'GET', timeout: timeoutMs, - signal: AbortSignal.timeout(Math.ceil(timeoutMs)), + signal, headers, }, - (res) => { - let body = ''; - res.setEncoding('utf8'); - res.on('data', (chunk) => { - body += chunk; - }); - res.on('end', () => { - const statusCode = res.statusCode ?? 500; - resolve({ - reachable: statusCode < 500, - statusCode, - ...readHealthPayload(body), - }); - }); - res.on('error', () => resolve({ reachable: false })); - res.on('aborted', () => resolve({ reachable: false })); - }, + (res) => collectHealthResponse(res, signal, resolve), ); req.on('timeout', () => { req.destroy(); - resolve({ reachable: false }); + resolve({ ...unreachable, timedOut: true }); }); req.on('error', () => { - resolve({ reachable: false }); + resolve({ ...unreachable, timedOut: signal.aborted }); }); req.end(); }); } +function collectHealthResponse( + res: http.IncomingMessage, + signal: AbortSignal, + resolve: (result: HealthProbeResult) => void, +): void { + const unreachable: HealthProbeResult = { health: { reachable: false }, timedOut: false }; + let body = ''; + res.setEncoding('utf8'); + res.on('data', (chunk) => { + body += chunk; + }); + res.on('end', () => { + const statusCode = res.statusCode ?? 500; + resolve({ + health: { reachable: statusCode < 500, statusCode, ...readHealthPayload(body) }, + timedOut: false, + }); + }); + res.on('error', () => resolve({ ...unreachable, timedOut: signal.aborted })); + res.on('aborted', () => resolve({ ...unreachable, timedOut: signal.aborted })); +} + function readHealthPayload(body: string): Omit { try { const parsed = JSON.parse(body) as { upstream?: unknown }; @@ -254,7 +279,13 @@ async function retryAfterRemoteInstanceMismatch( timeoutMs, deadline, ); - const health = await readRemoteDaemonHealth(info, probeTimeoutMs); + let probeTimedOut = false; + const health = await readRemoteDaemonHealth(info, probeTimeoutMs, () => { + probeTimedOut = true; + }); + if (probeTimedOut && timeoutMs !== undefined && isCallerDeadlineProbeBudget(probeTimeoutMs)) { + throw requestTimeoutError(info, req, statePaths, timeoutMs); + } const remainingMs = remainingRemoteRequestTimeoutMs(info, req, statePaths, timeoutMs, deadline); if (!health.reachable) { throw new AppError('COMMAND_FAILED', 'Remote daemon is unavailable', { @@ -272,6 +303,10 @@ async function retryAfterRemoteInstanceMismatch( } } +function isCallerDeadlineProbeBudget(probeTimeoutMs: number | undefined): boolean { + return probeTimeoutMs !== undefined && probeTimeoutMs <= REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS; +} + function remainingRemoteRequestTimeoutMs( info: DaemonInfo, req: DaemonRequest, @@ -282,7 +317,16 @@ function remainingRemoteRequestTimeoutMs( if (deadline === undefined || timeoutMs === undefined) return undefined; const remainingMs = deadline - performance.now(); if (remainingMs > 0) return remainingMs; - throw handleRequestTimeout({ + throw requestTimeoutError(info, req, statePaths, timeoutMs); +} + +function requestTimeoutError( + info: DaemonInfo, + req: DaemonRequest, + statePaths: DaemonPaths, + timeoutMs: number, +): AppError { + return handleRequestTimeout({ info, statePaths, ...timeoutRequestContext(req, true, timeoutMs), diff --git a/test/wire-compat/ledger.json b/test/wire-compat/ledger.json index 351e440b00..eb0be536c5 100644 --- a/test/wire-compat/ledger.json +++ b/test/wire-compat/ledger.json @@ -65,10 +65,10 @@ "src/daemon-client/daemon-client-rpc.ts#toDaemonHttpRpcError": "sha256:888246763c48670e7da893054d025744654f8715c3b4906312617a2b5028316b", "src/daemon-client/daemon-client-transport.ts#RemoteDaemonHealth": "sha256:df6b92e343a451a03b194e1e4aacac0c89b48a632723af24e087b377e73e05e3", "src/daemon-client/daemon-client-transport.ts#RemoteDaemonHealthLink": "sha256:464809c9bb14b44d098781722e9e9e2ff044c45464e3c7c4cb5690886e857fdf", - "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth": "sha256:6b64ae9b8e461a4b7dc1afb88465fcbaef908c31cc0be24499b2f440616e1cb5", + "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth": "sha256:4df0625b53e82e0d08225139a24d2dc234a9e18ce108b03d933d303d518e3731", "src/daemon-client/daemon-client-transport.ts#readHealthLink": "sha256:d7a84687d1c9b089460151a0532870e92a932d45f9567d248cea7accf8908593", "src/daemon-client/daemon-client-transport.ts#readHealthPayload": "sha256:4e85ffc3e35e02379c393e9312344757e003cf1f0ad9eb8d1f77d90c81c861f1", - "src/daemon-client/daemon-client-transport.ts#readRemoteDaemonHealth": "sha256:bcefa89fbb7fcbd6fee1b5ecb217955b9eee8d6fd2ad175fb653011edbd199d9", + "src/daemon-client/daemon-client-transport.ts#readRemoteDaemonHealth": "sha256:80bcf6b5987053e245b8e73cb21a687d158edc859a26605546746cdf0d64f443", "src/daemon/downloadable-artifact-http.ts#DownloadableArtifactHttpAuthorizer": "sha256:1b2702a929ca9170db2ca97c08e3ab67e17edb3ee75325a576c4c1b9cdbebb44", "src/daemon/downloadable-artifact-http.ts#DownloadableArtifactHttpRoute": "sha256:e63c4581ccde668913914149c9092d16ecf8a5cbd7ab33c6eb8617e77fc1015e", "src/daemon/downloadable-artifact-http.ts#handleArtifactDownload": "sha256:7f96d17b7c605230fb3cc21ceaa5b2e515d7445653ff95b2b0f6f15fa6214d4e", @@ -269,13 +269,13 @@ }, { "declaration": "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth", - "digest": "sha256:6b64ae9b8e461a4b7dc1afb88465fcbaef908c31cc0be24499b2f440616e1cb5", - "rationale": "#2650 gives every HTTP health probe a total-time cap (3s remote or 500ms local), including response-body reads. A restart retry also uses the RPC's remaining deadline. The health request and accepted payload fields are unchanged; stalled probes now end sooner." + "digest": "sha256:4df0625b53e82e0d08225139a24d2dc234a9e18ce108b03d933d303d518e3731", + "rationale": "#2650 gives every HTTP health probe a total-time cap (3s remote or 500ms local), including response-body reads. A restart retry also uses the RPC's remaining deadline. The health request and accepted payload fields are unchanged; stalled probes now end sooner. The health request and accepted payload do not change; only the internal probe outcome now records whether its own timer fired." }, { "declaration": "src/daemon-client/daemon-client-transport.ts#readRemoteDaemonHealth", - "digest": "sha256:bcefa89fbb7fcbd6fee1b5ecb217955b9eee8d6fd2ad175fb653011edbd199d9", - "rationale": "#2198 checks every link a command RPC crosses for protocol skew. #2650 caps a retry health probe at the request's remaining deadline. Health payload parsing and protocol mismatch refusal are unchanged; released peers still receive the same RPCs." + "digest": "sha256:80bcf6b5987053e245b8e73cb21a687d158edc859a26605546746cdf0d64f443", + "rationale": "#2198 checks every link a command RPC crosses for protocol skew. #2650 caps a retry health probe at the request's remaining deadline. Health payload parsing and protocol mismatch refusal are unchanged; released peers still receive the same RPCs. The health request and accepted payload do not change; only the internal probe outcome now records whether its own timer fired." }, { "declaration": "src/daemon/server/http-server.ts#authorizeAuxiliaryHttpRequest",