From df5176710c9f21c669a4ca741dd84ae7d0ebbdef Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sat, 3 Oct 2026 23:22:38 +0200 Subject: [PATCH 01/14] feat(contracts,host-kit): declare the per-call AbortSignal and its requester-side guard MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Declare `signal?: AbortSignal` once at the owning contract types — the per-call overrides, the transport context, and the internal request envelope — so no per-site guard repeats it (#3178). Add `createRequestGuard` to @agent-device/host-kit/request beside the daemon-side cancellation machinery it answers for: an already-aborted call is refused with `dispatched: 'no'`, and an abort racing an in-flight send wins with the typed canceled-request error so a custom transport that ignores the signal cannot outlive the caller's promise. --- packages/contracts/src/client-connection.ts | 18 ++++- packages/contracts/src/request-envelope.ts | 5 ++ .../src/internal/request-guard.test.ts | 69 +++++++++++++++++++ .../host-kit/src/internal/request-guard.ts | 66 ++++++++++++++++++ packages/host-kit/src/request.ts | 1 + 5 files changed, 158 insertions(+), 1 deletion(-) create mode 100644 packages/host-kit/src/internal/request-guard.test.ts create mode 100644 packages/host-kit/src/internal/request-guard.ts diff --git a/packages/contracts/src/client-connection.ts b/packages/contracts/src/client-connection.ts index 7568768a33..72c62df725 100644 --- a/packages/contracts/src/client-connection.ts +++ b/packages/contracts/src/client-connection.ts @@ -17,6 +17,13 @@ import type { export type AgentDeviceDaemonTransportContext = { authToken?: string; + /** + * Cancels this one in-flight request. A built-in transport that sees an abort destroys the + * request's connection, which makes the daemon mark the request canceled; the promise rejects + * with the typed canceled-request error. A custom transport that ignores the signal keeps its + * own cancellation contract, and the client still rejects the caller's promise on abort. + */ + signal?: AbortSignal; }; export type AgentDeviceDaemonTransport = ( @@ -88,7 +95,16 @@ export type AgentDeviceRequestOverrides = Pick< | 'iosXctestrunFile' | 'iosXctestDerivedDataPath' | 'iosXctestEnvDir' ->; +> & { + /** + * Cancels this one call. Already aborted: the call rejects without sending anything + * (`details.dispatched: 'no'`). Aborted in flight: the request's connection closes, the daemon + * marks the request canceled, and the promise rejects with the typed canceled-request error + * (`details.reason: 'request_canceled'`). An abort is never a timeout: no runner sweep, no + * daemon reset. + */ + signal?: AbortSignal; +}; export type AgentDeviceIdentifiers = { session?: string; diff --git a/packages/contracts/src/request-envelope.ts b/packages/contracts/src/request-envelope.ts index fc222434e3..df4ce416ea 100644 --- a/packages/contracts/src/request-envelope.ts +++ b/packages/contracts/src/request-envelope.ts @@ -98,4 +98,9 @@ export type InternalRequestOptions = AgentDeviceClientConfig & leaseTtlMs?: number; provider?: string; providerSessionId?: string; + /** + * Cancels this one call in flight; never crosses the wire. The client hands it to the transport + * context and rejects its own promise on abort even when a custom transport ignores it. + */ + signal?: AbortSignal; }; diff --git a/packages/host-kit/src/internal/request-guard.test.ts b/packages/host-kit/src/internal/request-guard.test.ts new file mode 100644 index 0000000000..52a6622139 --- /dev/null +++ b/packages/host-kit/src/internal/request-guard.test.ts @@ -0,0 +1,69 @@ +import { test } from 'vitest'; +import assert from 'node:assert/strict'; +import { isRequestCanceledError } from '@agent-device/kernel/errors'; +import { createRequestGuard } from './request-guard.ts'; + +function canceledOutcome(promise: Promise): Promise { + return promise.then( + () => { + throw new Error('expected rejection'); + }, + (error: unknown) => error, + ); +} + +test('a guard without a signal never refuses and never interferes', async () => { + const guard = createRequestGuard({ signal: undefined, requestId: 'req-none' }); + guard.refuseIfAborted(); + assert.equal(await guard.guard(async () => 'sent'), 'sent'); +}); + +test('refuseIfAborted refuses an already-aborted call with dispatched no', async () => { + const controller = new AbortController(); + controller.abort(new Error('caller gave up')); + const guard = createRequestGuard({ signal: controller.signal, requestId: 'req-pre' }); + assert.throws( + () => guard.refuseIfAborted(), + (error: unknown) => + isRequestCanceledError(error) && + (error as { details?: Record }).details?.dispatched === 'no' && + (error as { details?: Record }).details?.requestId === 'req-pre', + ); +}); + +test('guard settles an in-flight send with the typed canceled error whatever the abort reason', async () => { + const controller = new AbortController(); + const guard = createRequestGuard({ signal: controller.signal, requestId: 'req-flight' }); + let settleSend: (() => void) | undefined; + const pending = guard.guard( + () => new Promise((resolve) => (settleSend = () => resolve('late'))), + ); + controller.abort(new Error('anything')); + const error = await canceledOutcome(pending); + assert.equal(isRequestCanceledError(error), true); + assert.equal((error as { details?: Record }).details?.dispatched, 'unknown'); + assert.equal((error as { details?: Record }).details?.requestId, 'req-flight'); + // A transport that ignores the signal may still resolve later; that must not + // turn the already-rejected caller promise into a second outcome. + settleSend?.(); +}); + +test('guard keeps a send error when no abort ever fires', async () => { + const controller = new AbortController(); + const guard = createRequestGuard({ signal: controller.signal }); + const failure = new Error('daemon refused'); + await assert.rejects( + guard.guard(async () => { + throw failure; + }), + (error: unknown) => error === failure, + ); +}); + +test('guard returns the send value when it settles before any abort', async () => { + const controller = new AbortController(); + const guard = createRequestGuard({ signal: controller.signal }); + const value = await guard.guard(async () => 'response'); + controller.abort(); + assert.equal(value, 'response'); +}); diff --git a/packages/host-kit/src/internal/request-guard.ts b/packages/host-kit/src/internal/request-guard.ts new file mode 100644 index 0000000000..c751aff151 --- /dev/null +++ b/packages/host-kit/src/internal/request-guard.ts @@ -0,0 +1,66 @@ +import { createRequestCanceledError, type AppError } from '@agent-device/kernel/errors'; + +/** + * The requester-side half of a per-call `AbortSignal`. The transport half closes the request's + * connection so the daemon marks the request canceled (its own cancellation machinery is + * {@link ./request-cancel.ts}); this half answers for the two cases a transport cannot: + * + * - a signal already aborted before anything was sent — nothing may leave the process, so the + * refusal carries `details.dispatched: 'no'`; + * - a custom transport that ignores the signal it was handed — the caller's promise must still + * settle when the abort fires, with `details.dispatched: 'unknown'` because the transport may + * already have written the request. + * + * Both outcomes reject with the typed canceled-request error (`details.reason: + * 'request_canceled'`) whatever reason the caller's own controller aborted with, so every layer + * that lets a cancellation through dispatches on the same reason — and never on the caller's + * arbitrary abort reason or on error text. An abort is never a timeout: nothing here reads or + * extends a request deadline, and no timeout path may produce this rejection. + */ + +export type RequestGuard = { + /** Refuses an already-aborted call before anything is sent. No-op without a signal. */ + refuseIfAborted(): void; + /** Settles `send`'s outcome against the signal, winning with the typed canceled error on abort. */ + guard(send: () => Promise): Promise; +}; + +const NO_REQUEST_GUARD: RequestGuard = { + refuseIfAborted: () => {}, + guard: async (send: () => Promise) => await send(), +}; + +export function createRequestGuard(params: { + signal: AbortSignal | undefined; + requestId?: string; +}): RequestGuard { + const { signal, requestId } = params; + if (!signal) return NO_REQUEST_GUARD; + const canceled = (dispatched: 'no' | 'unknown'): AppError => + createRequestCanceledError( + { requestId, dispatched }, + signal.reason instanceof Error ? signal.reason : undefined, + ); + return { + refuseIfAborted() { + if (signal.aborted) throw canceled('no'); + }, + guard: async (send: () => Promise): Promise => { + if (signal.aborted) throw canceled('unknown'); + return await new Promise((resolve, reject) => { + const onAbort = (): void => reject(canceled('unknown')); + signal.addEventListener('abort', onAbort, { once: true }); + void send().then( + (value) => { + signal.removeEventListener('abort', onAbort); + resolve(value); + }, + (error: unknown) => { + signal.removeEventListener('abort', onAbort); + reject(error); + }, + ); + }); + }, + }; +} diff --git a/packages/host-kit/src/request.ts b/packages/host-kit/src/request.ts index ee98a01e7b..07e5fc455c 100644 --- a/packages/host-kit/src/request.ts +++ b/packages/host-kit/src/request.ts @@ -9,3 +9,4 @@ export { throwIfRequestCanceled, } from './internal/request-cancel.ts'; export { emitRequestProgress, withRequestProgressSink } from './internal/request-progress.ts'; +export { createRequestGuard, type RequestGuard } from './internal/request-guard.ts'; From 88a0846d94c4cde83f3eae5e86d5acbc9631e719 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sat, 3 Oct 2026 23:27:25 +0200 Subject: [PATCH 02/14] feat(client): thread the per-call signal through the daemon client and both transports Wire the declared signal from the client's execute seam down to the transport context and the request guard, and from the built-in transport into both daemon wires (#3178): - `sendToDaemon` refuses an already-aborted call before daemon startup and guards each client-side phase; an abort during the artifact upload now combines with the lease beat's own signal so the upload stops for either owner. - The socket transport destroys its connection on abort, and the HTTP transport destroys its request; either way the daemon sees a client disconnect and marks the request canceled. Every settle path detaches the abort listener, and the timeout timer is cleared on abort so an abort never triggers the timeout path's runner sweep or daemon reset. - Retries and fallback attempts start behind a refusal check, so nothing is sent after an abort. --- packages/host-kit/src/request.ts | 7 +- src/agent-device-client.ts | 15 ++- src/daemon-client/daemon-client-transport.ts | 100 +++++++++++++++++-- src/daemon-client/daemon-client.ts | 49 ++++++--- 4 files changed, 147 insertions(+), 24 deletions(-) diff --git a/packages/host-kit/src/request.ts b/packages/host-kit/src/request.ts index 07e5fc455c..684ade7e65 100644 --- a/packages/host-kit/src/request.ts +++ b/packages/host-kit/src/request.ts @@ -9,4 +9,9 @@ export { throwIfRequestCanceled, } from './internal/request-cancel.ts'; export { emitRequestProgress, withRequestProgressSink } from './internal/request-progress.ts'; -export { createRequestGuard, type RequestGuard } from './internal/request-guard.ts'; +export { + abortedRequestError, + createRequestGuard, + refuseAbortedRequest, + type RequestGuard, +} from './internal/request-guard.ts'; diff --git a/src/agent-device-client.ts b/src/agent-device-client.ts index 104d169b48..f2848384af 100644 --- a/src/agent-device-client.ts +++ b/src/agent-device-client.ts @@ -74,6 +74,7 @@ import { isRecord, readSnapshotKeyboardBandFact } from '@agent-device/kernel/rec import { readResponseWarnings } from '@agent-device/kernel/success-text'; import { createLeaseClient } from './client/lease-client.ts'; import { normalizeScreenshotCaptureResult } from './client/screenshot-result.ts'; +import { createRequestGuard } from '@agent-device/host-kit/request'; export function createAgentDeviceClient( config: AgentDeviceClientConfig = {}, @@ -95,6 +96,11 @@ export function createAgentDeviceClient( input?: Record, ): Promise> => { const merged = mergeClientOptions(config, options); + const cancellation = createRequestGuard({ + signal: merged.signal, + requestId: merged.requestId, + }); + cancellation.refuseIfAborted(); const request = { session: resolveSessionName(merged.session), command, @@ -104,7 +110,14 @@ export function createAgentDeviceClient( runtime: merged.runtime, meta: buildMeta(merged), }; - const response = await transport(request, { authToken: merged.daemonAuthToken }); + // `signal` rides the transport context (it is a live object, never wire data), and the guard + // answers for a custom transport that ignores it: the caller's promise settles on abort either + // way. The built-in transport closes the request's connection, which is what makes the daemon + // mark the request canceled. + const response = await cancellation.guard( + async () => + await transport(request, { authToken: merged.daemonAuthToken, signal: merged.signal }), + ); if (!response.ok) { throwDaemonError(response.error); } diff --git a/src/daemon-client/daemon-client-transport.ts b/src/daemon-client/daemon-client-transport.ts index d09d5af386..5e7cd1a0b7 100644 --- a/src/daemon-client/daemon-client-transport.ts +++ b/src/daemon-client/daemon-client-transport.ts @@ -2,6 +2,7 @@ import type { RequestProgressSink } from '@agent-device/contracts/progress'; import net from 'node:net'; import { AppError } from '@agent-device/kernel/errors'; import { loadNodeHttpRequester, readNodeHttpResponseBody } from '@agent-device/host-kit/transport'; +import { abortedRequestError, refuseAbortedRequest } from '@agent-device/host-kit/request'; import type { DaemonRequest, DaemonResponse } from '../daemon/daemon-request.ts'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; import type { DaemonPaths, DaemonTransportPreference } from '../daemon-resolution.ts'; @@ -26,6 +27,13 @@ import { readVersion } from '@agent-device/host-kit/version'; type ResolvedDaemonTransport = 'socket' | 'http'; type SendRequestOptions = { onProgress?: RequestProgressSink; + /** + * The caller's per-request cancellation (#3178). Aborted at or before a send attempt, nothing + * leaves on that attempt; aborted in flight, the attempt's own connection is destroyed — which is + * what makes the daemon mark the request canceled — and the promise rejects with the typed + * canceled-request error. An abort is never a timeout: it never reaches `handleRequestTimeout`. + */ + signal?: AbortSignal; }; const LOCAL_DAEMON_HEALTHCHECK_TIMEOUT_MS = 500; @@ -240,6 +248,9 @@ export async function sendRequest( ): Promise { const transport = chooseTransport(info, preference); const deadline = typeof timeoutMs === 'number' ? performance.now() + timeoutMs : undefined; + // A canceled caller must not open a connection just to lose it, and a fallback or instance retry + // that begins after the abort must not send either — every attempt starts behind this check. + refuseAbortedRequest(options.signal, req.meta?.requestId); try { return await sendRequestWithTransport(info, req, statePaths, timeoutMs, transport, options); } catch (error) { @@ -452,17 +463,26 @@ async function sendSocketRequest( ): Promise { const port = info.port; if (!port) throw daemonEndpointUnavailableError('socket'); + const callerSignal = options.signal; + refuseAbortedRequest(callerSignal, req.meta?.requestId); return new Promise((resolve, reject) => { let requestWritten = false; + let settled = false; const socket = net.createConnection({ host: '127.0.0.1', port }, () => { + // An abort that landed while the connection was still opening must not write the request. + if (callerSignal?.aborted) { + settled = true; + rejectAborted(callerSignal); + return; + } requestWritten = true; socket.write(`${JSON.stringify(req)}\n`); }); - let settled = false; const timeoutHandle = typeof timeoutMs === 'number' ? setTimeout(() => { settled = true; + detachCallerAbort(); socket.destroy(); reject( handleRequestTimeout({ @@ -473,6 +493,27 @@ async function sendSocketRequest( ); }, timeoutMs) : undefined; + // Destroying the connection is what makes the daemon mark the request canceled; the rejection + // keeps the typed canceled-request error, never a socket-error or timeout shape. The timeout + // timer is cleared here because an abort is never a timeout: nothing may run the timeout's + // runner sweep or daemon reset after the caller canceled. + const rejectAborted = (signal: AbortSignal): void => { + detachCallerAbort(); + if (timeoutHandle) clearTimeout(timeoutHandle); + socket.destroy(); + reject(abortedRequestError(signal, requestWritten ? 'unknown' : 'no', req.meta?.requestId)); + }; + const onCallerAbort = (): void => { + if (settled) return; + settled = true; + rejectAborted(callerSignal!); + }; + const detachCallerAbort = (): void => { + if (callerSignal) callerSignal.removeEventListener('abort', onCallerAbort); + }; + if (callerSignal) { + callerSignal.addEventListener('abort', onCallerAbort, { once: true }); + } readDaemonSocketProgressResponse(socket, { req, @@ -483,10 +524,12 @@ async function sendSocketRequest( }, resolve: (response) => { settled = true; + detachCallerAbort(); resolve(response); }, reject: (error) => { settled = true; + detachCallerAbort(); reject(error); }, }); @@ -495,6 +538,7 @@ async function sendSocketRequest( if (settled) return; settled = true; if (timeoutHandle) clearTimeout(timeoutHandle); + detachCallerAbort(); reject( handleTransportError(err, req.meta?.requestId, false, { daemonSocketRequestWritten: requestWritten, @@ -563,8 +607,27 @@ async function sendHttpRequest( Object.assign(headers, buildRemoteInstancePreconditionHeaders(info)); } const transport = await loadNodeHttpRequester(rpcUrl.protocol); + const callerSignal = options.signal; + refuseAbortedRequest(callerSignal, req.meta?.requestId); return await new Promise((resolve, reject) => { + let settled = false; + let requestEnded = false; + const detachCallerAbort = (): void => { + if (callerSignal) callerSignal.removeEventListener('abort', onCallerAbort); + }; + const resolveOnce = (response: DaemonResponse | PromiseLike): void => { + if (settled) return; + settled = true; + detachCallerAbort(); + resolve(response); + }; + const rejectOnce = (error: unknown): void => { + if (settled) return; + settled = true; + detachCallerAbort(); + reject(error); + }; const request = transport.request( { protocol: rpcUrl.protocol, @@ -578,7 +641,7 @@ async function sendHttpRequest( if (isRemoteInstanceMismatchResponse(res.statusCode, res.headers ?? {})) { res.resume(); if (timeoutHandle) clearTimeout(timeoutHandle); - reject( + rejectOnce( new AppError('COMMAND_FAILED', 'Remote daemon instance changed', { reason: 'remote_instance_mismatch', daemonBaseUrl: info.baseUrl, @@ -590,7 +653,7 @@ async function sendHttpRequest( readDaemonHttpProgressResponse(res, { req, onProgress: options.onProgress, - reject, + reject: rejectOnce, clearTimeout: () => { if (timeoutHandle) clearTimeout(timeoutHandle); }, @@ -599,8 +662,8 @@ async function sendHttpRequest( info, req, stateDir: statePaths.baseDir, - resolve, - reject, + resolve: resolveOnce, + reject: rejectOnce, }); }, }); @@ -621,13 +684,13 @@ async function sendHttpRequest( info, req, stateDir: statePaths.baseDir, - resolve, - reject, + resolve: resolveOnce, + reject: rejectOnce, }); }) .catch((error: unknown) => { if (timeoutHandle) clearTimeout(timeoutHandle); - reject(error); + rejectOnce(error); }); }, ); @@ -637,7 +700,7 @@ async function sendHttpRequest( typeof timeoutMs === 'number' ? setTimeout(() => { request.destroy(); - reject( + rejectOnce( handleRequestTimeout({ info, statePaths, @@ -647,13 +710,30 @@ async function sendHttpRequest( }, timeoutMs) : undefined; + // Destroying the request closes this one connection, which is what makes the daemon mark the + // request canceled. The timeout timer is cleared because an abort is never a timeout: nothing + // may run the timeout's runner sweep or daemon reset after the caller canceled. The destroy + // surfaces as a request `error`, which `settled` keeps out of the already-typed rejection. + const onCallerAbort = (): void => { + if (settled) return; + if (timeoutHandle) clearTimeout(timeoutHandle); + request.destroy(); + rejectOnce( + abortedRequestError(callerSignal!, requestEnded ? 'unknown' : 'no', req.meta?.requestId), + ); + }; + if (callerSignal) { + callerSignal.addEventListener('abort', onCallerAbort, { once: true }); + } + request.on('error', (err) => { if (timeoutHandle) clearTimeout(timeoutHandle); - reject(handleTransportError(err, req.meta?.requestId, remote)); + rejectOnce(handleTransportError(err, req.meta?.requestId, remote)); }); request.write(rpcPayload); request.end(); + requestEnded = true; }); } diff --git a/src/daemon-client/daemon-client.ts b/src/daemon-client/daemon-client.ts index 7d777254fb..10aab698e2 100644 --- a/src/daemon-client/daemon-client.ts +++ b/src/daemon-client/daemon-client.ts @@ -10,6 +10,7 @@ import { emitDiagnostic, withDiagnosticTimer, } from '@agent-device/host-kit/diagnostics'; +import { createRequestGuard } from '@agent-device/host-kit/request'; import { INTERNAL_COMMANDS, PUBLIC_COMMANDS } from '@agent-device/command-registry/catalog'; import { resolveCommandTimeoutPolicy } from '@agent-device/command-registry/registry'; import { resolveCommandRequestTimeoutMs } from '@agent-device/command-registry/timeout-policy'; @@ -60,16 +61,28 @@ export async function sendToDaemon( resolveCommandTimeoutPolicy(requestWithoutAuthFlag.command), requestWithoutAuthFlag, ); - const daemon = await withDiagnosticTimer( - 'daemon_startup', - async () => await ensureDaemon(settings), - { requestId, session: req.session }, - ); + // The caller's signal covers every phase of this one request: an already-aborted call is + // refused before the daemon is even started or an artifact byte is uploaded, and an abort that + // arrives mid-phase stops the upload and the request. Each phase's own machinery (the upload + // client, the transports below) honors the signal; the guard turns any of those outcomes into + // the typed canceled-request error, so an abort never borrows a timeout's shape. + const cancellation = createRequestGuard({ signal: options.signal, requestId }); + cancellation.refuseIfAborted(); + const daemon = await cancellation.guard(async () => { + return await withDiagnosticTimer('daemon_startup', async () => await ensureDaemon(settings), { + requestId, + session: req.session, + }); + }); const info = daemon.info; - const preparedRemoteRequest = await protectArtifactUploadWithLeaseBeats( - info, - settings, - requestWithoutAuthFlag, + const preparedRemoteRequest = await cancellation.guard( + async () => + await protectArtifactUploadWithLeaseBeats( + info, + settings, + requestWithoutAuthFlag, + options.signal, + ), ); writeInstallInProgressNotice(requestWithoutAuthFlag.command); @@ -103,7 +116,7 @@ export async function sendToDaemon( settings.transportPreference, settings.paths, requestTimeoutMs, - { onProgress: options.onProgress }, + { onProgress: options.onProgress, signal: options.signal }, ), { requestId, command: req.command }, ); @@ -316,15 +329,27 @@ async function protectArtifactUploadWithLeaseBeats( info: DaemonInfo, settings: DaemonClientSettings, request: Omit, + callerSignal: AbortSignal | undefined, ): Promise { const leaseScope = leaseScopeFromRequest(request); if (!isRemoteDaemon(info) || !leaseScope.leaseId) { - return await prepareRemoteRequestArtifacts(request, info, new AbortController().signal); + return await prepareRemoteRequestArtifacts( + request, + info, + callerSignal ?? new AbortController().signal, + ); } const { buildUploadLeaseHeartbeat, runProtectedLeaseWork } = await import('./daemon-client-lease-beat.ts'); return await runProtectedLeaseWork({ heartbeat: buildUploadLeaseHeartbeat(info, settings, request), - task: (signal) => prepareRemoteRequestArtifacts(request, info, signal), + // The upload stops for either owner of its signal: the lease beat finding the lease gone, or + // the caller aborting this one call. + task: (leaseSignal) => + prepareRemoteRequestArtifacts( + request, + info, + callerSignal ? AbortSignal.any([leaseSignal, callerSignal]) : leaseSignal, + ), }); } From 215826fc825b0dfc47ee031b53323f897f1a0bd6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sat, 3 Oct 2026 23:33:50 +0200 Subject: [PATCH 03/14] refactor(host-kit): own the transport-facing canceled-request helpers in the guard module `abortedRequestError` and `refuseAbortedRequest` complete the request-cancel seam the daemon request transports consume, so every rejection a caller's abort produces is built by the module that declares the cancellation contract rather than restated per transport. --- .../host-kit/src/internal/request-guard.ts | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/packages/host-kit/src/internal/request-guard.ts b/packages/host-kit/src/internal/request-guard.ts index c751aff151..673403da9f 100644 --- a/packages/host-kit/src/internal/request-guard.ts +++ b/packages/host-kit/src/internal/request-guard.ts @@ -18,6 +18,33 @@ import { createRequestCanceledError, type AppError } from '@agent-device/kernel/ * extends a request deadline, and no timeout path may produce this rejection. */ +/** + * The typed canceled-request error for a caller's own abort: the reason the caller aborted with + * survives as the cause, while the rejection itself always dispatches on `reason: + * 'request_canceled'` with the delivery evidence the aborting layer can prove. A built-in transport + * rejects every abort through this, so a caller's arbitrary abort reason never escapes as the + * outcome of a daemon request. + */ +export function abortedRequestError( + signal: AbortSignal, + dispatched: 'no' | 'unknown', + requestId?: string, +): AppError { + return createRequestCanceledError( + { requestId, dispatched }, + signal.reason instanceof Error ? signal.reason : undefined, + ); +} + +/** + * Refuse a send attempt that starts while the caller's signal is already aborted: nothing may leave + * the process, so this never touches a connection and the refusal carries `details.dispatched: 'no'`. + */ +export function refuseAbortedRequest(signal: AbortSignal | undefined, requestId?: string): void { + if (!signal?.aborted) return; + throw abortedRequestError(signal, 'no', requestId); +} + export type RequestGuard = { /** Refuses an already-aborted call before anything is sent. No-op without a signal. */ refuseIfAborted(): void; From 9f9cac7eceff4d7c472c1cdbee8bd21926bb61d1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sat, 3 Oct 2026 23:34:35 +0200 Subject: [PATCH 04/14] test(client): cover the per-call AbortSignal on the client surface and both transports #3178 acceptance: an abort before send sends nothing (`dispatched: 'no'`); an abort mid-request closes that one request so the daemon marks it canceled (`dispatched: 'unknown'`, proven daemon-side through the request-cancel registry), and the daemon keeps serving follow-up requests. Client-level cases pin the guard half: the signal rides the transport context, and a custom transport that ignores it cannot outlive the caller's promise. Regression evidence: all eight wiring cases fail on the pre-fix sources (verified by running these files against fc7df6c4b's client/daemon-client/ transport) and pass with the fix. --- src/__tests__/client-abort-signal.test.ts | 91 +++++++ .../__tests__/daemon-client-abort.test.ts | 238 ++++++++++++++++++ 2 files changed, 329 insertions(+) create mode 100644 src/__tests__/client-abort-signal.test.ts create mode 100644 src/daemon-client/__tests__/daemon-client-abort.test.ts diff --git a/src/__tests__/client-abort-signal.test.ts b/src/__tests__/client-abort-signal.test.ts new file mode 100644 index 0000000000..7bdf69e245 --- /dev/null +++ b/src/__tests__/client-abort-signal.test.ts @@ -0,0 +1,91 @@ +/** + * #3178: every client call accepts `signal?: AbortSignal`. + * + * These pin the client's own half of the contract against a scripted transport: an already-aborted + * call never reaches the transport (`details.dispatched: 'no'`), the signal rides the transport + * context for a transport that honors it, and a custom transport that ignores the signal still + * loses the race — the caller's promise rejects with the typed canceled-request error + * (`details.dispatched: 'unknown'`). The transport-into-daemon half is covered by + * `src/daemon-client/__tests__/daemon-client-abort.test.ts`. + */ +import assert from 'node:assert/strict'; +import { test } from 'vitest'; +import { isRequestCanceledError } from '@agent-device/kernel/errors'; +import { createAgentDeviceClient } from '../agent-device-client.ts'; +import type { AgentDeviceDaemonTransportContext } from '@agent-device/contracts/client'; +import type { DaemonRequest, DaemonResponse } from '@agent-device/kernel/contracts'; +import { createTransport } from './client-transport-fixture.ts'; + +function canceledWith(error: unknown, dispatched: 'no' | 'unknown'): boolean { + return ( + isRequestCanceledError(error) && + (error as { details?: Record }).details?.dispatched === dispatched + ); +} + +test('a client call with an already-aborted signal never reaches the transport', async () => { + const { config, calls, transport } = createTransport(() => ({ ok: true, data: {} })); + const client = createAgentDeviceClient(config, { transport }); + const controller = new AbortController(); + controller.abort(); + await assert.rejects( + client.interactions.press({ ref: '@e12', signal: controller.signal }), + (error: unknown) => canceledWith(error, 'no'), + ); + assert.deepEqual(calls, []); +}); + +test('the signal rides the transport context so a built-in-style transport can close the request', async () => { + const contexts: Array = []; + const client = createAgentDeviceClient( + {}, + { + transport: async (_req, context) => { + contexts.push(context); + return { ok: true, data: {} } satisfies DaemonResponse; + }, + }, + ); + const controller = new AbortController(); + await client.command.wait({ durationMs: 1, signal: controller.signal }); + assert.equal(contexts.length, 1); + assert.equal(contexts[0]?.signal, controller.signal); +}); + +test('aborting a call whose transport ignores the signal still rejects the caller with the typed canceled error', async () => { + const client = createAgentDeviceClient( + {}, + { + transport: (req) => + new Promise((resolve) => { + // A custom transport that never inspects the context signal: the guard must still settle + // the caller's promise when the abort fires. + setTimeout(() => resolve({ ok: true, data: { ignored: req.command } }), 10_000); + }), + }, + ); + const controller = new AbortController(); + const call = client.interactions.press({ ref: '@e12', signal: controller.signal }); + setTimeout(() => controller.abort(), 10); + await assert.rejects(call, (error: unknown) => canceledWith(error, 'unknown')); +}); + +test('a signal on one call does not cancel another', async () => { + const client = createAgentDeviceClient( + {}, + { + transport: async (req: Omit) => { + if (req.command === 'wait') { + await new Promise((resolve) => setTimeout(resolve, 50)); + } + return { ok: true, data: {} }; + }, + }, + ); + const canceled = new AbortController(); + const doomed = client.command.wait({ durationMs: 5000, signal: canceled.signal }); + const survivor = client.command.wait({ durationMs: 1 }); + canceled.abort(); + await assert.rejects(doomed, (error: unknown) => isRequestCanceledError(error)); + assert.deepEqual(await survivor, {}); +}); diff --git a/src/daemon-client/__tests__/daemon-client-abort.test.ts b/src/daemon-client/__tests__/daemon-client-abort.test.ts new file mode 100644 index 0000000000..776ac35833 --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-abort.test.ts @@ -0,0 +1,238 @@ +/** + * #3178: the built-in daemon transports honor a per-call `AbortSignal`. + * + * These run the REAL daemon servers (`createSocketServer`, `createDaemonHttpServer`) against a + * request handler that waits on the daemon's own request-scoped signal, so each mid-flight case + * proves the whole chain: the client aborts → that one connection closes → the daemon marks the + * request canceled (`markRequestCanceled`, reached only through the transport's disconnect path) → + * the client rejects with the typed canceled-request error — and the daemon itself stays alive to + * serve a follow-up request. + */ +import assert from 'node:assert/strict'; +import { test } from 'vitest'; +import { isRequestCanceledError } from '@agent-device/kernel/errors'; +import { getRequestSignal, isRequestCanceled } from '@agent-device/host-kit/request'; +import type { DaemonInvokeFn, DaemonRequest, DaemonResponse } from '../../daemon/daemon-request.ts'; +import { createDaemonHttpServer } from '../../daemon/server/http-server.ts'; +import { createSocketServer, listenNetServer } from '../../daemon/server/transport.ts'; +import { sendRequest } from '../daemon-client-transport.ts'; +import { resolveDaemonPaths } from '../../daemon-resolution.ts'; +import { + closeLoopbackServer, + listenOnLoopback, + skipWhenLoopbackUnavailable, +} from '../../__tests__/test-utils/loopback.ts'; +import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; + +const TOKEN = 'abort-signal-token'; +const STATE_PATHS = resolveDaemonPaths(mkdtempForTestSync('agent-device-client-abort-')); + +type SeenRequest = { started: boolean; requestId?: string; canceled: boolean[] }; + +// `wait` models the long request the caller abandons; every other command answers immediately, so +// the same handler also proves the daemon still serves requests after one was canceled. +function canceledAwareHandler(seen: SeenRequest): DaemonInvokeFn { + return async (req: DaemonRequest): Promise => { + const requestId = req.meta?.requestId; + seen.requestId = requestId; + if (req.command !== 'wait') { + return { ok: true, data: {} }; + } + seen.started = true; + const signal = getRequestSignal(requestId); + if (!signal) { + return { ok: true, data: { answeredWithoutSignal: true } }; + } + return await new Promise((resolve, reject) => { + signal.addEventListener( + 'abort', + () => { + seen.canceled.push(isRequestCanceled(requestId)); + reject(signal.reason); + }, + { once: true }, + ); + }); + }; +} + +function canceledRequestError(error: unknown, dispatched: 'no' | 'unknown'): boolean { + return ( + isRequestCanceledError(error) && + (error as { details?: Record }).details?.dispatched === dispatched + ); +} + +test('socket transport: an already-aborted signal sends nothing and refuses with dispatched no', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + const connections: number[] = []; + const server = createSocketServer(async () => { + connections.push(1); + return { ok: true, data: {} }; + }); + try { + const port = await listenNetServer(server); + const controller = new AbortController(); + controller.abort(); + await assert.rejects( + sendRequest( + { port, token: TOKEN, pid: 1 }, + { + command: 'devices', + session: 'default', + positionals: [], + flags: {}, + meta: { requestId: 'req-socket-pre-abort' }, + }, + 'socket', + STATE_PATHS, + undefined, + { signal: controller.signal }, + ), + (error: unknown) => canceledRequestError(error, 'no'), + ); + assert.deepEqual(connections, []); + } finally { + await closeLoopbackServer(server); + } +}); + +test('socket transport: an abort mid-request closes the connection, the daemon marks the request canceled, and the daemon stays alive', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + const seen: SeenRequest = { started: false, canceled: [] }; + const server = createSocketServer(canceledAwareHandler(seen)); + try { + const port = await listenNetServer(server); + const info = { port, token: TOKEN, pid: 1 }; + const controller = new AbortController(); + const inFlight = sendRequest( + info, + { + command: 'wait', + session: 'default', + positionals: [], + flags: {}, + meta: { requestId: 'req-socket-abort-in-flight' }, + }, + 'socket', + STATE_PATHS, + undefined, + { signal: controller.signal }, + ); + // Abort only once the daemon has the request in hand, so this is genuinely a mid-request + // cancel rather than a race against connection setup. + await waitFor(() => seen.started, 'the daemon to start the request'); + controller.abort(); + await assert.rejects(inFlight, (error: unknown) => canceledRequestError(error, 'unknown')); + // The daemon's disconnect path ran: this request's id is registered-canceled, proven through + // the daemon-side registry rather than off the client's rejection. + await waitFor(() => seen.canceled.length > 0, 'daemon-side cancellation'); + assert.deepEqual(seen.canceled, [true]); + assert.equal(seen.requestId, 'req-socket-abort-in-flight'); + + // Daemon alive and serving: a follow-up request without any cancellation completes. + const followUp = await sendRequest( + info, + { + command: 'devices', + session: 'default', + positionals: [], + flags: {}, + meta: { requestId: 'req-socket-follow-up' }, + }, + 'socket', + STATE_PATHS, + 5000, + {}, + ); + assert.equal(followUp.ok, true); + } finally { + await closeLoopbackServer(server); + } +}); + +test('http transport: an already-aborted signal sends nothing and refuses with dispatched no', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + const received: number[] = []; + const server = await createDaemonHttpServer({ + handleRequest: async () => { + received.push(1); + return { ok: true, data: {} }; + }, + token: TOKEN, + }); + try { + const port = await listenOnLoopback(server); + const controller = new AbortController(); + controller.abort(); + await assert.rejects( + sendRequest( + { httpPort: port, token: TOKEN, pid: 1 }, + { + token: TOKEN, + command: 'devices', + session: 'default', + positionals: [], + flags: {}, + meta: { requestId: 'req-http-pre-abort' }, + }, + 'http', + STATE_PATHS, + undefined, + { signal: controller.signal }, + ), + (error: unknown) => canceledRequestError(error, 'no'), + ); + assert.deepEqual(received, []); + } finally { + await closeLoopbackServer(server); + } +}); + +test('http transport: an abort mid-request closes that request, the daemon marks it canceled, and the daemon stays alive', async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + const seen: SeenRequest = { started: false, canceled: [] }; + const server = await createDaemonHttpServer({ + handleRequest: canceledAwareHandler(seen), + token: TOKEN, + }); + try { + const port = await listenOnLoopback(server); + const info = { httpPort: port, token: TOKEN, pid: 1 }; + const controller = new AbortController(); + const inFlight = sendRequest( + info, + { + token: TOKEN, + command: 'wait', + session: 'default', + positionals: [], + flags: {}, + meta: { requestId: 'req-http-abort-in-flight' }, + }, + 'http', + STATE_PATHS, + undefined, + { signal: controller.signal }, + ); + await waitFor(() => seen.started, 'the daemon to start the request'); + controller.abort(); + await assert.rejects(inFlight, (error: unknown) => canceledRequestError(error, 'unknown')); + await waitFor(() => seen.canceled.length > 0, 'daemon-side cancellation'); + assert.deepEqual(seen.canceled, [true]); + assert.equal(seen.requestId, 'req-http-abort-in-flight'); + + const health = await fetch(`http://127.0.0.1:${port}/health`); + assert.equal(health.ok, true); + } finally { + await closeLoopbackServer(server); + } +}); + +async function waitFor(condition: () => boolean, what: string): Promise { + const deadline = Date.now() + 3000; + while (!condition()) { + if (Date.now() >= deadline) throw new Error(`Timed out waiting for ${what}.`); + await new Promise((resolve) => setTimeout(resolve, 10)); + } +} From c8ba9af11f71609f230b7cd4085833b1ee01c0c2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sat, 3 Oct 2026 23:35:40 +0200 Subject: [PATCH 05/14] docs(client): document the per-call signal on every client method #3178's Done-when: the public client docs list `signal`, including the already-aborted refusal (`dispatched: 'no'`), the mid-flight cancel (`reason: 'request_canceled'`), and the promise that an abort never takes the timeout path's runner cleanup or daemon reset. --- website/docs/docs/client-api.md | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/website/docs/docs/client-api.md b/website/docs/docs/client-api.md index d59051db73..e2702ad8f6 100644 --- a/website/docs/docs/client-api.md +++ b/website/docs/docs/client-api.md @@ -239,6 +239,15 @@ Results are daemon-shaped objects with typed known fields, so command semantics A failed interaction rejects with the same error the CLI prints. Read `error.details.dispatched` before you retry; [Commands](./commands.md) explains the two values. +Every client call accepts `signal?: AbortSignal` to cancel that one call: + +```ts +const controller = new AbortController(); +await client.interactions.press({ ref: '@e12', signal: controller.signal }); +``` + +A signal that is already aborted rejects the call without sending anything (`error.details.dispatched: 'no'`). Aborting while the request is in flight closes that request's connection, the daemon marks the request canceled, and the promise rejects with the typed canceled-request error (`error.details.reason: 'request_canceled'`, `error.details.dispatched: 'unknown'`). The daemon and the session stay alive for other requests. An abort is never a timeout: it never triggers the timeout path's runner cleanup or daemon reset. + ```ts await client.command.wait({ text: 'Continue', From e33c53d4537444183a64efc9c51d9983e210d483 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sat, 3 Oct 2026 23:38:05 +0200 Subject: [PATCH 06/14] test(daemon-client): carry the wire token in the abort fixtures' daemon requests --- src/daemon-client/__tests__/daemon-client-abort.test.ts | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/src/daemon-client/__tests__/daemon-client-abort.test.ts b/src/daemon-client/__tests__/daemon-client-abort.test.ts index 776ac35833..3c9dfbca8c 100644 --- a/src/daemon-client/__tests__/daemon-client-abort.test.ts +++ b/src/daemon-client/__tests__/daemon-client-abort.test.ts @@ -43,7 +43,8 @@ function canceledAwareHandler(seen: SeenRequest): DaemonInvokeFn { if (!signal) { return { ok: true, data: { answeredWithoutSignal: true } }; } - return await new Promise((resolve, reject) => { + // The long request never answers: only the client's disconnect ends it, through the abort. + return await new Promise((_resolve, reject) => { signal.addEventListener( 'abort', () => { @@ -78,6 +79,7 @@ test('socket transport: an already-aborted signal sends nothing and refuses with sendRequest( { port, token: TOKEN, pid: 1 }, { + token: TOKEN, command: 'devices', session: 'default', positionals: [], @@ -108,6 +110,7 @@ test('socket transport: an abort mid-request closes the connection, the daemon m const inFlight = sendRequest( info, { + token: TOKEN, command: 'wait', session: 'default', positionals: [], @@ -134,6 +137,7 @@ test('socket transport: an abort mid-request closes the connection, the daemon m const followUp = await sendRequest( info, { + token: TOKEN, command: 'devices', session: 'default', positionals: [], From db4ee91e4afaea371da3b87acda13c4fc9e9e475 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sun, 4 Oct 2026 00:00:11 +0200 Subject: [PATCH 07/14] test(daemon-client): prove an abort is never a timeout at the transport seam Add the paired regression cases for #3178's never-a-timeout promise: with the timeout seam mocked, an abort with an armed budget leaves the seam never called, while the same budget without a signal reaches it. Counterfactual verified: deleting the abort path's `clearTimeout` fails only the abort case. Clear the ignoring-transport test's late resolve so it keeps no timer alive. --- src/__tests__/client-abort-signal.test.ts | 11 +- .../daemon-client-abort-timeout.test.ts | 108 ++++++++++++++++++ 2 files changed, 117 insertions(+), 2 deletions(-) create mode 100644 src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts diff --git a/src/__tests__/client-abort-signal.test.ts b/src/__tests__/client-abort-signal.test.ts index 7bdf69e245..f9e6f6073f 100644 --- a/src/__tests__/client-abort-signal.test.ts +++ b/src/__tests__/client-abort-signal.test.ts @@ -53,14 +53,20 @@ test('the signal rides the transport context so a built-in-style transport can c }); test('aborting a call whose transport ignores the signal still rejects the caller with the typed canceled error', async () => { + let lingering: ReturnType | undefined; const client = createAgentDeviceClient( {}, { transport: (req) => new Promise((resolve) => { // A custom transport that never inspects the context signal: the guard must still settle - // the caller's promise when the abort fires. - setTimeout(() => resolve({ ok: true, data: { ignored: req.command } }), 10_000); + // the caller's promise when the abort fires. The late resolve is cleared once the caller + // has been rejected, so the worker keeps no timer alive for its own promise. + lingering = setTimeout( + () => resolve({ ok: true, data: { ignored: req.command } }), + 10_000, + ); + lingering.unref?.(); }), }, ); @@ -68,6 +74,7 @@ test('aborting a call whose transport ignores the signal still rejects the calle const call = client.interactions.press({ ref: '@e12', signal: controller.signal }); setTimeout(() => controller.abort(), 10); await assert.rejects(call, (error: unknown) => canceledWith(error, 'unknown')); + if (lingering) clearTimeout(lingering); }); test('a signal on one call does not cancel another', async () => { diff --git a/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts b/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts new file mode 100644 index 0000000000..31cae735e6 --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts @@ -0,0 +1,108 @@ +/** + * #3178: an abort is never a timeout. `handleRequestTimeout` is the one seam that sweeps runner + * processes and resets the local daemon, so an abort that settles a request must also disarm that + * request's timeout timer. The seam is mocked here — the subject is the transport's timer + * discipline, not the sweep — and each case outlives the armed budget, so the assertion proves the + * timer was cleared rather than merely beaten to the settle. + * + * The abort case aborts right after `sendRequest` arms its timer, so the clear is what is tested. + * The paired no-signal case lets the same budget expire and must reach the mocked seam, which + * makes the abort case's empty-calls assertion evidence instead of an absent hook: delete the + * abort path's `clearTimeout` and only the abort case fails. + */ +import assert from 'node:assert/strict'; +import { test, vi } from 'vitest'; +import { AppError, isRequestCanceledError } from '@agent-device/kernel/errors'; +import { getRequestSignal } from '@agent-device/host-kit/request'; +import type { DaemonInvokeFn, DaemonRequest, DaemonResponse } from '../../daemon/daemon-request.ts'; +import { createSocketServer, listenNetServer } from '../../daemon/server/transport.ts'; +import { sendRequest } from '../daemon-client-transport.ts'; +import { resolveDaemonPaths } from '../../daemon-resolution.ts'; +import { + closeLoopbackServer, + skipWhenLoopbackUnavailable, + type SkippableTestContext, +} from '../../__tests__/test-utils/loopback.ts'; +import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; + +const { handleRequestTimeoutCalls } = vi.hoisted(() => ({ + handleRequestTimeoutCalls: [] as unknown[], +})); + +vi.mock('../daemon-client-timeout.ts', () => ({ + handleRequestTimeout: (params: unknown) => { + handleRequestTimeoutCalls.push(params); + return new AppError('COMMAND_FAILED', 'Daemon request timed out', { + reason: 'daemon_transport_timeout', + }); + }, +})); + +const STATE_PATHS = resolveDaemonPaths(mkdtempForTestSync('agent-device-abort-timeout-')); +const TIMEOUT_MS = 15; +// Past the armed budget, so a timer the abort forgot to clear fires inside the case. +const OUTLIVE_MS = 80; + +// A request the daemon never answers: the only things that can end it are the client's abort and +// the transport's own budget, which is what makes the timer's fate observable. +function hangingHandler(): DaemonInvokeFn { + return async (req: DaemonRequest): Promise => { + const signal = getRequestSignal(req.meta?.requestId); + return await new Promise((_resolve, reject) => { + signal?.addEventListener('abort', () => reject(signal.reason), { once: true }); + }); + }; +} + +async function runTimerDiscipline( + t: SkippableTestContext, + outcome: 'aborted' | 'timed-out', +): Promise { + if (await skipWhenLoopbackUnavailable(t)) return; + const aborted = outcome === 'aborted'; + handleRequestTimeoutCalls.length = 0; + const server = createSocketServer(hangingHandler()); + try { + const port = await listenNetServer(server); + const controller = new AbortController(); + const request = sendRequest( + { port, token: 't', pid: 1 }, + { + token: 't', + command: 'wait', + session: 'default', + positionals: [], + flags: {}, + meta: { requestId: `req-${outcome}-timeout-timer` }, + }, + 'socket', + STATE_PATHS, + TIMEOUT_MS, + aborted ? { signal: controller.signal } : {}, + ); + if (aborted) { + // `sendRequest` armed the timeout timer synchronously by returning; this abort must disarm it. + controller.abort(); + await assert.rejects(request, (error: unknown) => isRequestCanceledError(error)); + } else { + await assert.rejects( + request, + (error: unknown) => + error instanceof AppError && error.details?.reason === 'daemon_transport_timeout', + ); + assert.equal(handleRequestTimeoutCalls.length, 1); + } + await new Promise((resolve) => setTimeout(resolve, OUTLIVE_MS)); + if (aborted) assert.deepEqual(handleRequestTimeoutCalls, []); + } finally { + await closeLoopbackServer(server); + } +} + +test('an abort with a timeout armed clears the timer instead of running the timeout path', async (t) => { + await runTimerDiscipline(t, 'aborted'); +}); + +test('the same budget without a signal reaches the timeout seam', async (t) => { + await runTimerDiscipline(t, 'timed-out'); +}); From 10eb9d6a3c9eb477b75f24acabdaf6274f51e5d5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sun, 4 Oct 2026 08:55:23 +0200 Subject: [PATCH 08/14] fix(client): own the per-call guard beside its transport and answer the review - Move the requester-side AbortSignal guard out of @agent-device/host-kit/request into daemon-client-transport.ts, the module the CLI eager closure already loads, so the caller-side half and the transport half share one owner and the ADR-0019 loading-shape probe stops counting 9 extra modules against cli.ts and the replay-port facade. - Generate the request id in execute before guarding, so a canceled call names itself in request.meta and matches the daemon's diagnostic for the same request. - Report dispatched 'no' when a guard is entered already aborted: the refusal proves send never ran, so 'unknown' would suppress a safe retry. - Keep a non-Error abort reason as the cause unchanged. - Return before handleTransportError once a request is settled, so the error an intentional destroy emits cannot log daemon_request_socket_error. - Cut the remote-retry health probe on the caller's signal and recheck the signal before timeout and retry handling, so a canceled call never reads as timed out. - Forward the caller's abort into the lease beat, which owns its own connection. - Prove the signal reaches the daemon through the DEFAULT transport: one loopback test per transport drives createAgentDeviceClient with no injected transport and asserts the daemon saw the cancel. - Prove the RPC handler still serves after an HTTP cancellation, not just /health. - Qualify the signal contract and the docs with the two phases the guarantee does not cover: a response-artifact download already underway, and a canceled one-shot replay's existing daemon cleanup. --- packages/contracts/src/client-connection.ts | 6 + .../src/internal/request-guard.test.ts | 69 -------- .../host-kit/src/internal/request-guard.ts | 93 ----------- packages/host-kit/src/request.ts | 6 - src/__tests__/client-abort-signal.test.ts | 7 +- src/agent-device-client.ts | 14 +- ...mon-client-abort-default-transport.test.ts | 154 ++++++++++++++++++ .../__tests__/daemon-client-abort.test.ts | 31 +++- .../daemon-client-lease-beat.test.ts | 40 ++++- .../daemon-client-request-guard.test.ts | 112 +++++++++++++ src/daemon-client/daemon-client-lease-beat.ts | 36 +++- src/daemon-client/daemon-client-transport.ts | 124 ++++++++++++-- src/daemon-client/daemon-client.ts | 11 +- test/wire-compat/ledger.json | 8 +- website/docs/docs/client-api.md | 4 +- 15 files changed, 501 insertions(+), 214 deletions(-) delete mode 100644 packages/host-kit/src/internal/request-guard.test.ts delete mode 100644 packages/host-kit/src/internal/request-guard.ts create mode 100644 src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts create mode 100644 src/daemon-client/__tests__/daemon-client-request-guard.test.ts diff --git a/packages/contracts/src/client-connection.ts b/packages/contracts/src/client-connection.ts index 72c62df725..52d6bc78ce 100644 --- a/packages/contracts/src/client-connection.ts +++ b/packages/contracts/src/client-connection.ts @@ -102,6 +102,12 @@ export type AgentDeviceRequestOverrides = Pick< * marks the request canceled, and the promise rejects with the typed canceled-request error * (`details.reason: 'request_canceled'`). An abort is never a timeout: no runner sweep, no * daemon reset. + * + * The guarantee covers the daemon request, and the built-in transports enforce it; a custom + * transport receives the signal on its context and may implement cancellation differently. Two + * phases run outside it: a response-artifact download started after the response begins is not + * canceled, and a canceled one-shot replay still runs the existing cleanup that may tear down a + * daemon this client started. */ signal?: AbortSignal; }; diff --git a/packages/host-kit/src/internal/request-guard.test.ts b/packages/host-kit/src/internal/request-guard.test.ts deleted file mode 100644 index 52a6622139..0000000000 --- a/packages/host-kit/src/internal/request-guard.test.ts +++ /dev/null @@ -1,69 +0,0 @@ -import { test } from 'vitest'; -import assert from 'node:assert/strict'; -import { isRequestCanceledError } from '@agent-device/kernel/errors'; -import { createRequestGuard } from './request-guard.ts'; - -function canceledOutcome(promise: Promise): Promise { - return promise.then( - () => { - throw new Error('expected rejection'); - }, - (error: unknown) => error, - ); -} - -test('a guard without a signal never refuses and never interferes', async () => { - const guard = createRequestGuard({ signal: undefined, requestId: 'req-none' }); - guard.refuseIfAborted(); - assert.equal(await guard.guard(async () => 'sent'), 'sent'); -}); - -test('refuseIfAborted refuses an already-aborted call with dispatched no', async () => { - const controller = new AbortController(); - controller.abort(new Error('caller gave up')); - const guard = createRequestGuard({ signal: controller.signal, requestId: 'req-pre' }); - assert.throws( - () => guard.refuseIfAborted(), - (error: unknown) => - isRequestCanceledError(error) && - (error as { details?: Record }).details?.dispatched === 'no' && - (error as { details?: Record }).details?.requestId === 'req-pre', - ); -}); - -test('guard settles an in-flight send with the typed canceled error whatever the abort reason', async () => { - const controller = new AbortController(); - const guard = createRequestGuard({ signal: controller.signal, requestId: 'req-flight' }); - let settleSend: (() => void) | undefined; - const pending = guard.guard( - () => new Promise((resolve) => (settleSend = () => resolve('late'))), - ); - controller.abort(new Error('anything')); - const error = await canceledOutcome(pending); - assert.equal(isRequestCanceledError(error), true); - assert.equal((error as { details?: Record }).details?.dispatched, 'unknown'); - assert.equal((error as { details?: Record }).details?.requestId, 'req-flight'); - // A transport that ignores the signal may still resolve later; that must not - // turn the already-rejected caller promise into a second outcome. - settleSend?.(); -}); - -test('guard keeps a send error when no abort ever fires', async () => { - const controller = new AbortController(); - const guard = createRequestGuard({ signal: controller.signal }); - const failure = new Error('daemon refused'); - await assert.rejects( - guard.guard(async () => { - throw failure; - }), - (error: unknown) => error === failure, - ); -}); - -test('guard returns the send value when it settles before any abort', async () => { - const controller = new AbortController(); - const guard = createRequestGuard({ signal: controller.signal }); - const value = await guard.guard(async () => 'response'); - controller.abort(); - assert.equal(value, 'response'); -}); diff --git a/packages/host-kit/src/internal/request-guard.ts b/packages/host-kit/src/internal/request-guard.ts deleted file mode 100644 index 673403da9f..0000000000 --- a/packages/host-kit/src/internal/request-guard.ts +++ /dev/null @@ -1,93 +0,0 @@ -import { createRequestCanceledError, type AppError } from '@agent-device/kernel/errors'; - -/** - * The requester-side half of a per-call `AbortSignal`. The transport half closes the request's - * connection so the daemon marks the request canceled (its own cancellation machinery is - * {@link ./request-cancel.ts}); this half answers for the two cases a transport cannot: - * - * - a signal already aborted before anything was sent — nothing may leave the process, so the - * refusal carries `details.dispatched: 'no'`; - * - a custom transport that ignores the signal it was handed — the caller's promise must still - * settle when the abort fires, with `details.dispatched: 'unknown'` because the transport may - * already have written the request. - * - * Both outcomes reject with the typed canceled-request error (`details.reason: - * 'request_canceled'`) whatever reason the caller's own controller aborted with, so every layer - * that lets a cancellation through dispatches on the same reason — and never on the caller's - * arbitrary abort reason or on error text. An abort is never a timeout: nothing here reads or - * extends a request deadline, and no timeout path may produce this rejection. - */ - -/** - * The typed canceled-request error for a caller's own abort: the reason the caller aborted with - * survives as the cause, while the rejection itself always dispatches on `reason: - * 'request_canceled'` with the delivery evidence the aborting layer can prove. A built-in transport - * rejects every abort through this, so a caller's arbitrary abort reason never escapes as the - * outcome of a daemon request. - */ -export function abortedRequestError( - signal: AbortSignal, - dispatched: 'no' | 'unknown', - requestId?: string, -): AppError { - return createRequestCanceledError( - { requestId, dispatched }, - signal.reason instanceof Error ? signal.reason : undefined, - ); -} - -/** - * Refuse a send attempt that starts while the caller's signal is already aborted: nothing may leave - * the process, so this never touches a connection and the refusal carries `details.dispatched: 'no'`. - */ -export function refuseAbortedRequest(signal: AbortSignal | undefined, requestId?: string): void { - if (!signal?.aborted) return; - throw abortedRequestError(signal, 'no', requestId); -} - -export type RequestGuard = { - /** Refuses an already-aborted call before anything is sent. No-op without a signal. */ - refuseIfAborted(): void; - /** Settles `send`'s outcome against the signal, winning with the typed canceled error on abort. */ - guard(send: () => Promise): Promise; -}; - -const NO_REQUEST_GUARD: RequestGuard = { - refuseIfAborted: () => {}, - guard: async (send: () => Promise) => await send(), -}; - -export function createRequestGuard(params: { - signal: AbortSignal | undefined; - requestId?: string; -}): RequestGuard { - const { signal, requestId } = params; - if (!signal) return NO_REQUEST_GUARD; - const canceled = (dispatched: 'no' | 'unknown'): AppError => - createRequestCanceledError( - { requestId, dispatched }, - signal.reason instanceof Error ? signal.reason : undefined, - ); - return { - refuseIfAborted() { - if (signal.aborted) throw canceled('no'); - }, - guard: async (send: () => Promise): Promise => { - if (signal.aborted) throw canceled('unknown'); - return await new Promise((resolve, reject) => { - const onAbort = (): void => reject(canceled('unknown')); - signal.addEventListener('abort', onAbort, { once: true }); - void send().then( - (value) => { - signal.removeEventListener('abort', onAbort); - resolve(value); - }, - (error: unknown) => { - signal.removeEventListener('abort', onAbort); - reject(error); - }, - ); - }); - }, - }; -} diff --git a/packages/host-kit/src/request.ts b/packages/host-kit/src/request.ts index 684ade7e65..ee98a01e7b 100644 --- a/packages/host-kit/src/request.ts +++ b/packages/host-kit/src/request.ts @@ -9,9 +9,3 @@ export { throwIfRequestCanceled, } from './internal/request-cancel.ts'; export { emitRequestProgress, withRequestProgressSink } from './internal/request-progress.ts'; -export { - abortedRequestError, - createRequestGuard, - refuseAbortedRequest, - type RequestGuard, -} from './internal/request-guard.ts'; diff --git a/src/__tests__/client-abort-signal.test.ts b/src/__tests__/client-abort-signal.test.ts index f9e6f6073f..33af448513 100644 --- a/src/__tests__/client-abort-signal.test.ts +++ b/src/__tests__/client-abort-signal.test.ts @@ -92,7 +92,12 @@ test('a signal on one call does not cancel another', async () => { const canceled = new AbortController(); const doomed = client.command.wait({ durationMs: 5000, signal: canceled.signal }); const survivor = client.command.wait({ durationMs: 1 }); + // Deferred so the doomed call is genuinely in flight when the abort fires: a synchronous abort + // would land before `execute` installs the guard and reject through the pre-abort `no` path that + // the first test already covers. The guard answers for a transport that ignores the signal, so + // this exercises the in-flight `unknown` rejection while the survivor runs untouched. + await new Promise((resolve) => setTimeout(resolve, 10)); canceled.abort(); - await assert.rejects(doomed, (error: unknown) => isRequestCanceledError(error)); + await assert.rejects(doomed, (error: unknown) => canceledWith(error, 'unknown')); assert.deepEqual(await survivor, {}); }); diff --git a/src/agent-device-client.ts b/src/agent-device-client.ts index f2848384af..1273790a39 100644 --- a/src/agent-device-client.ts +++ b/src/agent-device-client.ts @@ -35,6 +35,7 @@ import { type SessionRuntimeHints, } from '@agent-device/kernel/contracts'; import { AppError, throwDaemonError } from '@agent-device/kernel/errors'; +import { createRequestId } from '@agent-device/host-kit/diagnostics'; import { buildMeta, normalizeDeployResult, @@ -74,7 +75,7 @@ import { isRecord, readSnapshotKeyboardBandFact } from '@agent-device/kernel/rec import { readResponseWarnings } from '@agent-device/kernel/success-text'; import { createLeaseClient } from './client/lease-client.ts'; import { normalizeScreenshotCaptureResult } from './client/screenshot-result.ts'; -import { createRequestGuard } from '@agent-device/host-kit/request'; +import { createRequestGuard } from './daemon-client/daemon-client-transport.ts'; export function createAgentDeviceClient( config: AgentDeviceClientConfig = {}, @@ -96,10 +97,11 @@ export function createAgentDeviceClient( input?: Record, ): Promise> => { const merged = mergeClientOptions(config, options); - const cancellation = createRequestGuard({ - signal: merged.signal, - requestId: merged.requestId, - }); + // The id is generated before the guard so a canceled call names itself in `request.meta`, and + // the daemon's diagnostics for the request the transport goes on to cancel carry the same id + // the caller's rejection carries. + const requestId = merged.requestId ?? createRequestId(); + const cancellation = createRequestGuard({ signal: merged.signal, requestId }); cancellation.refuseIfAborted(); const request = { session: resolveSessionName(merged.session), @@ -108,7 +110,7 @@ export function createAgentDeviceClient( ...(input ? { input } : {}), flags: buildRequestFlags(merged, metadataFlags), runtime: merged.runtime, - meta: buildMeta(merged), + meta: { ...buildMeta(merged), requestId }, }; // `signal` rides the transport context (it is a live object, never wire data), and the guard // answers for a custom transport that ignores it: the caller's promise settles on abort either diff --git a/src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts b/src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts new file mode 100644 index 0000000000..f2d18b49c4 --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts @@ -0,0 +1,154 @@ +/** + * #3178: the real client route, through the default transport. + * + * `daemon-client-abort.test.ts` drives `sendRequest` directly, and `client-abort-signal.test.ts` + * injects a fake transport, so neither proves that the caller's signal actually reaches the + * transport from `createAgentDeviceClient`. That property lives at one line — the signal handed to + * `sendRequest` in `daemon-client.ts` — and a test that skips that wiring cannot see it disappear: + * delete the hand-off and a fake-transport test still passes, because the client-level guard + * rejects the caller with `request_canceled` while the daemon keeps running the command on the + * device. That is exactly the bug #3178 reports. + * + * Each case builds a client with NO injected transport, points its state dir at a loopback daemon + * (a real `createSocketServer` / `createDaemonHttpServer` published through `daemon.json`), aborts + * a long `wait` mid-flight, and asserts the DAEMON observed the cancellation — plus that the id in + * the caller's rejection is the id the daemon canceled, which is what lets a caller correlate a + * canceled call with the daemon's own diagnostics. + */ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import { test } from 'vitest'; +import { isRequestCanceledError } from '@agent-device/kernel/errors'; +import { getRequestSignal, isRequestCanceled } from '@agent-device/host-kit/request'; +import { readProcessStartTime } from '@agent-device/host-kit/process'; +import { readVersion } from '@agent-device/host-kit/version'; +import type { DaemonInvokeFn, DaemonRequest, DaemonResponse } from '../../daemon/daemon-request.ts'; +import { createDaemonHttpServer } from '../../daemon/server/http-server.ts'; +import { createSocketServer, listenNetServer } from '../../daemon/server/transport.ts'; +import { createAgentDeviceClient } from '../../agent-device-client.ts'; +import { currentDaemonCodeSignature } from '../../__tests__/test-utils/daemon-http-fixture.ts'; +import { + closeLoopbackServer, + listenOnLoopback, + skipWhenLoopbackUnavailable, + type SkippableTestContext, +} from '../../__tests__/test-utils/loopback.ts'; +import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; + +const TOKEN = 'default-transport-abort-token'; + +type SeenRequest = { started: boolean; canceled: boolean[]; requestId?: string }; + +// `wait` is the long request the caller abandons; it never answers on its own, so only the client's +// disconnect ends it — through the daemon's own request-scoped cancellation. +function hangingWaitHandler(seen: SeenRequest): DaemonInvokeFn { + return async (req: DaemonRequest): Promise => { + const requestId = req.meta?.requestId; + seen.started = true; + seen.requestId = requestId; + const signal = getRequestSignal(requestId); + if (!signal) { + return { ok: true, data: { answeredWithoutSignal: true } }; + } + return await new Promise((_resolve, reject) => { + signal.addEventListener( + 'abort', + () => { + seen.canceled.push(isRequestCanceled(requestId)); + reject(signal.reason); + }, + { once: true }, + ); + }); + }; +} + +// The `daemon.json` a running daemon publishes, shaped so the client's takeover ladder reuses it: +// matching version and code signature, then the live probe the loopback server answers. +function publishLoopbackDaemonInfo(stateDir: string, info: Record): void { + fs.mkdirSync(stateDir, { recursive: true }); + fs.writeFileSync( + `${stateDir}/daemon.json`, + `${JSON.stringify({ + token: TOKEN, + pid: process.pid, + version: readVersion(), + codeSignature: currentDaemonCodeSignature(), + processStartTime: readProcessStartTime(process.pid) ?? undefined, + ...info, + })}\n`, + ); +} + +async function waitFor(condition: () => boolean, what: string): Promise { + const deadline = Date.now() + 3000; + while (!condition()) { + if (Date.now() >= deadline) throw new Error(`Timed out waiting for ${what}.`); + await new Promise((resolve) => setTimeout(resolve, 10)); + } +} + +async function runDefaultTransportAbort( + t: SkippableTestContext, + transport: 'socket' | 'http', + serve: ( + handleRequest: DaemonInvokeFn, + ) => Promise<{ server: Parameters[0]; port: number }>, +): Promise { + if (await skipWhenLoopbackUnavailable(t)) return; + const seen: SeenRequest = { started: false, canceled: [] }; + const stateDir = mkdtempForTestSync(`agent-device-abort-default-${transport}-`); + const { server, port } = await serve(hangingWaitHandler(seen)); + try { + publishLoopbackDaemonInfo( + stateDir, + transport === 'socket' + ? { port, transport: 'socket' } + : { httpPort: port, transport: 'http' }, + ); + // No injected transport: the client's own `sendToDaemon` and default transport have to carry the + // caller signal all the way to the connection they open. + const client = createAgentDeviceClient({ stateDir, daemonTransport: transport }); + const controller = new AbortController(); + + const call = client.command.wait({ durationMs: 5000, signal: controller.signal }); + await waitFor(() => seen.started, 'the daemon to start the request'); + controller.abort(); + + const rejection = await call.then( + () => { + throw new Error('expected the aborted call to reject'); + }, + (error: unknown) => error, + ); + assert.equal(isRequestCanceledError(rejection), true); + const details = (rejection as { details?: Record }).details ?? {}; + assert.equal(details.dispatched, 'unknown'); + // The proof the fake-transport tests cannot offer: the cancellation crossed the wire, because + // the daemon's own request registry saw it. Delete the signal hand-off in `daemon-client.ts` + // and this stays empty while the caller still rejects through its guard — the #3178 bug. + await waitFor(() => seen.canceled.length > 0, 'daemon-side cancellation'); + assert.deepEqual(seen.canceled, [true]); + // The caller's rejection names the same request the daemon canceled, so a caller can correlate + // the two records. + assert.equal(typeof details.requestId, 'string'); + assert.equal(details.requestId, seen.requestId); + } finally { + await closeLoopbackServer(server); + fs.rmSync(stateDir, { recursive: true, force: true }); + } +} + +test('the default socket transport carries the caller signal: the daemon cancels the request the client aborted', async (t) => { + await runDefaultTransportAbort(t, 'socket', async (handleRequest) => { + const server = createSocketServer(handleRequest); + return { server, port: await listenNetServer(server) }; + }); +}); + +test('the default HTTP transport carries the caller signal: the daemon cancels the request the client aborted', async (t) => { + await runDefaultTransportAbort(t, 'http', async (handleRequest) => { + const server = await createDaemonHttpServer({ handleRequest, token: TOKEN }); + return { server, port: await listenOnLoopback(server) }; + }); +}); diff --git a/src/daemon-client/__tests__/daemon-client-abort.test.ts b/src/daemon-client/__tests__/daemon-client-abort.test.ts index 3c9dfbca8c..2837d82c09 100644 --- a/src/daemon-client/__tests__/daemon-client-abort.test.ts +++ b/src/daemon-client/__tests__/daemon-client-abort.test.ts @@ -27,7 +27,13 @@ import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; const TOKEN = 'abort-signal-token'; const STATE_PATHS = resolveDaemonPaths(mkdtempForTestSync('agent-device-client-abort-')); -type SeenRequest = { started: boolean; requestId?: string; canceled: boolean[] }; +type SeenRequest = { + started: boolean; + requestId?: string; + canceled: boolean[]; + /** The command of the last non-`wait` request the RPC handler itself served. */ + servedCommand?: string; +}; // `wait` models the long request the caller abandons; every other command answers immediately, so // the same handler also proves the daemon still serves requests after one was canceled. @@ -36,6 +42,7 @@ function canceledAwareHandler(seen: SeenRequest): DaemonInvokeFn { const requestId = req.meta?.requestId; seen.requestId = requestId; if (req.command !== 'wait') { + seen.servedCommand = req.command; return { ok: true, data: {} }; } seen.started = true; @@ -150,6 +157,7 @@ test('socket transport: an abort mid-request closes the connection, the daemon m {}, ); assert.equal(followUp.ok, true); + assert.equal(seen.servedCommand, 'devices'); } finally { await closeLoopbackServer(server); } @@ -226,8 +234,25 @@ test('http transport: an abort mid-request closes that request, the daemon marks assert.deepEqual(seen.canceled, [true]); assert.equal(seen.requestId, 'req-http-abort-in-flight'); - const health = await fetch(`http://127.0.0.1:${port}/health`); - assert.equal(health.ok, true); + const followUp = await sendRequest( + info, + { + token: TOKEN, + command: 'devices', + session: 'default', + positionals: [], + flags: {}, + meta: { requestId: 'req-http-follow-up' }, + }, + 'http', + STATE_PATHS, + 5000, + {}, + ); + assert.equal(followUp.ok, true); + // The follow-up ran through the RPC path `/health` does not: an aborted request's cancel must + // leave `handleRequest` serving, not just the health endpoint answering. + assert.equal(seen.servedCommand, 'devices'); } finally { await closeLoopbackServer(server); } diff --git a/src/daemon-client/__tests__/daemon-client-lease-beat.test.ts b/src/daemon-client/__tests__/daemon-client-lease-beat.test.ts index 858feae913..7d1a01d995 100644 --- a/src/daemon-client/__tests__/daemon-client-lease-beat.test.ts +++ b/src/daemon-client/__tests__/daemon-client-lease-beat.test.ts @@ -375,6 +375,37 @@ describe('runProtectedLeaseWork', () => { assert.equal(sawAbort, true, 'the upload is told to stop before the bytes finish'); }); + test('a caller abort stops the beats and cancels the beat already in flight', async () => { + vi.useFakeTimers(); + // #3178: a beat owns a connection of its own, so a caller that has given up has to reach it — + // otherwise a renewal lands for a request nobody is waiting on anymore. + const caller = new AbortController(); + const beatSignals: AbortSignal[] = []; + const upload = deferred(); + const running = runProtectedLeaseWork({ + task: () => upload.promise, + heartbeat: (_budgetMs, signal) => { + beatSignals.push(signal); + return new Promise(() => undefined); + }, + callerSignal: caller.signal, + }); + + await vi.advanceTimersByTimeAsync(0); + assert.equal(beatSignals.length, 1, 'the opening beat started'); + assert.equal(beatSignals[0]?.aborted, false); + + caller.abort(); + await vi.advanceTimersByTimeAsync(20_000); + assert.equal(beatSignals[0]?.aborted, true, 'the in-flight beat lost its connection'); + assert.equal(beatSignals.length, 1, 'the caller abort ends the schedule, not just the beat'); + + // The phase belongs to the caller: the loop keeps no timer alive past the abort, so the upload + // still decides the outcome. + upload.resolve('uploaded'); + assert.equal(await running, 'uploaded'); + }); + test('a phase that throws synchronously still stops the beats', async () => { vi.useFakeTimers(); const heartbeat = vi.fn(async () => ({ ok: true })); @@ -620,10 +651,10 @@ describe('buildUploadLeaseHeartbeat', () => { installRequest, ); assert.ok(beat); - await beat!(5_000); + await beat!(5_000, new AbortController().signal); // Two beats, because a beat that times out is canceled under its own id: sharing one would // let a later beat inherit an earlier cancellation and stop renewing a live lease. - await beat!(5_000); + await beat!(5_000, new AbortController().signal); } finally { // The server keeps the connection open, and `close()` waits for it. for (const connection of connections) connection.destroy(); @@ -666,7 +697,10 @@ describe('buildUploadLeaseHeartbeat', () => { installRequest, ); assert.ok(beat); - const rejected = assert.rejects((async () => await beat!(1_000))(), /timed out/i); + const rejected = assert.rejects( + (async () => await beat!(1_000, new AbortController().signal))(), + /timed out/i, + ); await vi.advanceTimersByTimeAsync(1_000); await rejected; } finally { diff --git a/src/daemon-client/__tests__/daemon-client-request-guard.test.ts b/src/daemon-client/__tests__/daemon-client-request-guard.test.ts new file mode 100644 index 0000000000..2746df91ca --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-request-guard.test.ts @@ -0,0 +1,112 @@ +/** + * #3178: the requester-side half of the per-call `AbortSignal`, owned by the transport module that + * enforces the other half. These pin the guard's contract: what it refuses before sending + * (`details.dispatched: 'no'`), what it answers when an abort lands while a send is outstanding + * (`details.dispatched: 'unknown'`), and that the caller's own abort reason survives as the cause + * whatever its shape — never as the rejection's reason, which always dispatches as + * `request_canceled`. + */ +import { test } from 'vitest'; +import assert from 'node:assert/strict'; +import { isRequestCanceledError } from '@agent-device/kernel/errors'; +import { createRequestGuard } from '../daemon-client-transport.ts'; + +function canceledOutcome(promise: Promise): Promise { + return promise.then( + () => { + throw new Error('expected rejection'); + }, + (error: unknown) => error, + ); +} + +function detailsOf(error: unknown): Record { + return (error as { details?: Record }).details ?? {}; +} + +test('a guard without a signal never refuses and never interferes', async () => { + const guard = createRequestGuard({ signal: undefined, requestId: 'req-none' }); + guard.refuseIfAborted(); + assert.equal(await guard.guard(async () => 'sent'), 'sent'); +}); + +test('refuseIfAborted refuses an already-aborted call with dispatched no', async () => { + const controller = new AbortController(); + controller.abort(new Error('caller gave up')); + const guard = createRequestGuard({ signal: controller.signal, requestId: 'req-pre' }); + assert.throws( + () => guard.refuseIfAborted(), + (error: unknown) => + isRequestCanceledError(error) && + detailsOf(error).dispatched === 'no' && + detailsOf(error).requestId === 'req-pre', + ); +}); + +test('a guard called with an already-aborted signal refuses before send with dispatched no', async () => { + // Load-bearing branch: an abort listener attached to an already-aborted signal never fires, so + // without this check the guard would pend forever instead of rejecting. The refusal also proves + // `send` never ran, so the honest evidence is 'no' rather than 'unknown'. + const controller = new AbortController(); + controller.abort(); + const guard = createRequestGuard({ signal: controller.signal, requestId: 'req-late' }); + const error = await canceledOutcome( + guard.guard(async () => { + throw new Error('send must not run'); + }), + ); + assert.equal(isRequestCanceledError(error), true); + assert.equal(detailsOf(error).dispatched, 'no'); + assert.equal(detailsOf(error).requestId, 'req-late'); +}); + +test('guard settles an in-flight send with the typed canceled error whatever the abort reason', async () => { + const controller = new AbortController(); + const guard = createRequestGuard({ signal: controller.signal, requestId: 'req-flight' }); + let settleSend: (() => void) | undefined; + const pending = guard.guard( + () => new Promise((resolve) => (settleSend = () => resolve('late'))), + ); + controller.abort(new Error('anything')); + const error = await canceledOutcome(pending); + assert.equal(isRequestCanceledError(error), true); + assert.equal(detailsOf(error).dispatched, 'unknown'); + assert.equal(detailsOf(error).requestId, 'req-flight'); + // A transport that ignores the signal may still resolve later; that must not + // turn the already-rejected caller promise into a second outcome. + settleSend?.(); +}); + +test('guard keeps a non-Error abort reason as the cause unchanged', async () => { + // `AbortController.abort()` accepts any value, and the helper promises the reason survives as + // the cause: a bare string or a number must not be dropped for want of being an Error, while the + // rejection itself still dispatches on request_canceled. + for (const reason of ['gave-up', 42]) { + const controller = new AbortController(); + controller.abort(reason); + const guard = createRequestGuard({ signal: controller.signal }); + const error = await canceledOutcome(guard.guard(() => new Promise(() => undefined))); + assert.equal(isRequestCanceledError(error), true); + assert.equal((error as { cause?: unknown }).cause, reason); + } +}); + +test('guard keeps a send error when no abort ever fires', async () => { + const controller = new AbortController(); + const guard = createRequestGuard({ signal: controller.signal }); + const failure = new Error('daemon refused'); + await assert.rejects( + guard.guard(async () => { + throw failure; + }), + (error: unknown) => error === failure, + ); +}); + +test('guard returns the send value when it settles before any abort', async () => { + const controller = new AbortController(); + const guard = createRequestGuard({ signal: controller.signal }); + const value = await guard.guard(async () => 'response'); + controller.abort(); + assert.equal(value, 'response'); +}); diff --git a/src/daemon-client/daemon-client-lease-beat.ts b/src/daemon-client/daemon-client-lease-beat.ts index adfd4c9f07..cb489c8010 100644 --- a/src/daemon-client/daemon-client-lease-beat.ts +++ b/src/daemon-client/daemon-client-lease-beat.ts @@ -104,27 +104,41 @@ export async function runProtectedLeaseWork( * One renewal. `budgetMs` is how long this beat may take before the loop abandons it: the window * the beat protects, so a slow round trip still gets a chance to answer inside the lease it is * renewing. Absent when the request names no lease to renew, which is the ordinary unleased - * install. + * install. `signal` is this beat's own cancellation: the caller's abort arrives on it, so a beat + * already in flight when the caller gave up closes its connection instead of finishing a + * renewal for work nobody is waiting on anymore. */ - heartbeat?: ((budgetMs: number) => Promise) | undefined; + heartbeat?: ((budgetMs: number, signal: AbortSignal) => Promise) | undefined; task: (signal: AbortSignal) => Promise; + /** The caller's per-call signal (#3178), observed between and during beats. */ + callerSignal?: AbortSignal | undefined; }>, ): Promise { - const { heartbeat } = options; + const { heartbeat, callerSignal } = options; if (!heartbeat) return await options.task(new AbortController().signal); const control = new AbortController(); + // A beat owns a connection of its own, so the caller's abort is forwarded to it rather than left + // to the phase settle: without this, a beat already in flight renews the lease after the caller + // has given up. + const beatControl = new AbortController(); // Until a beat names the window, the loop assumes the shortest window the daemon will accept: a // beat that budgets itself on a longer window than the lease actually has would outlive it. let windowMs = MIN_LEASE_WINDOW_MS; let intervalMs = leaseBeatIntervalMs(windowMs); let timer: ReturnType | undefined; let stopped = false; + const stopBeating = (): void => { + stopped = true; + if (timer) clearTimeout(timer); + beatControl.abort(); + }; let terminalError: unknown; let reportTerminal: ((error: unknown) => void) | undefined; const terminal = new Promise((_, reject) => { reportTerminal = reject; }); + if (callerSignal) callerSignal.addEventListener('abort', stopBeating, { once: true }); const runBeat = (): void => { // Armed while this beat is still outstanding: a beat that never settles is abandoned on @@ -138,7 +152,9 @@ export async function runProtectedLeaseWork( const settle = (async () => { const budgetMs = windowMs; try { - const renewed = leaseWindowFromHeartbeatResponse(await heartbeat(budgetMs)); + const renewed = leaseWindowFromHeartbeatResponse( + await heartbeat(budgetMs, beatControl.signal), + ); // An answer that names no window keeps the cadence it was asked at: the loop only ever // moves on evidence of how long the lease is good for, and never on the absence of it. if (renewed === undefined) return; @@ -167,6 +183,9 @@ export async function runProtectedLeaseWork( reportTerminal?.(error); return; } + // A beat the loop itself canceled — for a caller that has already given up — failed for a + // reason that says nothing about the lease, so it is not worth a warning. + if (stopped) return; emitDiagnostic({ level: 'warn', phase: 'lease_heartbeat_failed', @@ -190,8 +209,8 @@ export async function runProtectedLeaseWork( // timer below is always cleared. (async () => await Promise.race([options.task(control.signal), terminal]))(), ); - stopped = true; - if (timer) clearTimeout(timer); + stopBeating(); + if (callerSignal) callerSignal.removeEventListener('abort', stopBeating); // No outstanding beat is awaited here: a beat on a half-open connection would hold a finished // upload behind its own budget for no decision the phase still has to make. // A beat that ended the protection outranks a phase that settled meanwhile, from either side: the @@ -278,7 +297,7 @@ export function buildUploadLeaseHeartbeat( info: DaemonInfo, settings: DaemonClientSettings, request: Omit, -): ((budgetMs: number) => Promise) | undefined { +): ((budgetMs: number, signal: AbortSignal) => Promise) | undefined { if (!isRemoteDaemon(info)) return undefined; const leaseScope = leaseScopeFromRequest(request); if (!leaseScope.leaseId) return undefined; @@ -286,7 +305,7 @@ export function buildUploadLeaseHeartbeat( resolveCommandTimeoutPolicy(INTERNAL_COMMANDS.leaseHeartbeat), { positionals: [] }, ); - return async (budgetMs) => + return async (budgetMs, signal) => await sendRequest( info, buildLeaseHeartbeatRequest(leaseScope, { @@ -300,5 +319,6 @@ export function buildUploadLeaseHeartbeat( // The beat's own budget governs; the command's heartbeat policy only ever caps it, and an // unbounded policy leaves the budget standing on its own. policyTimeoutMs === undefined ? budgetMs : Math.min(policyTimeoutMs, budgetMs), + { signal }, ); } diff --git a/src/daemon-client/daemon-client-transport.ts b/src/daemon-client/daemon-client-transport.ts index 5e7cd1a0b7..17cd72ade2 100644 --- a/src/daemon-client/daemon-client-transport.ts +++ b/src/daemon-client/daemon-client-transport.ts @@ -1,8 +1,7 @@ import type { RequestProgressSink } from '@agent-device/contracts/progress'; import net from 'node:net'; -import { AppError } from '@agent-device/kernel/errors'; +import { AppError, createRequestCanceledError } from '@agent-device/kernel/errors'; import { loadNodeHttpRequester, readNodeHttpResponseBody } from '@agent-device/host-kit/transport'; -import { abortedRequestError, refuseAbortedRequest } from '@agent-device/host-kit/request'; import type { DaemonRequest, DaemonResponse } from '../daemon/daemon-request.ts'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; import type { DaemonPaths, DaemonTransportPreference } from '../daemon-resolution.ts'; @@ -36,6 +35,93 @@ type SendRequestOptions = { signal?: AbortSignal; }; +/** + * The requester-side half of a per-call `AbortSignal` (#3178), kept beside the transport that + * enforces it so cancellation has one owner across both halves. The transport half closes the + * request's connection (which is what makes the daemon mark the request canceled); this half + * answers for the cases a transport cannot observe: + * + * - a signal already aborted before anything is sent — nothing may leave the process, so the refusal + * carries `details.dispatched: 'no'`; + * - an abort during a phase with no connection to close, or a custom transport that ignores the + * signal — the caller's promise still settles, with `details.dispatched: 'unknown'` because the + * request may already have been written. + * + * Both reject with the typed canceled-request error (`details.reason: 'request_canceled'`) whatever + * reason the caller aborted with, so every layer dispatches on the same reason and never on the + * caller's arbitrary abort reason or on error text. An abort is never a timeout: nothing here reads + * or extends a request deadline, and no timeout path may produce this rejection. + */ + +/** + * The typed canceled-request error for a caller's own abort. The reason the caller aborted with + * survives unchanged as the cause — `AbortController.abort()` accepts any value — while the + * rejection itself dispatches on `reason: 'request_canceled'` with the delivery evidence the + * aborting layer can prove. A built-in transport rejects every abort through this, so a caller's + * arbitrary abort reason never escapes as the outcome of a daemon request. + */ +export function abortedRequestError( + signal: AbortSignal, + dispatched: 'no' | 'unknown', + requestId?: string, +): AppError { + return createRequestCanceledError({ requestId, dispatched }, signal.reason); +} + +/** + * Refuses a send attempt that starts while the caller's signal is already aborted: nothing may leave + * the process, so this never touches a connection and the refusal carries `details.dispatched: 'no'`. + */ +export function refuseAbortedRequest(signal: AbortSignal | undefined, requestId?: string): void { + if (!signal?.aborted) return; + throw abortedRequestError(signal, 'no', requestId); +} + +export type RequestGuard = { + /** Refuses an already-aborted call before anything is sent. No-op without a signal. */ + refuseIfAborted(): void; + /** Settles `send`'s outcome against the signal, winning with the typed canceled error on abort. */ + guard(send: () => Promise): Promise; +}; + +const NO_REQUEST_GUARD: RequestGuard = { + refuseIfAborted: () => {}, + guard: async (send: () => Promise) => await send(), +}; + +export function createRequestGuard(params: { + signal: AbortSignal | undefined; + requestId?: string; +}): RequestGuard { + const { signal, requestId } = params; + if (!signal) return NO_REQUEST_GUARD; + return { + refuseIfAborted() { + if (signal.aborted) throw abortedRequestError(signal, 'no', requestId); + }, + guard: async (send: () => Promise): Promise => { + // Checked before `send` runs, so nothing has been dispatched yet — the honest evidence is + // 'no'. An abort listener attached to an already-aborted signal never fires, so this check is + // what keeps a late `guard` call rejecting rather than pending forever. + if (signal.aborted) throw abortedRequestError(signal, 'no', requestId); + return await new Promise((resolve, reject) => { + const onAbort = (): void => reject(abortedRequestError(signal, 'unknown', requestId)); + signal.addEventListener('abort', onAbort, { once: true }); + void send().then( + (value) => { + signal.removeEventListener('abort', onAbort); + resolve(value); + }, + (error: unknown) => { + signal.removeEventListener('abort', onAbort); + reject(error); + }, + ); + }); + }, + }; +} + const LOCAL_DAEMON_HEALTHCHECK_TIMEOUT_MS = 500; const REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS = 3000; const DAEMON_ENDPOINT_UNAVAILABLE_REASON = 'daemon_endpoint_unavailable'; @@ -125,8 +211,9 @@ function canConnectHttp(info: DaemonInfo): Promise { export async function readRemoteDaemonHealth( info: DaemonInfo, probeTimeoutMs?: number, + callerSignal?: AbortSignal, ): Promise { - const health = await readDaemonHttpHealth(info, probeTimeoutMs); + const health = await readDaemonHttpHealth(info, probeTimeoutMs, callerSignal); 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. @@ -152,6 +239,7 @@ export async function readRemoteDaemonHealth( async function readDaemonHttpHealth( info: DaemonInfo, probeTimeoutMs?: number, + callerSignal?: AbortSignal, ): Promise { const endpoint = info.baseUrl ? buildDaemonHttpUrl(info.baseUrl, 'health') @@ -166,11 +254,14 @@ async function readDaemonHttpHealth( probeTimeoutMs ?? Number.POSITIVE_INFINITY, ); if (timeoutMs <= 0) return { reachable: false, timedOut: true }; - const signal = AbortSignal.timeout(Math.ceil(timeoutMs)); + // `timedOut` is keyed on the probe's own budget alone: a caller abort must never read as a + // timeout, because a timed-out probe on an RPC-capped budget is answered as the RPC timing out. + const timeoutSignal = AbortSignal.timeout(Math.ceil(timeoutMs)); + const signal = callerSignal ? AbortSignal.any([timeoutSignal, callerSignal]) : timeoutSignal; return await new Promise((resolve) => { const headers = info.baseUrl ? buildDaemonHttpAuthHeaders(info.token) : {}; const unreachable = (): RemoteDaemonHealth => - signal.aborted ? { reachable: false, timedOut: true } : { reachable: false }; + timeoutSignal.aborted ? { reachable: false, timedOut: true } : { reachable: false }; const req = transport.request( { protocol: url.protocol, @@ -289,7 +380,11 @@ async function retryAfterRemoteInstanceMismatch( timeoutMs, deadline, ); - const health = await readRemoteDaemonHealth(info, probeTimeoutMs); + const health = await readRemoteDaemonHealth(info, probeTimeoutMs, options.signal); + // An abort that landed during the probe is this caller's cancellation, not the probe's outcome: + // it is answered before the timed-out-probe and unreachable-daemon branches, so a canceled call + // never borrows the timeout's shape or a daemon-unavailable error. + refuseAbortedRequest(options.signal, req.meta?.requestId); // The probe's timer starts from the event loop's cached clock, so it can expire while // performance.now() is still short of the deadline: a probe the RPC deadline capped that ran out // of time is the RPC timing out. @@ -493,8 +588,7 @@ async function sendSocketRequest( ); }, timeoutMs) : undefined; - // Destroying the connection is what makes the daemon mark the request canceled; the rejection - // keeps the typed canceled-request error, never a socket-error or timeout shape. The timeout + // Destroying the connection is what makes the daemon mark the request canceled. The timeout // timer is cleared here because an abort is never a timeout: nothing may run the timeout's // runner sweep or daemon reset after the caller canceled. const rejectAborted = (signal: AbortSignal): void => { @@ -699,7 +793,6 @@ async function sendHttpRequest( const timeoutHandle = typeof timeoutMs === 'number' ? setTimeout(() => { - request.destroy(); rejectOnce( handleRequestTimeout({ info, @@ -707,20 +800,22 @@ async function sendHttpRequest( ...timeoutRequestContext(req, remote, timeoutMs), }), ); + request.destroy(); }, timeoutMs) : undefined; // Destroying the request closes this one connection, which is what makes the daemon mark the // request canceled. The timeout timer is cleared because an abort is never a timeout: nothing - // may run the timeout's runner sweep or daemon reset after the caller canceled. The destroy - // surfaces as a request `error`, which `settled` keeps out of the already-typed rejection. + // may run the timeout's runner sweep or daemon reset after the caller canceled. The settle goes + // first so the `error` the destroy surfaces arrives after `settled` and cannot rewrite the + // typed cancellation into a transport failure. const onCallerAbort = (): void => { if (settled) return; if (timeoutHandle) clearTimeout(timeoutHandle); - request.destroy(); rejectOnce( abortedRequestError(callerSignal!, requestEnded ? 'unknown' : 'no', req.meta?.requestId), ); + request.destroy(); }; if (callerSignal) { callerSignal.addEventListener('abort', onCallerAbort, { once: true }); @@ -728,6 +823,11 @@ async function sendHttpRequest( request.on('error', (err) => { if (timeoutHandle) clearTimeout(timeoutHandle); + // `request.destroy()` surfaces an `error` after an intentional abort or timeout has already + // settled this promise. Returning before `handleTransportError` keeps that expected teardown + // from emitting a `daemon_request_socket_error` diagnostic; a genuine transport error still + // arrives before the settle and rejects. + if (settled) return; rejectOnce(handleTransportError(err, req.meta?.requestId, remote)); }); diff --git a/src/daemon-client/daemon-client.ts b/src/daemon-client/daemon-client.ts index 10aab698e2..461c91ac8d 100644 --- a/src/daemon-client/daemon-client.ts +++ b/src/daemon-client/daemon-client.ts @@ -10,7 +10,6 @@ import { emitDiagnostic, withDiagnosticTimer, } from '@agent-device/host-kit/diagnostics'; -import { createRequestGuard } from '@agent-device/host-kit/request'; import { INTERNAL_COMMANDS, PUBLIC_COMMANDS } from '@agent-device/command-registry/catalog'; import { resolveCommandTimeoutPolicy } from '@agent-device/command-registry/registry'; import { resolveCommandRequestTimeoutMs } from '@agent-device/command-registry/timeout-policy'; @@ -29,7 +28,7 @@ import { type DaemonClientSettings, type EnsuredDaemon, } from './daemon-client-lifecycle.ts'; -import { sendRequest } from './daemon-client-transport.ts'; +import { createRequestGuard, sendRequest } from './daemon-client-transport.ts'; import { isRemoteDaemon, type DaemonInfo } from './daemon-client-metadata.ts'; import { leaseScopeFromRequest } from '@agent-device/contracts/lease-scope'; @@ -61,11 +60,8 @@ export async function sendToDaemon( resolveCommandTimeoutPolicy(requestWithoutAuthFlag.command), requestWithoutAuthFlag, ); - // The caller's signal covers every phase of this one request: an already-aborted call is - // refused before the daemon is even started or an artifact byte is uploaded, and an abort that - // arrives mid-phase stops the upload and the request. Each phase's own machinery (the upload - // client, the transports below) honors the signal; the guard turns any of those outcomes into - // the typed canceled-request error, so an abort never borrows a timeout's shape. + // The caller's signal covers every phase of this one request, and the guard turns any phase's + // cancellation into the typed canceled-request error — so an abort never borrows a timeout's shape. const cancellation = createRequestGuard({ signal: options.signal, requestId }); cancellation.refuseIfAborted(); const daemon = await cancellation.guard(async () => { @@ -351,5 +347,6 @@ async function protectArtifactUploadWithLeaseBeats( info, callerSignal ? AbortSignal.any([leaseSignal, callerSignal]) : leaseSignal, ), + callerSignal, }); } diff --git a/test/wire-compat/ledger.json b/test/wire-compat/ledger.json index 80e13dc3e9..34fe4dccc2 100644 --- a/test/wire-compat/ledger.json +++ b/test/wire-compat/ledger.json @@ -67,10 +67,10 @@ "src/daemon-client/daemon-client-rpc.ts#toDaemonHttpRpcError": "sha256:888246763c48670e7da893054d025744654f8715c3b4906312617a2b5028316b", "src/daemon-client/daemon-client-transport.ts#RemoteDaemonHealth": "sha256:3bac36fa97090b273afe1128103e1bf476d05441fd9de41317f1b78775b88304", "src/daemon-client/daemon-client-transport.ts#RemoteDaemonHealthLink": "sha256:7702598468b4c82b4ffa93064cd6bdfec938c1c527f7e6007c0584f3884a6eea", - "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth": "sha256:eef0e153eb6f0bc02175e464dd6d98559b500b65ea8c74ba2203cfef09ec09a0", + "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth": "sha256:a483fbdd845db39ffea596061ff04e9db4acdafd20b593dfed74e7c82a93e522", "src/daemon-client/daemon-client-transport.ts#readHealthLink": "sha256:f56404b94d73da57de7248719c9136b760d2be8b4a53cde5315a2311b088425b", "src/daemon-client/daemon-client-transport.ts#readHealthPayload": "sha256:4e85ffc3e35e02379c393e9312344757e003cf1f0ad9eb8d1f77d90c81c861f1", - "src/daemon-client/daemon-client-transport.ts#readRemoteDaemonHealth": "sha256:bcefa89fbb7fcbd6fee1b5ecb217955b9eee8d6fd2ad175fb653011edbd199d9", + "src/daemon-client/daemon-client-transport.ts#readRemoteDaemonHealth": "sha256:df174d1ccf73851ca58bd75533fb4660424003b98bd9e41fe942cf3ab3698e5f", "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", @@ -276,8 +276,8 @@ }, { "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:df174d1ccf73851ca58bd75533fb4660424003b98bd9e41fe942cf3ab3698e5f", + "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. #3178 threads the caller's optional AbortSignal through so a canceled call cuts a probe already in flight instead of waiting it out: the signal is client-local, never serialized, and the GET /health request and payload are unchanged." }, { "declaration": "src/daemon/server/http-server.ts#authorizeAuxiliaryHttpRequest", diff --git a/website/docs/docs/client-api.md b/website/docs/docs/client-api.md index e2702ad8f6..b1845ce9bc 100644 --- a/website/docs/docs/client-api.md +++ b/website/docs/docs/client-api.md @@ -239,14 +239,14 @@ Results are daemon-shaped objects with typed known fields, so command semantics A failed interaction rejects with the same error the CLI prints. Read `error.details.dispatched` before you retry; [Commands](./commands.md) explains the two values. -Every client call accepts `signal?: AbortSignal` to cancel that one call: +Every client call that dispatches to the daemon accepts `signal?: AbortSignal` to cancel that one call: ```ts const controller = new AbortController(); await client.interactions.press({ ref: '@e12', signal: controller.signal }); ``` -A signal that is already aborted rejects the call without sending anything (`error.details.dispatched: 'no'`). Aborting while the request is in flight closes that request's connection, the daemon marks the request canceled, and the promise rejects with the typed canceled-request error (`error.details.reason: 'request_canceled'`, `error.details.dispatched: 'unknown'`). The daemon and the session stay alive for other requests. An abort is never a timeout: it never triggers the timeout path's runner cleanup or daemon reset. +With the built-in transport, a signal that is already aborted rejects the call without sending anything (`error.details.dispatched: 'no'`). Aborting while the request is in flight closes that request's connection, the daemon marks the request canceled, and the promise rejects with the typed canceled-request error (`error.details.reason: 'request_canceled'`, `error.details.dispatched: 'unknown'`). A custom transport receives the signal on its context and may cancel differently. The daemon and the session stay alive for other requests, and an abort is never a timeout: it never triggers the timeout path's runner cleanup or daemon reset. Cancellation covers the daemon request itself, so a response-artifact download already underway is not stopped, and a canceled one-shot replay still runs the existing cleanup that can tear down a daemon this client started. ```ts await client.command.wait({ From e2c10e1b282466ed947b0f6ccae89285e0484af1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sun, 4 Oct 2026 09:03:27 +0200 Subject: [PATCH 09/14] refactor(daemon-client): keep the guard's helpers module-local and name the beat's outcomes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The code-quality audit found two exports with no consumer outside the transport module and two functions whose new branches pushed them over the complexity threshold. Nothing imported the helpers, so the guard's error factory and its send-attempt refusal stay module-private beside the two callers that use them. The lease beat now says what a failed beat means — one path ends the protection, one reports a transient miss — and the retry path names the budget a deadline-capped probe proves, instead of inlining both decisions. --- src/daemon-client/daemon-client-lease-beat.ts | 56 ++++++++++--------- src/daemon-client/daemon-client-transport.ts | 35 +++++++----- 2 files changed, 53 insertions(+), 38 deletions(-) diff --git a/src/daemon-client/daemon-client-lease-beat.ts b/src/daemon-client/daemon-client-lease-beat.ts index cb489c8010..dcce3011e5 100644 --- a/src/daemon-client/daemon-client-lease-beat.ts +++ b/src/daemon-client/daemon-client-lease-beat.ts @@ -148,6 +148,34 @@ export async function runProtectedLeaseWork( if (timer) clearTimeout(timer); timer = setTimeout(runBeat, delayMs); }; + // A beat that finds the lease gone (or finds this client can never renew it) ends the + // protection: the upload is pointed at a device nobody owns, so it is stopped rather than + // allowed to finish bytes nobody will use. A beat that ends after the phase settled cannot + // change that outcome, but a lease this client just learned is gone is still worth one + // diagnostic on the way out. + const reportTerminalBeatFailure = (error: unknown): void => { + terminalError = error; + if (stopped) { + emitDiagnostic({ + level: 'warn', + phase: 'lease_lost_after_phase', + data: { message: error instanceof Error ? error.message : String(error) }, + }); + return; + } + control.abort(); + reportTerminal?.(error); + }; + const reportTransientBeatFailure = (error: unknown): void => { + // A beat the loop itself canceled — for a caller that has already given up — failed for a + // reason that says nothing about the lease, so it is not worth a warning. + if (stopped) return; + emitDiagnostic({ + level: 'warn', + phase: 'lease_heartbeat_failed', + data: { message: error instanceof Error ? error.message : String(error) }, + }); + }; arm(intervalMs); const settle = (async () => { const budgetMs = windowMs; @@ -157,8 +185,7 @@ export async function runProtectedLeaseWork( ); // An answer that names no window keeps the cadence it was asked at: the loop only ever // moves on evidence of how long the lease is good for, and never on the absence of it. - if (renewed === undefined) return; - if (renewed === windowMs) return; + if (renewed === undefined || renewed === windowMs) return; const cadence = leaseBeatIntervalMs(renewed); windowMs = renewed; intervalMs = cadence; @@ -166,31 +193,10 @@ export async function runProtectedLeaseWork( arm(cadence); } catch (error) { if (isTerminalLeaseBeatError(error)) { - terminalError = error; - if (stopped) { - // The phase settled first; the outcome it returned already stands, but a lease this - // client just learned is gone is worth one diagnostic on the way out. - emitDiagnostic({ - level: 'warn', - phase: 'lease_lost_after_phase', - data: { message: error instanceof Error ? error.message : String(error) }, - }); - return; - } - // The upload is the only thing still consuming this phase's time, and it is pointed at a - // device this client can no longer renew. Stop it rather than finish bytes nobody owns. - control.abort(); - reportTerminal?.(error); + reportTerminalBeatFailure(error); return; } - // A beat the loop itself canceled — for a caller that has already given up — failed for a - // reason that says nothing about the lease, so it is not worth a warning. - if (stopped) return; - emitDiagnostic({ - level: 'warn', - phase: 'lease_heartbeat_failed', - data: { message: error instanceof Error ? error.message : String(error) }, - }); + reportTransientBeatFailure(error); } })(); // A beat the loop has moved on from is still listened to, and nothing awaits it: its outcome is diff --git a/src/daemon-client/daemon-client-transport.ts b/src/daemon-client/daemon-client-transport.ts index 17cd72ade2..e1ac9b1bb4 100644 --- a/src/daemon-client/daemon-client-transport.ts +++ b/src/daemon-client/daemon-client-transport.ts @@ -60,7 +60,7 @@ type SendRequestOptions = { * aborting layer can prove. A built-in transport rejects every abort through this, so a caller's * arbitrary abort reason never escapes as the outcome of a daemon request. */ -export function abortedRequestError( +function abortedRequestError( signal: AbortSignal, dispatched: 'no' | 'unknown', requestId?: string, @@ -72,12 +72,12 @@ export function abortedRequestError( * Refuses a send attempt that starts while the caller's signal is already aborted: nothing may leave * the process, so this never touches a connection and the refusal carries `details.dispatched: 'no'`. */ -export function refuseAbortedRequest(signal: AbortSignal | undefined, requestId?: string): void { +function refuseAbortedRequest(signal: AbortSignal | undefined, requestId?: string): void { if (!signal?.aborted) return; throw abortedRequestError(signal, 'no', requestId); } -export type RequestGuard = { +type RequestGuard = { /** Refuses an already-aborted call before anything is sent. No-op without a signal. */ refuseIfAborted(): void; /** Settles `send`'s outcome against the signal, winning with the typed canceled error on abort. */ @@ -385,19 +385,12 @@ async function retryAfterRemoteInstanceMismatch( // it is answered before the timed-out-probe and unreachable-daemon branches, so a canceled call // never borrows the timeout's shape or a daemon-unavailable error. refuseAbortedRequest(options.signal, req.meta?.requestId); - // The probe's timer starts from the event loop's cached clock, so it can expire while - // performance.now() is still short of the deadline: a probe the RPC deadline capped that ran out - // of time is the RPC timing out. - if ( - health.timedOut && - timeoutMs !== undefined && - probeTimeoutMs !== undefined && - probeTimeoutMs <= REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS - ) { + const timedOutRpcBudgetMs = deadlineCappedProbeTimeoutBudgetMs(health, timeoutMs, probeTimeoutMs); + if (timedOutRpcBudgetMs !== undefined) { throw handleRequestTimeout({ info, statePaths, - ...timeoutRequestContext(req, true, timeoutMs), + ...timeoutRequestContext(req, true, timedOutRpcBudgetMs), }); } const remainingMs = remainingRemoteRequestTimeoutMs(info, req, statePaths, timeoutMs, deadline); @@ -417,6 +410,22 @@ async function retryAfterRemoteInstanceMismatch( } } +/** + * The RPC budget a probe that ran out of time proves, or `undefined` when the probe's timeout was + * its own. The probe's timer starts from the event loop's cached clock, so it can expire while + * `performance.now()` is still short of the deadline; a probe capped by the RPC deadline is + * therefore the deadline speaking, not a slow daemon. + */ +function deadlineCappedProbeTimeoutBudgetMs( + health: RemoteDaemonHealth, + timeoutMs: number | undefined, + probeTimeoutMs: number | undefined, +): number | undefined { + if (!health.timedOut || timeoutMs === undefined) return undefined; + if (probeTimeoutMs === undefined) return undefined; + return probeTimeoutMs <= REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS ? timeoutMs : undefined; +} + function remainingRemoteRequestTimeoutMs( info: DaemonInfo, req: DaemonRequest, From e8a34278375ef3704af796b6251b788b86f565d8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sun, 4 Oct 2026 09:24:37 +0200 Subject: [PATCH 10/14] test(provider-scenarios): resolve alert interactors per call from the request's execution metadata MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The scenario gateway built its Apple interactor once at bind time with an empty runner context, so an alert leg asked the runner for a request id it never had. While the client sent no meta.requestId the unkeyed provider scope still matched, but every call now carries one: the request-scoped scripted provider is keyed by (deviceId, requestId), the empty lookup missed it, and the leg fell through to the local XCTest runner — which spawns real tooling and hangs the in-process harness. Per-call resolution mirrors the production alert binding, which projects the input's execution metadata into the interactor context. --- .../ios-alert-settings.test.ts | 38 ++++++++++++++++--- 1 file changed, 33 insertions(+), 5 deletions(-) diff --git a/test/integration/provider-scenarios/ios-alert-settings.test.ts b/test/integration/provider-scenarios/ios-alert-settings.test.ts index 847635c179..b54250be80 100644 --- a/test/integration/provider-scenarios/ios-alert-settings.test.ts +++ b/test/integration/provider-scenarios/ios-alert-settings.test.ts @@ -283,11 +283,39 @@ function createRecordingPlatformRuntimeGateway(params: { ), // R59 does the same for `alert`: the scenario's gateway states and serves the four // legs, reusing the Apple family's own module so the runner transcript this scenario - // scripts — including its retry and poll windows — is what actually runs. - readAlert: async (input) => await interactor.readAlert!(alertOptions(input)), - awaitAlert: async (input) => await interactor.awaitAlert!(alertOptions(input)), - acceptAlert: async (input) => await interactor.acceptAlert!(alertOptions(input)), - dismissAlert: async (input) => await interactor.dismissAlert!(alertOptions(input)), + // scripts — including its retry and poll windows — is what actually runs. Each leg + // resolves its interactor per call from the input's execution metadata, exactly like + // the production alert binding: the request-scoped runner provider is keyed by the + // request id that metadata carries, so a bind-time interactor would fall back to the + // local runner instead of reaching the scripted provider. + readAlert: async (input) => + await ( + await applePlugin.createInteractor(device, { + ...input.execution, + appBundleId: input.appBundleId, + }) + ).readAlert!(alertOptions(input)), + awaitAlert: async (input) => + await ( + await applePlugin.createInteractor(device, { + ...input.execution, + appBundleId: input.appBundleId, + }) + ).awaitAlert!(alertOptions(input)), + acceptAlert: async (input) => + await ( + await applePlugin.createInteractor(device, { + ...input.execution, + appBundleId: input.appBundleId, + }) + ).acceptAlert!(alertOptions(input)), + dismissAlert: async (input) => + await ( + await applePlugin.createInteractor(device, { + ...input.execution, + appBundleId: input.appBundleId, + }) + ).dismissAlert!(alertOptions(input)), appLogReattach: async () => ({ status: 'missing' }), appLogCleanup: async () => ({ status: 'already-missing' }), resolveOpenTarget: async (input) => ({ From e079b2b92fbef6178285b27546a8b62bdbcd82a6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sun, 4 Oct 2026 09:26:41 +0200 Subject: [PATCH 11/14] refactor(test): one per-call alert-leg resolver for the iOS scenario gateway --- .../ios-alert-settings.test.ts | 52 +++++++------------ 1 file changed, 19 insertions(+), 33 deletions(-) diff --git a/test/integration/provider-scenarios/ios-alert-settings.test.ts b/test/integration/provider-scenarios/ios-alert-settings.test.ts index b54250be80..52bc61fb5b 100644 --- a/test/integration/provider-scenarios/ios-alert-settings.test.ts +++ b/test/integration/provider-scenarios/ios-alert-settings.test.ts @@ -200,6 +200,20 @@ function createRecordingPlatformRuntimeGateway(params: { } const appleOs = device.appleOs; const interactor = await applePlugin.createInteractor(device, {}); + // Each alert leg resolves its interactor per call from the input's execution metadata, + // exactly like the production alert binding. The request-scoped runner provider is keyed by + // the request id that metadata carries, so a bind-time interactor with an empty runner + // context would miss the scripted provider and fall through to the local XCTest runner. + const runAlertLeg = async ( + leg: 'readAlert' | 'awaitAlert' | 'acceptAlert' | 'dismissAlert', + input: AlertRuntimeInput, + ): Promise> => { + const alertInteractor = await applePlugin.createInteractor(device, { + ...input.execution, + appBundleId: input.appBundleId, + }); + return await alertInteractor[leg]!(alertOptions(input)); + }; return { device, owner, @@ -283,39 +297,11 @@ function createRecordingPlatformRuntimeGateway(params: { ), // R59 does the same for `alert`: the scenario's gateway states and serves the four // legs, reusing the Apple family's own module so the runner transcript this scenario - // scripts — including its retry and poll windows — is what actually runs. Each leg - // resolves its interactor per call from the input's execution metadata, exactly like - // the production alert binding: the request-scoped runner provider is keyed by the - // request id that metadata carries, so a bind-time interactor would fall back to the - // local runner instead of reaching the scripted provider. - readAlert: async (input) => - await ( - await applePlugin.createInteractor(device, { - ...input.execution, - appBundleId: input.appBundleId, - }) - ).readAlert!(alertOptions(input)), - awaitAlert: async (input) => - await ( - await applePlugin.createInteractor(device, { - ...input.execution, - appBundleId: input.appBundleId, - }) - ).awaitAlert!(alertOptions(input)), - acceptAlert: async (input) => - await ( - await applePlugin.createInteractor(device, { - ...input.execution, - appBundleId: input.appBundleId, - }) - ).acceptAlert!(alertOptions(input)), - dismissAlert: async (input) => - await ( - await applePlugin.createInteractor(device, { - ...input.execution, - appBundleId: input.appBundleId, - }) - ).dismissAlert!(alertOptions(input)), + // scripts — including its retry and poll windows — is what actually runs. + readAlert: async (input) => await runAlertLeg('readAlert', input), + awaitAlert: async (input) => await runAlertLeg('awaitAlert', input), + acceptAlert: async (input) => await runAlertLeg('acceptAlert', input), + dismissAlert: async (input) => await runAlertLeg('dismissAlert', input), appLogReattach: async () => ({ status: 'missing' }), appLogCleanup: async () => ({ status: 'already-missing' }), resolveOpenTarget: async (input) => ({ From 942bfab4a6aface4321d931e366c6eb3ae9adb2f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sun, 4 Oct 2026 09:32:27 +0200 Subject: [PATCH 12/14] test(daemon-client): destroy tracked sockets when a loopback abort case fails net.Server.close() waits for accepted connections, and a hanging-handler request keeps its connection open until the client's abort destroys it. When an assertion before that fails, the daemon-side handler waits forever and close() hangs the lane, turning the regression signal the abort tests exist to provide into a test timeout. Track the accepted sockets and destroy them in teardown at the three socket suites whose handlers never answer. http.Server already has closeAllConnections(); net.Server has no equivalent. --- src/__tests__/test-utils/loopback.ts | 23 +++++++++++++++++++ ...mon-client-abort-default-transport.test.ts | 18 +++++++++++---- .../daemon-client-abort-timeout.test.ts | 5 ++++ .../__tests__/daemon-client-abort.test.ts | 6 +++++ 4 files changed, 47 insertions(+), 5 deletions(-) diff --git a/src/__tests__/test-utils/loopback.ts b/src/__tests__/test-utils/loopback.ts index 323d2eb638..0f8d1cf5aa 100644 --- a/src/__tests__/test-utils/loopback.ts +++ b/src/__tests__/test-utils/loopback.ts @@ -79,6 +79,29 @@ function closeHttpConnections(server: LoopbackServer): void { maybeHttpServer.closeIdleConnections?.(); } +/** + * `net.Server.close()` waits for every accepted connection to finish, and a request handler that + * never answers keeps its connection open forever — `http.Server` has `closeAllConnections()` and + * `net.Server` has no equivalent. Track the accepted sockets so a test whose assertion already + * failed can destroy them instead of hanging in teardown: a regression in cancellation wiring must + * surface as the failed assertion, not as the test timeout swallowing it. + */ +export function trackLoopbackSockets(server: net.Server): () => void { + const sockets = new Set(); + server.on('connection', (socket) => { + sockets.add(socket); + socket.on('close', () => { + sockets.delete(socket); + }); + }); + return () => { + for (const socket of sockets) { + socket.destroy(); + } + sockets.clear(); + }; +} + export function waitForHttpOk(url: string, timeoutMs: number): Promise { const deadline = Date.now() + timeoutMs; return new Promise((resolve, reject) => { diff --git a/src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts b/src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts index f2d18b49c4..dc2fbf4bee 100644 --- a/src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts +++ b/src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts @@ -31,6 +31,7 @@ import { closeLoopbackServer, listenOnLoopback, skipWhenLoopbackUnavailable, + trackLoopbackSockets, type SkippableTestContext, } from '../../__tests__/test-utils/loopback.ts'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; @@ -91,14 +92,16 @@ async function waitFor(condition: () => boolean, what: string): Promise { async function runDefaultTransportAbort( t: SkippableTestContext, transport: 'socket' | 'http', - serve: ( - handleRequest: DaemonInvokeFn, - ) => Promise<{ server: Parameters[0]; port: number }>, + serve: (handleRequest: DaemonInvokeFn) => Promise<{ + server: Parameters[0]; + port: number; + destroyConnections?: () => void; + }>, ): Promise { if (await skipWhenLoopbackUnavailable(t)) return; const seen: SeenRequest = { started: false, canceled: [] }; const stateDir = mkdtempForTestSync(`agent-device-abort-default-${transport}-`); - const { server, port } = await serve(hangingWaitHandler(seen)); + const { server, port, destroyConnections } = await serve(hangingWaitHandler(seen)); try { publishLoopbackDaemonInfo( stateDir, @@ -134,6 +137,7 @@ async function runDefaultTransportAbort( assert.equal(typeof details.requestId, 'string'); assert.equal(details.requestId, seen.requestId); } finally { + destroyConnections?.(); await closeLoopbackServer(server); fs.rmSync(stateDir, { recursive: true, force: true }); } @@ -142,7 +146,11 @@ async function runDefaultTransportAbort( test('the default socket transport carries the caller signal: the daemon cancels the request the client aborted', async (t) => { await runDefaultTransportAbort(t, 'socket', async (handleRequest) => { const server = createSocketServer(handleRequest); - return { server, port: await listenNetServer(server) }; + // A canceled request's connection is destroyed by the client, but if an assertion above fails + // before that, the handler waits forever on an open connection and `net.Server.close()` would + // never return. Destroying the tracked sockets turns that into the failed assertion it is. + const destroyConnections = trackLoopbackSockets(server); + return { server, port: await listenNetServer(server), destroyConnections }; }); }); 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 31cae735e6..d88325338c 100644 --- a/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts +++ b/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts @@ -21,6 +21,7 @@ import { resolveDaemonPaths } from '../../daemon-resolution.ts'; import { closeLoopbackServer, skipWhenLoopbackUnavailable, + trackLoopbackSockets, type SkippableTestContext, } from '../../__tests__/test-utils/loopback.ts'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; @@ -62,6 +63,9 @@ async function runTimerDiscipline( const aborted = outcome === 'aborted'; handleRequestTimeoutCalls.length = 0; const server = createSocketServer(hangingHandler()); + // The hanging handler never answers, so whichever way the request ends, a failed assertion + // before that would leave the connection open and `net.Server.close()` would hang the lane. + const destroySockets = trackLoopbackSockets(server); try { const port = await listenNetServer(server); const controller = new AbortController(); @@ -95,6 +99,7 @@ async function runTimerDiscipline( await new Promise((resolve) => setTimeout(resolve, OUTLIVE_MS)); if (aborted) assert.deepEqual(handleRequestTimeoutCalls, []); } finally { + destroySockets(); await closeLoopbackServer(server); } } diff --git a/src/daemon-client/__tests__/daemon-client-abort.test.ts b/src/daemon-client/__tests__/daemon-client-abort.test.ts index 2837d82c09..37e3805b34 100644 --- a/src/daemon-client/__tests__/daemon-client-abort.test.ts +++ b/src/daemon-client/__tests__/daemon-client-abort.test.ts @@ -20,6 +20,7 @@ import { resolveDaemonPaths } from '../../daemon-resolution.ts'; import { closeLoopbackServer, listenOnLoopback, + trackLoopbackSockets, skipWhenLoopbackUnavailable, } from '../../__tests__/test-utils/loopback.ts'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; @@ -110,6 +111,10 @@ test('socket transport: an abort mid-request closes the connection, the daemon m if (await skipWhenLoopbackUnavailable(t)) return; const seen: SeenRequest = { started: false, canceled: [] }; const server = createSocketServer(canceledAwareHandler(seen)); + // The mid-request case's connection only closes through the client's abort; if an assertion + // above that fails, the daemon-side handler waits forever and `net.Server.close()` would hang + // the whole lane instead of reporting the failure. + const destroySockets = trackLoopbackSockets(server); try { const port = await listenNetServer(server); const info = { port, token: TOKEN, pid: 1 }; @@ -159,6 +164,7 @@ test('socket transport: an abort mid-request closes the connection, the daemon m assert.equal(followUp.ok, true); assert.equal(seen.servedCommand, 'devices'); } finally { + destroySockets(); await closeLoopbackServer(server); } }); From b431157c59d09722191e963f271aa742cc4f710c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sun, 4 Oct 2026 09:32:27 +0200 Subject: [PATCH 13/14] chore(gates): ack the caller-signal digest on the HTTP health probe The compatibleChanges entry must carry the post-change digest; the declaration moved when the probe began composing the #3178 caller signal with its own timeout budget. The /health request and accepted payload are untouched, so protocol 2 peers parse it exactly as before. --- test/wire-compat/ledger.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/test/wire-compat/ledger.json b/test/wire-compat/ledger.json index 34fe4dccc2..c6cf7cbb49 100644 --- a/test/wire-compat/ledger.json +++ b/test/wire-compat/ledger.json @@ -431,8 +431,8 @@ }, { "declaration": "src/daemon-client/daemon-client-transport.ts#readDaemonHttpHealth", - "digest": "sha256:eef0e153eb6f0bc02175e464dd6d98559b500b65ea8c74ba2203cfef09ec09a0", - "rationale": "A restart retry must tell a probe that ran out of time from one that failed, so the client can report the RPC deadline instead of 'Remote daemon is unavailable'. `timedOut` is client-local: the prober sets it on its own timeout and abort paths and never reads it from a /health payload. The health request and the accepted payload fields are unchanged, so a released daemon or proxy is probed and parsed exactly as before." + "digest": "sha256:a483fbdd845db39ffea596061ff04e9db4acdafd20b593dfed74e7c82a93e522", + "rationale": "A restart retry must tell a probe that ran out of time from one that failed, so the client can report the RPC deadline instead of 'Remote daemon is unavailable'. `timedOut` is client-local: the prober sets it on its own timeout and abort paths and never reads it from a /health payload. The #3178 caller signal joins the probe's own timeout in-process (AbortSignal.any) and never enters the /health request or the accepted payload fields, so a released daemon or proxy is probed and parsed exactly as before." } ] } From bb8e19de1bbd7c1ee9d5dbcca913aacb4378ff87 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sun, 4 Oct 2026 10:15:58 +0200 Subject: [PATCH 14/14] docs(agents): name both request-cancellation seams The caller's per-call abort guard moved beside the transports it serves, while `@agent-device/host-kit/request` kept the daemon-side cancel registry and progress. The declaration-site line still credited the whole seam to host-kit, pointing the next reader at a module that no longer exports the guard. --- AGENTS.md | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 00ac5f7f8d..f01ce73b7b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -56,9 +56,11 @@ Read the declaration rather than maintaining a prose copy: - common command input fields, and which surface may write an input key (model, operator, retired): `src/commands/common-input-fields.ts` and `src/commands/input-audience.ts` -Shared selector parsing and matching belongs in `@agent-device/selectors`; request cancellation -and progress in `@agent-device/host-kit/request`; cross-layer contracts in `packages/contracts/src`; -CLI flags in `src/commands/cli-grammar`; cross-surface schema composition in `src/commands/schema`. +Shared selector parsing and matching belongs in `@agent-device/selectors`; daemon-side request +cancellation and progress in `@agent-device/host-kit/request`, and the caller's per-call abort guard +beside the transports it serves in `src/daemon-client/daemon-client-transport.ts`; cross-layer +contracts in `packages/contracts/src`; CLI flags in `src/commands/cli-grammar`; cross-surface schema +composition in `src/commands/schema`. Resolve registry completeness failures at the missing declaration. Diagnose other gate failures at their reported invariant; do not suppress them or add an allowlist to get a pass. Build interaction