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 diff --git a/packages/contracts/src/client-connection.ts b/packages/contracts/src/client-connection.ts index 7568768a33..52d6bc78ce 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,22 @@ 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. + * + * 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; +}; 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/src/__tests__/client-abort-signal.test.ts b/src/__tests__/client-abort-signal.test.ts new file mode 100644 index 0000000000..33af448513 --- /dev/null +++ b/src/__tests__/client-abort-signal.test.ts @@ -0,0 +1,103 @@ +/** + * #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 () => { + 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. 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?.(); + }), + }, + ); + 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')); + if (lingering) clearTimeout(lingering); +}); + +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 }); + // 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) => canceledWith(error, 'unknown')); + assert.deepEqual(await survivor, {}); +}); 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/agent-device-client.ts b/src/agent-device-client.ts index 104d169b48..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,6 +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 './daemon-client/daemon-client-transport.ts'; export function createAgentDeviceClient( config: AgentDeviceClientConfig = {}, @@ -95,6 +97,12 @@ export function createAgentDeviceClient( input?: Record, ): Promise> => { const merged = mergeClientOptions(config, options); + // 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), command, @@ -102,9 +110,16 @@ export function createAgentDeviceClient( ...(input ? { input } : {}), flags: buildRequestFlags(merged, metadataFlags), runtime: merged.runtime, - meta: buildMeta(merged), + meta: { ...buildMeta(merged), requestId }, }; - 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/__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..dc2fbf4bee --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-abort-default-transport.test.ts @@ -0,0 +1,162 @@ +/** + * #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, + trackLoopbackSockets, + 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; + 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, destroyConnections } = 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 { + destroyConnections?.(); + 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); + // 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 }; + }); +}); + +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-timeout.test.ts b/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts new file mode 100644 index 0000000000..d88325338c --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-abort-timeout.test.ts @@ -0,0 +1,113 @@ +/** + * #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, + trackLoopbackSockets, + 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()); + // 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(); + 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 { + destroySockets(); + 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'); +}); 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..37e3805b34 --- /dev/null +++ b/src/daemon-client/__tests__/daemon-client-abort.test.ts @@ -0,0 +1,273 @@ +/** + * #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, + trackLoopbackSockets, + 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[]; + /** 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. +function canceledAwareHandler(seen: SeenRequest): DaemonInvokeFn { + return async (req: DaemonRequest): Promise => { + const requestId = req.meta?.requestId; + seen.requestId = requestId; + if (req.command !== 'wait') { + seen.servedCommand = req.command; + return { ok: true, data: {} }; + } + seen.started = true; + const signal = getRequestSignal(requestId); + if (!signal) { + return { ok: true, data: { answeredWithoutSignal: true } }; + } + // The long request never answers: only the client's disconnect ends it, through the abort. + 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 }, + { + token: TOKEN, + 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)); + // 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 }; + const controller = new AbortController(); + const inFlight = sendRequest( + info, + { + token: TOKEN, + 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, + { + token: TOKEN, + command: 'devices', + session: 'default', + positionals: [], + flags: {}, + meta: { requestId: 'req-socket-follow-up' }, + }, + 'socket', + STATE_PATHS, + 5000, + {}, + ); + assert.equal(followUp.ok, true); + assert.equal(seen.servedCommand, 'devices'); + } finally { + destroySockets(); + 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 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); + } +}); + +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)); + } +} 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..dcce3011e5 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 @@ -134,15 +148,44 @@ 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; 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; - if (renewed === windowMs) return; + if (renewed === undefined || renewed === windowMs) return; const cadence = leaseBeatIntervalMs(renewed); windowMs = renewed; intervalMs = cadence; @@ -150,28 +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; } - 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 @@ -190,8 +215,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 +303,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 +311,7 @@ export function buildUploadLeaseHeartbeat( resolveCommandTimeoutPolicy(INTERNAL_COMMANDS.leaseHeartbeat), { positionals: [] }, ); - return async (budgetMs) => + return async (budgetMs, signal) => await sendRequest( info, buildLeaseHeartbeatRequest(leaseScope, { @@ -300,5 +325,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 d09d5af386..e1ac9b1bb4 100644 --- a/src/daemon-client/daemon-client-transport.ts +++ b/src/daemon-client/daemon-client-transport.ts @@ -1,6 +1,6 @@ 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 type { DaemonRequest, DaemonResponse } from '../daemon/daemon-request.ts'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; @@ -26,8 +26,102 @@ 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; }; +/** + * 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. + */ +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'`. + */ +function refuseAbortedRequest(signal: AbortSignal | undefined, requestId?: string): void { + if (!signal?.aborted) return; + throw abortedRequestError(signal, 'no', requestId); +} + +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'; @@ -117,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. @@ -144,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') @@ -158,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, @@ -240,6 +339,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) { @@ -278,20 +380,17 @@ async function retryAfterRemoteInstanceMismatch( timeoutMs, deadline, ); - const health = await readRemoteDaemonHealth(info, probeTimeoutMs); - // The probe's timer starts from the event loop's cached clock, so it can expire while - // performance.now() is still short of the deadline: a probe the RPC deadline capped that ran out - // of time is the RPC timing out. - if ( - health.timedOut && - timeoutMs !== undefined && - probeTimeoutMs !== undefined && - probeTimeoutMs <= REMOTE_DAEMON_HEALTHCHECK_TIMEOUT_MS - ) { + 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); + 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); @@ -311,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, @@ -452,17 +567,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 +597,26 @@ async function sendSocketRequest( ); }, timeoutMs) : undefined; + // 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 => { + 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 +627,12 @@ async function sendSocketRequest( }, resolve: (response) => { settled = true; + detachCallerAbort(); resolve(response); }, reject: (error) => { settled = true; + detachCallerAbort(); reject(error); }, }); @@ -495,6 +641,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 +710,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 +744,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 +756,7 @@ async function sendHttpRequest( readDaemonHttpProgressResponse(res, { req, onProgress: options.onProgress, - reject, + reject: rejectOnce, clearTimeout: () => { if (timeoutHandle) clearTimeout(timeoutHandle); }, @@ -599,8 +765,8 @@ async function sendHttpRequest( info, req, stateDir: statePaths.baseDir, - resolve, - reject, + resolve: resolveOnce, + reject: rejectOnce, }); }, }); @@ -621,13 +787,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); }); }, ); @@ -636,24 +802,47 @@ async function sendHttpRequest( const timeoutHandle = typeof timeoutMs === 'number' ? setTimeout(() => { - request.destroy(); - reject( + rejectOnce( handleRequestTimeout({ info, statePaths, ...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 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); + rejectOnce( + abortedRequestError(callerSignal!, requestEnded ? 'unknown' : 'no', req.meta?.requestId), + ); + request.destroy(); + }; + if (callerSignal) { + callerSignal.addEventListener('abort', onCallerAbort, { once: true }); + } + request.on('error', (err) => { if (timeoutHandle) clearTimeout(timeoutHandle); - reject(handleTransportError(err, req.meta?.requestId, remote)); + // `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)); }); 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..461c91ac8d 100644 --- a/src/daemon-client/daemon-client.ts +++ b/src/daemon-client/daemon-client.ts @@ -28,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'; @@ -60,16 +60,25 @@ 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, 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 () => { + 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 +112,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 +325,28 @@ 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, + ), + callerSignal, }); } diff --git a/test/integration/provider-scenarios/ios-alert-settings.test.ts b/test/integration/provider-scenarios/ios-alert-settings.test.ts index 847635c179..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, @@ -284,10 +298,10 @@ 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)), + 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) => ({ diff --git a/test/wire-compat/ledger.json b/test/wire-compat/ledger.json index 80e13dc3e9..c6cf7cbb49 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", @@ -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." } ] } diff --git a/website/docs/docs/client-api.md b/website/docs/docs/client-api.md index d59051db73..b1845ce9bc 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 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 }); +``` + +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({ text: 'Continue',