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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
62 changes: 61 additions & 1 deletion src/daemon-client/__tests__/daemon-client-transport.test.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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';
Expand Down
94 changes: 69 additions & 25 deletions src/daemon-client/daemon-client-transport.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -94,11 +95,20 @@ function canConnectHttp(info: DaemonInfo): Promise<boolean> {
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<RemoteDaemonHealth> {
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.
Expand All @@ -125,21 +135,30 @@ async function readDaemonHttpHealth(
info: DaemonInfo,
probeTimeoutMs?: number,
): Promise<RemoteDaemonHealth> {
return (await probeDaemonHttpHealth(info, probeTimeoutMs)).health;
}

async function probeDaemonHttpHealth(

@cubic-dev-ai cubic-dev-ai Bot Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This refactor moves the wire-relevant health request tokens out of the ledger-covered readDaemonHttpHealth declaration. Previously buildDaemonHttpUrl, transport.request({method: 'GET', ...}), the statusCode < 500 threshold, and body-read/abort handling all lived inside readDaemonHttpHealth's digested body; now they live in probeDaemonHttpHealth/collectHealthResponse, neither of which is listed in test/wire-compat/surface.ts's /health consumer group. The re-pinned digest for readDaemonHttpHealth now covers only a delegation, so a future wire-breaking change to the health request shape inside these helpers will no longer move the ledger digest — coverage the /health boundary previously had is silently lost. Add the new helpers to the surface (with a wire-mutations proof if they introduce a new break class) or note the moved tokens as reviewer-owned uncovered, per the wire-compat README's 'never claim coverage you don't provide'.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At src/daemon-client/daemon-client-transport.ts, line 141:

<comment>This refactor moves the wire-relevant health request tokens out of the ledger-covered `readDaemonHttpHealth` declaration. Previously `buildDaemonHttpUrl`, `transport.request({method: 'GET', ...})`, the `statusCode < 500` threshold, and body-read/abort handling all lived inside `readDaemonHttpHealth`'s digested body; now they live in `probeDaemonHttpHealth`/`collectHealthResponse`, neither of which is listed in `test/wire-compat/surface.ts`'s /health consumer group. The re-pinned digest for `readDaemonHttpHealth` now covers only a delegation, so a future wire-breaking change to the health request shape inside these helpers will no longer move the ledger digest — coverage the /health boundary previously had is silently lost. Add the new helpers to the surface (with a `wire-mutations` proof if they introduce a new break class) or note the moved tokens as reviewer-owned `uncovered`, per the wire-compat README's 'never claim coverage you don't provide'.</comment>

<file context>
@@ -125,21 +135,30 @@ async function readDaemonHttpHealth(
+  return (await probeDaemonHttpHealth(info, probeTimeoutMs)).health;
+}
+
+async function probeDaemonHttpHealth(
+  info: DaemonInfo,
+  probeTimeoutMs?: number,
</file context>
Fix with cubic

info: DaemonInfo,
probeTimeoutMs?: number,
): Promise<HealthProbeResult> {
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,
Expand All @@ -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<RemoteDaemonHealth, 'reachable' | 'statusCode'> {
try {
const parsed = JSON.parse(body) as { upstream?: unknown };
Expand Down Expand Up @@ -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', {
Expand All @@ -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,
Expand All @@ -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),
Expand Down
12 changes: 6 additions & 6 deletions test/wire-compat/ledger.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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",
Expand Down
Loading