diff --git a/apps/rush-cli-client/src/clientCancellation.ts b/apps/rush-cli-client/src/clientCancellation.ts new file mode 100644 index 0000000000..4c60a964eb --- /dev/null +++ b/apps/rush-cli-client/src/clientCancellation.ts @@ -0,0 +1,40 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. Licensed under the MIT license. +// See LICENSE in the project root for license information. + +import * as os from 'node:os'; + +import type { DaemonClientOutcome } from '@rushstack/rush-client-core'; + +/** Signals that cancel a daemon-routed command. */ +export const CANCELLATION_SIGNALS: ReadonlyArray = ['SIGINT', 'SIGTERM', 'SIGHUP']; + +const SIGNAL_EXIT_CODE_BASE: number = 128; + +/** + * Returns the conventional shell exit code for a process terminated by `signal` (128 + signal number), + * e.g. 130 for SIGINT and 143 for SIGTERM. + */ +export function getSignalExitCode(signal: NodeJS.Signals): number { + return SIGNAL_EXIT_CODE_BASE + (os.constants.signals[signal] ?? os.constants.signals.SIGINT); +} + +/** Formats the notice printed when a daemon-routed command is cancelled. */ +export function formatCancellationMessage(commandName: string): string { + return `rush-client: ${commandName} cancelled.\n`; +} + +/** + * Returns whether a daemon outcome represents a cancelled command. A result is cancelled when the daemon reports it + * as aborted, even if an operation failure determines its semantic outcome. A completed (non-aborted) result wins + * over a late signal, and a rejection is never reported as a cancellation. + */ +export function isCancelledOutcome(outcome: DaemonClientOutcome, signalled: boolean): boolean { + switch (outcome.kind) { + case 'result': + return outcome.result.aborted; + case 'rejected': + return false; + default: + return signalled; + } +} diff --git a/apps/rush-cli-client/src/launchClient.ts b/apps/rush-cli-client/src/launchClient.ts index 1145164c4f..6ae500af55 100644 --- a/apps/rush-cli-client/src/launchClient.ts +++ b/apps/rush-cli-client/src/launchClient.ts @@ -26,6 +26,12 @@ import { executeDaemonCommandAsync } from './daemonCommands'; import { formatAdmissionFailure, getConfiguredAdmission } from './ClientAdmissionControls'; import { ClientOperationRenderer } from './ClientOperationRenderer'; import type { AgentProgressRenderer } from './AgentProgressRenderer'; +import { + CANCELLATION_SIGNALS, + formatCancellationMessage, + getSignalExitCode, + isCancelledOutcome +} from './clientCancellation'; import { getDaemonConnectionOptionsAsync } from './daemonConnectionOptions'; import { readUseRushReporter } from './outputSelection'; import { selectClientRoute, type IClientRoute } from './routing'; @@ -148,9 +154,13 @@ export async function launchClientAsync( return; } const abort: AbortController = new AbortController(); - const onSignal = (): void => abort.abort(); - process.on('SIGINT', onSignal); - process.on('SIGTERM', onSignal); + let cancellationSignal: NodeJS.Signals | undefined; + // Windows test harnesses emit signals without a name; treat those as Ctrl+C. + const onSignal = (signal?: NodeJS.Signals): void => { + cancellationSignal ??= signal ?? 'SIGINT'; + abort.abort(); + }; + for (const signal of CANCELLATION_SIGNALS) process.on(signal, onSignal); const renderer: ClientOperationRenderer = new ClientOperationRenderer({ requestId: request.requestId, colorLevel: terminal.supportsColor ? 1 : 0, @@ -166,7 +176,7 @@ export async function launchClientAsync( writeAsync: (bytes, stream) => writeStreamAsync(stream === 'stderr' ? process.stderr : process.stdout, bytes) }); - let outcome: DaemonClientOutcome; + let outcome: DaemonClientOutcome | undefined; const discoveryLines: string[] = []; const writeDiscoveryAsync = async (): Promise => { if (discoveryLines.length > 0) { @@ -215,16 +225,27 @@ export async function launchClientAsync( } : undefined }); + } catch (error) { + // After cancellation, a transport failure (e.g. the cancellation deadline) still means "cancelled". + if (!abort.signal.aborted || !(error instanceof DaemonClientError)) throw error; + outcome = undefined; } finally { - process.removeListener('SIGINT', onSignal); - process.removeListener('SIGTERM', onSignal); + for (const signal of CANCELLATION_SIGNALS) process.removeListener(signal, onSignal); try { await renderer.closeAsync(); } finally { await client.closeAsync(); } } - if (outcome.kind === 'result') { + if (outcome === undefined || isCancelledOutcome(outcome, abort.signal.aborted)) { + const exitCode: number = getSignalExitCode(cancellationSignal ?? 'SIGINT'); + agentRenderer?.finish({ exitCode, errorMessage: 'cancelled' }); + process.exitCode = exitCode; + // After SIGHUP the terminal may be gone; the exit code is what matters. + await writeStreamAsync(process.stderr, Buffer.from(formatCancellationMessage(route.commandName))).catch( + () => undefined + ); + } else if (outcome.kind === 'result') { agentRenderer?.finish(outcome.result); process.exitCode = outcome.result.exitCode; const diagnostic: string | undefined = getResultDiagnostic(outcome.result); @@ -239,9 +260,6 @@ export async function launchClientAsync( } else if (outcome.kind === 'rejected') { agentRenderer?.finish({ exitCode: 1, errorMessage: `daemon rejected the request (${outcome.rejection.code})` }); throw new Error(`Daemon rejected the request (${outcome.rejection.code}): ${outcome.rejection.message}`); - } else if (abort.signal.aborted) { - agentRenderer?.finish({ exitCode: 130, errorMessage: 'cancelled' }); - process.exitCode = 130; } else { agentRenderer?.dispose(); process.stderr.write(`rush-client: ${outcome.message ?? outcome.reason}; using in-process Rush.\n`); diff --git a/apps/rush-cli-client/src/test/launchClient.test.ts b/apps/rush-cli-client/src/test/launchClient.test.ts index d454c99c4e..39b091f31f 100644 --- a/apps/rush-cli-client/src/test/launchClient.test.ts +++ b/apps/rush-cli-client/src/test/launchClient.test.ts @@ -11,6 +11,7 @@ import { setTimeout as delayAsync } from 'node:timers/promises'; import { Rush } from '@microsoft/rush-lib'; import { DaemonClient, connectOrStartDaemonAsync, getDaemonLogFilePath } from '@rushstack/rush-client-core'; +import type { DaemonClientOutcome } from '@rushstack/rush-client-core'; import { RushDaemonHost, WorkspaceSession } from '@rushstack/rush-daemon'; import { removeTestFolderAsync, @@ -20,6 +21,12 @@ import { captureTestDaemonListenerAsync } from '@rushstack/rush-daemon/lib/test/ import { readDaemonLockfile } from '@rushstack/rush-daemon-transport'; import { getDaemonConnectionOptions } from '../daemonConnectionOptions'; +import { + CANCELLATION_SIGNALS, + formatCancellationMessage, + getSignalExitCode, + isCancelledOutcome +} from '../clientCancellation'; interface IInvocationResult { readonly code: number | undefined; @@ -493,3 +500,47 @@ describe('standalone rushx fallback', () => { expect(result.stderr).toContain('--no-daemon cannot be combined with daemon start'); }); }); + +describe('daemon client cancellation exit codes', () => { + const abortedResult: DaemonClientOutcome = { + kind: 'result', + result: { requestId: 'r', outcome: 'aborted', exitCode: 1, aborted: true } + }; + + it('maps cancellation signals to 128 + signal number', () => { + expect(getSignalExitCode('SIGINT')).toBe(130); + expect(getSignalExitCode('SIGTERM')).toBe(143); + expect(getSignalExitCode('SIGHUP')).toBe(129); + expect(CANCELLATION_SIGNALS).toEqual(['SIGINT', 'SIGTERM', 'SIGHUP']); + }); + + it('treats an aborted daemon result as cancelled instead of copying its exit code', () => { + expect(isCancelledOutcome(abortedResult, true)).toBe(true); + // Ctrl+C read from a raw-mode TTY cancels without a process signal. + expect(isCancelledOutcome(abortedResult, false)).toBe(true); + expect(formatCancellationMessage('build')).toBe('rush-client: build cancelled.\n'); + }); + + it('treats a cancelled result as cancelled even when a failure outcome takes precedence', () => { + const cancelledWithFailure: DaemonClientOutcome = { + kind: 'result', + result: { requestId: 'r', outcome: 'failure', exitCode: 1, aborted: true } + }; + expect(isCancelledOutcome(cancelledWithFailure, true)).toBe(true); + }); + + it('keeps completed results and rejections when a signal arrives late', () => { + const succeeded: DaemonClientOutcome = { + kind: 'result', + result: { requestId: 'r', outcome: 'success', exitCode: 0, aborted: false } + }; + const rejected: DaemonClientOutcome = { + kind: 'rejected', + rejection: { requestId: 'r', code: 'unsupportedProtocolVersion', message: 'no' } + } as unknown as DaemonClientOutcome; + expect(isCancelledOutcome(succeeded, true)).toBe(false); + expect(isCancelledOutcome(rejected, true)).toBe(false); + expect(isCancelledOutcome({ kind: 'fallback', reason: 'unsupported' }, true)).toBe(true); + expect(isCancelledOutcome({ kind: 'fallback', reason: 'unsupported' }, false)).toBe(false); + }); +}); diff --git a/apps/rush-cli-client/src/test/persistentIpcCancellation.test.ts b/apps/rush-cli-client/src/test/persistentIpcCancellation.test.ts index ddc6c5eb5d..883d2f83ab 100644 --- a/apps/rush-cli-client/src/test/persistentIpcCancellation.test.ts +++ b/apps/rush-cli-client/src/test/persistentIpcCancellation.test.ts @@ -55,8 +55,9 @@ describe('public-client cancellation of an admitted Node operation', () => { } if (process.platform === 'win32') client.send('SIGINT'); else client.kill('SIGINT'); - // Phased cancellation preserves the daemon's existing aborted-result exit code. - expect(await closed).toEqual([1, null]); + // Cancellation terminates the client like a native signal (128 + SIGINT). + expect(await closed).toEqual([130, null]); + expect(stderr).toContain('rush-client: build cancelled.'); expect(stderr).not.toMatch(/using in-process|not retried|timed out/i); expect(fixture.events().filter((event) => event.kind === 'ready')).toHaveLength(1); expect(fixture.events().filter((event) => event.kind === 'complete')).toHaveLength(2); diff --git a/common/changes/@microsoft/rush/rushd-cancel-hard-abort_2026-09-23.json b/common/changes/@microsoft/rush/rushd-cancel-hard-abort_2026-09-23.json new file mode 100644 index 0000000000..9cddb07fcf --- /dev/null +++ b/common/changes/@microsoft/rush/rushd-cancel-hard-abort_2026-09-23.json @@ -0,0 +1,11 @@ +{ + "changes": [ + { + "packageName": "@microsoft/rush", + "comment": "Add an opt-in hard abort for operation graphs (`abortCurrentIterationAsync({ terminateRunning: true })`) that terminates running shell operation process trees and reports them as Aborted; daemon engine graphs enable it.", + "type": "none" + } + ], + "packageName": "@microsoft/rush", + "email": "selarkin@microsoft.com" +} diff --git a/common/changes/@rushstack/rush-cli-client/rushd-cancel-hard-abort_2026-09-23.json b/common/changes/@rushstack/rush-cli-client/rushd-cancel-hard-abort_2026-09-23.json new file mode 100644 index 0000000000..8487bcf692 --- /dev/null +++ b/common/changes/@rushstack/rush-cli-client/rushd-cancel-hard-abort_2026-09-23.json @@ -0,0 +1,11 @@ +{ + "changes": [ + { + "packageName": "@rushstack/rush-cli-client", + "comment": "Exit with 128+signal (130/143/129) and print a cancellation notice when a daemon-routed command is cancelled; handle SIGHUP like SIGTERM.", + "type": "patch" + } + ], + "packageName": "@rushstack/rush-cli-client", + "email": "selarkin@microsoft.com" +} diff --git a/common/changes/@rushstack/rush-daemon/rushd-cancel-hard-abort_2026-09-23.json b/common/changes/@rushstack/rush-daemon/rushd-cancel-hard-abort_2026-09-23.json new file mode 100644 index 0000000000..0fc8c070a3 --- /dev/null +++ b/common/changes/@rushstack/rush-daemon/rushd-cancel-hard-abort_2026-09-23.json @@ -0,0 +1,11 @@ +{ + "changes": [ + { + "packageName": "@rushstack/rush-daemon", + "comment": "Terminate running operations when the last live client cancels a phased request, and answer a cancelling client immediately when other clients still need the shared work.", + "type": "patch" + } + ], + "packageName": "@rushstack/rush-daemon", + "email": "selarkin@microsoft.com" +} diff --git a/common/reviews/api/rush-lib.api.md b/common/reviews/api/rush-lib.api.md index 1b712f18dd..7b82e2af5d 100644 --- a/common/reviews/api/rush-lib.api.md +++ b/common/reviews/api/rush-lib.api.md @@ -721,7 +721,9 @@ export interface IOperationExecutionResult extends IBaseOperationExecutionResult // @alpha export interface IOperationGraph { readonly abortController: AbortController; - abortCurrentIterationAsync(): Promise; + abortCurrentIterationAsync(options?: { + terminateRunning?: boolean; + }): Promise; addTerminalDestination(destination: TerminalWritable): void; allowOversubscription: boolean; closeRunnersAsync(operations?: Iterable): Promise; @@ -823,6 +825,7 @@ export interface IOperationRunner { // @beta export interface IOperationRunnerContext { + readonly abortSignal?: AbortSignal; collatedWriter: CollatedWriter; // @internal createChildProcessReporter(): _IOperationChildProcessReporter | undefined; diff --git a/libraries/rush-daemon/src/PhasedRequestRouter.ts b/libraries/rush-daemon/src/PhasedRequestRouter.ts index c943eee794..dda743a498 100644 --- a/libraries/rush-daemon/src/PhasedRequestRouter.ts +++ b/libraries/rush-daemon/src/PhasedRequestRouter.ts @@ -101,6 +101,12 @@ const OBSERVED_STATUS_OVERRIDES_RETAINED: ReadonlySet = new Set OperationStatus.Blocked, OperationStatus.Skipped ]); +const IN_PROGRESS_STATUSES: ReadonlySet = new Set([ + OperationStatus.Waiting, + OperationStatus.Ready, + OperationStatus.Queued, + OperationStatus.Executing +]); /** * Routes one caller-resolved phased request through a real warm workspace operation graph. @@ -422,6 +428,12 @@ class PhasedRequestBatchCoordinator { this.#graph, participants.map((entry: IBatchEntry) => entry.selection) ); + for (const entry of batch) { + if (!participants.includes(entry)) { + // Clients that cancelled before execution must not wait for the participants' work. + this.#finishDetachedEntry(entry); + } + } for (const entry of participants) { entry.participated = true; const activeOperationIds: ReadonlySet = new Set( @@ -531,6 +543,16 @@ class PhasedRequestBatchCoordinator { } } + if ( + entry.executionStarted && + this.#currentBatch?.includes(entry) && + this.#hasLiveBatchParticipant() + ) { + // Other live participants still need the shared work: detach this client and answer it now. + this.#finishDetachedEntry(entry); + return; + } + if ( entry.executionStarted && this.#currentBatch && @@ -541,10 +563,43 @@ class PhasedRequestBatchCoordinator { } } + #finishDetachedEntry(entry: IBatchEntry): void { + entry.finishPromise ??= this.#produceResultAsync( + entry, + entry.participated, + undefined, + [], + undefined, + true + ).catch((error: unknown) => { + if (!entry.completed) { + this.#completeEntry(entry); + entry.reject(error); + } + }); + } + #isEntryLive(entry: IBatchEntry): boolean { return !entry.abortRequested && !entry.client.abortSignal.aborted && entry.outputError === undefined; } + /** + * Whether a live client still needs the current batch's iteration, including compatible requests that were + * accepted into the pending queue and will join the batch once the execution lease is acquired. + */ + #hasLiveBatchParticipant(): boolean { + if (this.#currentBatch?.some((candidate: IBatchEntry) => this.#needsIteration(candidate))) { + return true; + } + return ( + this.#acceptingCurrentBatch && + this.#pending.some( + (candidate: IBatchEntry) => + candidate.exclusivityClass === RequestExclusivityClass.SharedBuild && this.#isEntryLive(candidate) + ) + ); + } + /** Whether a live participant still waits for the running iteration to produce its result. */ #needsIteration(entry: IBatchEntry): boolean { return entry.finishPromise === undefined && this.#isEntryLive(entry); @@ -590,7 +645,8 @@ class PhasedRequestBatchCoordinator { } #requestIterationAbort(): void { - const abortPromise: Promise = this.#graph.abortCurrentIterationAsync(); + // Nobody needs the running work any more, so terminate in-flight operations instead of awaiting them. + const abortPromise: Promise = this.#graph.abortCurrentIterationAsync({ terminateRunning: true }); this.#abortTail = Promise.all([this.#abortTail, abortPromise]) .then(() => undefined) .catch((error: unknown) => { @@ -984,12 +1040,14 @@ function collectOperationOutcomes( const retained: IOperationExecutionResult | undefined = graph.resultByOperation.get(operation); let status: string | undefined; let errorMessage: string | undefined; - if ( + if (iterationInProgress && observed !== undefined) { + // While the iteration still runs, retained results may predate this iteration, and work this client + // stopped observing before it finished (a detached cancellation) was abandoned. + status = IN_PROGRESS_STATUSES.has(observed.status) ? OperationStatus.Aborted : observed.status; + errorMessage = observed.executionResult.error?.message; + } else if ( observed !== undefined && - // While the iteration still runs, retained results may predate this iteration. - (iterationInProgress || - retained === undefined || - OBSERVED_STATUS_OVERRIDES_RETAINED.has(observed.status)) + (retained === undefined || OBSERVED_STATUS_OVERRIDES_RETAINED.has(observed.status)) ) { status = observed.status; errorMessage = observed.executionResult.error?.message; @@ -998,6 +1056,10 @@ function collectOperationOutcomes( errorMessage = retained?.error?.message ?? observed?.executionResult.error?.message; } status ??= fillMissingAsAborted ? OperationStatus.Aborted : undefined; + if (fillMissingAsAborted && status !== undefined && IN_PROGRESS_STATUSES.has(status)) { + // The client stopped observing before this operation finished, e.g. because it was terminated. + status = OperationStatus.Aborted; + } if (status === undefined) { continue; } diff --git a/libraries/rush-daemon/src/test/NativeEngineTestCommands.ts b/libraries/rush-daemon/src/test/NativeEngineTestCommands.ts index b334e12e96..324916efbe 100644 --- a/libraries/rush-daemon/src/test/NativeEngineTestCommands.ts +++ b/libraries/rush-daemon/src/test/NativeEngineTestCommands.ts @@ -61,6 +61,8 @@ export async function createNativeScriptGateAsync( const server: net.Server = net.createServer((socket) => { sockets.add(socket); socket.once('close', () => sockets.delete(socket)); + // A daemon shutdown or cancellation may terminate the gated script, which resets its connection. + socket.on('error', () => undefined); entered.resolve(); }); await new Promise((resolve, reject) => { diff --git a/libraries/rush-daemon/src/test/PhasedRequestBatching.test.ts b/libraries/rush-daemon/src/test/PhasedRequestBatching.test.ts index 3d67e07eb5..1cde7c35a6 100644 --- a/libraries/rush-daemon/src/test/PhasedRequestBatching.test.ts +++ b/libraries/rush-daemon/src/test/PhasedRequestBatching.test.ts @@ -431,7 +431,7 @@ describe('shared phased request batching', () => { expect(fixture.runners.get(OPERATION_C)?.runCount).toBe(1); }); - it('reports authoritative retained status when a client cancels during a shared operation', async () => { + it('answers a client that cancels during a shared operation immediately, without waiting for the batch', async () => { const operationStarted: IDeferred = createDeferred(); const releaseOperation: IDeferred = createDeferred(); const fixture: ITestRoutingFixture = createFixture({ @@ -448,12 +448,14 @@ describe('shared phased request batching', () => { await operationStarted.promise; cancelledClient.abortController.abort(); + // The shared operation is still running for the other client. + const cancelledResult: IDaemonPhasedRequestResult = await cancelled; releaseOperation.resolve(); - const [cancelledResult, continuingResult] = await Promise.all([cancelled, continuing]); + const continuingResult: IDaemonPhasedRequestResult = await continuing; expect(cancelledResult).toMatchObject({ aborted: true, outcome: 'aborted' }); expect(cancelledResult.operationResults).toEqual([ - expect.objectContaining({ operationId: OPERATION_A, status: OperationStatus.Success }) + expect.objectContaining({ operationId: OPERATION_A, status: OperationStatus.Aborted }) ]); expect(continuingResult).toMatchObject({ exitCode: 0, outcome: 'success' }); expect(fixture.runners.get(OPERATION_A)?.runCount).toBe(1); @@ -486,7 +488,7 @@ describe('shared phased request batching', () => { ]); }); - it('preserves failure precedence when a client cancels during a failing shared operation', async () => { + it('keeps failure for the continuing client when another client cancels during a failing shared operation', async () => { const operationStarted: IDeferred = createDeferred(); const releaseOperation: IDeferred = createDeferred(); const fixture: ITestRoutingFixture = createFixture({ @@ -506,14 +508,16 @@ describe('shared phased request batching', () => { await operationStarted.promise; cancelledClient.abortController.abort(); + const cancelledResult: IDaemonPhasedRequestResult = await cancelled; releaseOperation.resolve(); - const [cancelledResult, continuingResult] = await Promise.all([cancelled, continuing]); + const continuingResult: IDaemonPhasedRequestResult = await continuing; - expect(cancelledResult).toMatchObject({ aborted: true, exitCode: 1, outcome: 'failure' }); - expect(cancelledResult.operationResults).toEqual([ + // The cancelled client detached before the shared operation failed. + expect(cancelledResult).toMatchObject({ aborted: true, outcome: 'aborted' }); + expect(continuingResult).toMatchObject({ aborted: false, exitCode: 1, outcome: 'failure' }); + expect(continuingResult.operationResults).toEqual([ expect.objectContaining({ operationId: OPERATION_A, status: OperationStatus.Failure }) ]); - expect(continuingResult).toMatchObject({ aborted: false, exitCode: 1, outcome: 'failure' }); expect(fixture.runners.get(OPERATION_A)?.runCount).toBe(1); }); diff --git a/libraries/rush-daemon/src/test/PhasedRequestCancellation.test.ts b/libraries/rush-daemon/src/test/PhasedRequestCancellation.test.ts new file mode 100644 index 0000000000..7642e99cb4 --- /dev/null +++ b/libraries/rush-daemon/src/test/PhasedRequestCancellation.test.ts @@ -0,0 +1,189 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. Licensed under the MIT license. +// See LICENSE in the project root for license information. + +import type { ITerminal } from '@rushstack/terminal'; +import type { IDaemonPhasedRequest, IDaemonPhasedRequestResult } from '@rushstack/rush-daemon-protocol'; +import { type IOperationRunnerContext, OperationStatus } from '@microsoft/rush-lib'; + +import { PhasedRequestRouter } from '../PhasedRequestRouter'; +import { + TEST_ENGINE_SHAPE, + TestOperationRunner, + TestPhasedRequestClient, + createRoutingFixture +} from './PhasedRequestRouterTestUtilities'; +import type { ITestRoutingFixture } from './PhasedRequestRouterTestUtilities'; + +const OPERATION_A: string = 'project-a (_phase:test)'; +const OPERATION_B: string = 'project-b (_phase:test)'; +const PROMPT_CANCELLATION_MS: number = 1000; + +function createRequest(requestId: string, operationId: string): IDaemonPhasedRequest { + return { + commandName: 'build', + commandOrigin: 'built-in', + engineShape: TEST_ENGINE_SHAPE, + environment: {}, + operationSelection: [{ enabledState: true, operationId }], + requestId + }; +} + +interface IHangingOperation { + readonly started: Promise; + readonly terminated: Promise; + readonly signals: AbortSignal[]; + readonly release: () => void; +} + +/** An operation that only finishes when released, or when its hard-abort signal fires (like a killed process). */ +function createHangingOperation(): { + hanging: IHangingOperation; + actionAsync: (terminal: ITerminal, context: IOperationRunnerContext) => Promise; +} { + let onStarted: () => void = () => undefined; + let onTerminated: () => void = () => undefined; + let release: () => void = () => undefined; + const released: Promise = new Promise((resolve) => (release = resolve)); + const signals: AbortSignal[] = []; + const hanging: IHangingOperation = { + started: new Promise((resolve) => (onStarted = resolve)), + terminated: new Promise((resolve) => (onTerminated = resolve)), + signals, + release: () => release() + }; + const actionAsync = async ( + terminal: ITerminal, + context: IOperationRunnerContext + ): Promise => { + const { abortSignal } = context; + if (!abortSignal) throw new Error('Expected a hard-abort signal for daemon operations.'); + signals.push(abortSignal); + onStarted(); + const aborted: Promise = new Promise((resolve) => + abortSignal.addEventListener('abort', () => resolve(), { once: true }) + ); + await Promise.race([aborted, released]); + if (abortSignal.aborted) { + onTerminated(); + return OperationStatus.Aborted; + } + }; + return { hanging, actionAsync }; +} + +function createFixture( + actionAsync: (terminal: ITerminal, context: IOperationRunnerContext) => Promise +): ITestRoutingFixture { + return createRoutingFixture( + new Map([ + [OPERATION_A, new TestOperationRunner(OPERATION_A, OperationStatus.Success, actionAsync)], + [OPERATION_B, new TestOperationRunner(OPERATION_B)] + ]), + [], + { supportsTerminateRunning: true } + ); +} + +describe('phased request client cancellation', () => { + it('terminates running operations when the last client cancels and releases the graph promptly', async () => { + const { hanging, actionAsync } = createHangingOperation(); + const fixture: ITestRoutingFixture = createFixture(actionAsync); + const abortSpy: jest.SpyInstance = jest.spyOn(fixture.graph, 'abortCurrentIterationAsync'); + const router: PhasedRequestRouter = new PhasedRequestRouter(fixture.session); + const client: TestPhasedRequestClient = new TestPhasedRequestClient('one'); + const cancelled: Promise = router.executeAsync( + createRequest('cancelled', OPERATION_A), + client + ); + await hanging.started; + + const cancelledAt: number = Date.now(); + client.abortController.abort(); + const result: IDaemonPhasedRequestResult = await cancelled; + + expect(Date.now() - cancelledAt).toBeLessThan(PROMPT_CANCELLATION_MS); + await hanging.terminated; + expect(abortSpy).toHaveBeenCalledWith({ terminateRunning: true }); + expect(result).toMatchObject({ aborted: true, outcome: 'aborted' }); + expect(result.operationResults).toEqual([ + expect.objectContaining({ operationId: OPERATION_A, status: OperationStatus.Aborted }) + ]); + // The aborted operation is not retained, and the next client is admitted immediately. + expect(fixture.graph.resultByOperation.size).toBe(0); + const next: IDaemonPhasedRequestResult = await router.executeAsync( + createRequest('next', OPERATION_B), + new TestPhasedRequestClient('two') + ); + expect(next).toMatchObject({ exitCode: 0, outcome: 'success' }); + }); + + it('only detaches a cancelling client while another live client still needs the running work', async () => { + const { hanging, actionAsync } = createHangingOperation(); + const fixture: ITestRoutingFixture = createFixture(actionAsync); + const abortSpy: jest.SpyInstance = jest.spyOn(fixture.graph, 'abortCurrentIterationAsync'); + const router: PhasedRequestRouter = new PhasedRequestRouter(fixture.session); + const cancelledClient: TestPhasedRequestClient = new TestPhasedRequestClient('one'); + const cancelled: Promise = router.executeAsync( + createRequest('cancelled', OPERATION_A), + cancelledClient + ); + const continuing: Promise = router.executeAsync( + createRequest('continuing', OPERATION_A), + new TestPhasedRequestClient('two') + ); + await hanging.started; + const abortCallsBeforeCancellation: number = abortSpy.mock.calls.length; + + cancelledClient.abortController.abort(); + const cancelledResult: IDaemonPhasedRequestResult = await cancelled; + + expect(cancelledResult).toMatchObject({ aborted: true, outcome: 'aborted' }); + expect(abortSpy).toHaveBeenCalledTimes(abortCallsBeforeCancellation); + expect(hanging.signals.map((signal: AbortSignal) => signal.aborted)).toEqual([false]); + + hanging.release(); + expect(await continuing).toMatchObject({ aborted: false, exitCode: 0, outcome: 'success' }); + expect(fixture.runners.get(OPERATION_A)?.runCount).toBe(1); + }); + + it('answers a cancelling client at once when an accepted pending client will join the batch', async () => { + const { hanging, actionAsync } = createHangingOperation(); + const fixture: ITestRoutingFixture = createFixture(actionAsync); + const abortSpy: jest.SpyInstance = jest.spyOn(fixture.graph, 'abortCurrentIterationAsync'); + let onReconciling: () => void = () => undefined; + const reconciling: Promise = new Promise((resolve) => (onReconciling = resolve)); + let releaseReconcile: () => void = () => undefined; + const reconcileReleased: Promise = new Promise((resolve) => (releaseReconcile = resolve)); + fixture.session.onReconcileAsync = async () => { + onReconciling(); + await reconcileReleased; + }; + const router: PhasedRequestRouter = new PhasedRequestRouter(fixture.session); + const cancelledClient: TestPhasedRequestClient = new TestPhasedRequestClient('one'); + const cancelled: Promise = router.executeAsync( + createRequest('cancelled', OPERATION_A), + cancelledClient + ); + await reconciling; + // Accepted while the batch is still being prepared, so it joins once preparation finishes. + const continuing: Promise = router.executeAsync( + createRequest('continuing', OPERATION_A), + new TestPhasedRequestClient('two') + ); + // Let the continuing request finish preparation and enter the pending queue. + for (let tick: number = 0; tick < 20; tick++) { + await new Promise((resolve) => setImmediate(resolve)); + } + + cancelledClient.abortController.abort(); + expect(await cancelled).toMatchObject({ aborted: true, outcome: 'aborted' }); + + releaseReconcile(); + await hanging.started; + expect(hanging.signals.map((signal: AbortSignal) => signal.aborted)).toEqual([false]); + expect(abortSpy).not.toHaveBeenCalledWith({ terminateRunning: true }); + hanging.release(); + expect(await continuing).toMatchObject({ aborted: false, exitCode: 0, outcome: 'success' }); + }); +}); diff --git a/libraries/rush-daemon/src/test/PhasedRequestRouterTestUtilities.ts b/libraries/rush-daemon/src/test/PhasedRequestRouterTestUtilities.ts index 41f17440d7..58a52c478d 100644 --- a/libraries/rush-daemon/src/test/PhasedRequestRouterTestUtilities.ts +++ b/libraries/rush-daemon/src/test/PhasedRequestRouterTestUtilities.ts @@ -137,14 +137,16 @@ export class TestOperationRunner implements IOperationRunner { public closeCount: number = 0; public runCount: number = 0; - readonly #actionAsync: ((terminal: ITerminal) => Promise) | undefined; + readonly #actionAsync: + | ((terminal: ITerminal, context: IOperationRunnerContext) => Promise) + | undefined; readonly #status: OperationStatus; public readonly name: string; public constructor( name: string, status: OperationStatus = OperationStatus.Success, - actionAsync?: (terminal: ITerminal) => Promise + actionAsync?: (terminal: ITerminal, context: IOperationRunnerContext) => Promise ) { this.name = name; this.#status = status; @@ -160,8 +162,8 @@ export class TestOperationRunner implements IOperationRunner { this.runCount++; return context.runWithTerminalAsync( async (terminal: ITerminal): Promise => { - await this.#actionAsync?.(terminal); - return this.#status; + const status: void | OperationStatus = await this.#actionAsync?.(terminal, context); + return status ?? this.#status; }, { createLogFile: false, logFileSuffix: '' } ); @@ -216,7 +218,8 @@ export class TestRoutingWorkspaceSession implements IWorkspaceSession { export function createRoutingFixture( runnerById: ReadonlyMap, - dependencies: ReadonlyArray = [] + dependencies: ReadonlyArray = [], + graphOptionOverrides: Partial = {} ): ITestRoutingFixture { const operations: Map = new Map(); const runners: Map = new Map(runnerById); @@ -252,7 +255,8 @@ export function createRoutingFixture( destinations: [new MockWritable()], parallelism: 1, pauseNextIteration: false, - quietMode: false + quietMode: false, + ...graphOptionOverrides }; // The package's bundled public declarations and deep-import declarations describe the same runtime classes, // but TypeScript assigns them distinct recursive identities. diff --git a/libraries/rush-lib/src/cli/scriptActions/PhasedScriptAction.ts b/libraries/rush-lib/src/cli/scriptActions/PhasedScriptAction.ts index 89a0dbd3f7..b04cb177b4 100644 --- a/libraries/rush-lib/src/cli/scriptActions/PhasedScriptAction.ts +++ b/libraries/rush-lib/src/cli/scriptActions/PhasedScriptAction.ts @@ -815,6 +815,7 @@ export class PhasedScriptAction extends BaseScriptAction i getInputsSnapshotAsync: getGraphInputsSnapshotAsync, abortController: this.sessionAbortController, closeRunnersOnAbort: !onEngine, + supportsTerminateRunning: !!onEngine, telemetry: executionTelemetryHandler }; @@ -1066,7 +1067,7 @@ async function disposeEngineGraphAsync( graph.abortController.abort(); const errors: unknown[] = []; for (const cleanupAsync of [ - () => graph.abortCurrentIterationAsync(), + () => graph.abortCurrentIterationAsync({ terminateRunning: true }), () => graph.closeRunnersAsync(), async () => { await cobuildConfiguration?.destroyLockProviderAsync(); diff --git a/libraries/rush-lib/src/logic/operations/IOperationGraph.ts b/libraries/rush-lib/src/logic/operations/IOperationGraph.ts index 5575bf77ab..e6b4316fbe 100644 --- a/libraries/rush-lib/src/logic/operations/IOperationGraph.ts +++ b/libraries/rush-lib/src/logic/operations/IOperationGraph.ts @@ -106,8 +106,11 @@ export interface IOperationGraph { /** * Abort the current execution iteration, if any. Operations that have already started * will run to completion; only operations that have not yet begun will be aborted. + * + * If `options.terminateRunning` is true and the graph supports it, operations that are already running are also + * signaled to terminate (via `IOperationRunnerContext.abortSignal`) and are reported as `Aborted`. */ - abortCurrentIterationAsync(): Promise; + abortCurrentIterationAsync(options?: { terminateRunning?: boolean }): Promise; /** * Cleans up any resources used by the operation runners, if applicable. diff --git a/libraries/rush-lib/src/logic/operations/IOperationRunner.ts b/libraries/rush-lib/src/logic/operations/IOperationRunner.ts index 1e820a0af4..f422083d85 100644 --- a/libraries/rush-lib/src/logic/operations/IOperationRunner.ts +++ b/libraries/rush-lib/src/logic/operations/IOperationRunner.ts @@ -63,6 +63,14 @@ export interface IOperationRunnerContext { */ readonly shouldRunnerPersist: boolean; + /** + * When defined, this signal is aborted if the host requests termination of in-progress work + * (see `IOperationGraph.abortCurrentIterationAsync` with `terminateRunning: true`). Runners that observe it + * should promptly stop any work they started (for example, kill their child process tree) and return + * `OperationStatus.Aborted`. It is `undefined` when the graph does not support terminating running operations. + */ + readonly abortSignal?: AbortSignal; + /** * The environment in which the operation is being executed. * A return value of `undefined` indicates that it should inherit the environment from the parent process. diff --git a/libraries/rush-lib/src/logic/operations/OperationExecutionRecord.ts b/libraries/rush-lib/src/logic/operations/OperationExecutionRecord.ts index ed7c448268..129d5c04c5 100644 --- a/libraries/rush-lib/src/logic/operations/OperationExecutionRecord.ts +++ b/libraries/rush-lib/src/logic/operations/OperationExecutionRecord.ts @@ -51,6 +51,8 @@ export interface IOperationExecutionRecordContext { invalidate?: (operations: Iterable, reason: string) => void; inputsSnapshot: IInputsSnapshot | undefined; maxParallelism: number; + /** Aborted when the host requests termination of running operations in this iteration. */ + terminateSignal?: AbortSignal; /** * Optional structured event sink for dual-emit. When present, every status @@ -245,6 +247,10 @@ export class OperationExecutionRecord implements IOperationRunnerContext, IOpera return this.#context.createEnvironment?.(this); } + public get abortSignal(): AbortSignal | undefined { + return this.#context.terminateSignal; + } + public getInvalidateCallback(): (reason: string) => void { const invalidateFn: ((operations: Iterable, reason: string) => void) | undefined = this.#context.invalidate; diff --git a/libraries/rush-lib/src/logic/operations/OperationGraph.ts b/libraries/rush-lib/src/logic/operations/OperationGraph.ts index 17a57096d0..a3ed25bb1c 100644 --- a/libraries/rush-lib/src/logic/operations/OperationGraph.ts +++ b/libraries/rush-lib/src/logic/operations/OperationGraph.ts @@ -54,6 +54,12 @@ export interface IOperationGraphOptions { abortController: AbortController; /** Hosts with awaited lifetime cleanup can disable the legacy fire-and-forget abort cleanup. */ closeRunnersOnAbort?: boolean; + /** + * If true, runners receive `IOperationRunnerContext.abortSignal`, and + * `abortCurrentIterationAsync({ terminateRunning: true })` terminates operations that are already running. + * Runners may isolate their child processes (e.g. in a separate process group) to support this. + */ + supportsTerminateRunning?: boolean; isWatch?: boolean; pauseNextIteration?: boolean; @@ -82,6 +88,7 @@ interface IStatefulExecutionContext { */ interface IExecutionIterationContext extends IOperationExecutionRecordContext { abortController: AbortController; + terminateController: AbortController | undefined; terminal: CollatedTerminal; records: Map; @@ -176,6 +183,7 @@ export class OperationGraph implements IOperationGraph { // Immutable properties from options readonly #isWatch: boolean; + readonly #supportsTerminateRunning: boolean; readonly #telemetry: IOperationGraphTelemetry | undefined; readonly #getInputsSnapshotAsync: (() => Promise) | undefined; @@ -220,11 +228,13 @@ export class OperationGraph implements IOperationGraph { abortController, isWatch = false, pauseNextIteration = false, + supportsTerminateRunning = false, telemetry, getInputsSnapshotAsync } = options; this.operations = operations; + this.#supportsTerminateRunning = supportsTerminateRunning; this.#maxParallelism = maxParallelism; this.#parallelism = coerceParallelism(parallelism, maxParallelism, 1); @@ -627,9 +637,12 @@ export class OperationGraph implements IOperationGraph { return true; } - public async abortCurrentIterationAsync(): Promise { + public async abortCurrentIterationAsync(options?: { terminateRunning?: boolean }): Promise { const iteration: IExecutionIterationContext | undefined = this.#currentIteration; if (iteration) { + if (options?.terminateRunning) { + iteration.terminateController?.abort(); + } iteration.abortController.abort(); try { await iteration.promise; @@ -709,10 +722,16 @@ export class OperationGraph implements IOperationGraph { return hooks.createEnvironmentForOperation.call({ ...process.env }, record); } + const terminateController: AbortController | undefined = this.#supportsTerminateRunning + ? new AbortController() + : undefined; + // Convert the developer graph to the mutable execution graph const iterationContext: IExecutionIterationContext = { iterationId: this.#nextIterationId++, abortController, + terminateController, + terminateSignal: terminateController?.signal, startTime, streamCollator, terminal, diff --git a/libraries/rush-lib/src/logic/operations/ShellOperationRunner.ts b/libraries/rush-lib/src/logic/operations/ShellOperationRunner.ts index c76fce6931..627236e23b 100644 --- a/libraries/rush-lib/src/logic/operations/ShellOperationRunner.ts +++ b/libraries/rush-lib/src/logic/operations/ShellOperationRunner.ts @@ -4,7 +4,7 @@ import type * as child_process from 'node:child_process'; import * as path from 'node:path'; -import { Path } from '@rushstack/node-core-library'; +import { Path, SubprocessTerminator } from '@rushstack/node-core-library'; import { type ITerminal, type ITerminalProvider, TerminalProviderSeverity } from '@rushstack/terminal'; import type { IPhase } from '../../api/CommandLineConfiguration'; @@ -103,7 +103,7 @@ export class ShellOperationRunner implements IOperationRunner { const { rushConfiguration, projectFolder } = this.#rushProject; - const { environment: initialEnvironment } = context; + const { environment: initialEnvironment, abortSignal } = context; const childProcessReporter: IOperationChildProcessReporter | undefined = !IS_WINDOWS && isHeftCommand(commandToRun) ? context.createChildProcessReporter() : undefined; @@ -117,8 +117,28 @@ export class ShellOperationRunner implements IOperationRunner { }, initialEnvironment, additionalEnvironment: childProcessReporter?.environment, - stdio: childProcessReporter?.stdio + stdio: childProcessReporter?.stdio, + // Isolate the process tree so that a hard abort can terminate it. + connectSubprocessTerminator: abortSignal !== undefined }); + const terminateProcessTree: () => void = () => { + try { + if (!IS_WINDOWS && subProcess.pid !== undefined && typeof subProcess.exitCode === 'number') { + // The shell already exited, but descendants in its process group may still hold its stdio open. + // killProcessTree() is a no-op in that state, so signal the process group directly. + killExitedProcessGroup(subProcess.pid); + } else { + SubprocessTerminator.killProcessTree(subProcess, SubprocessTerminator.RECOMMENDED_OPTIONS); + } + } catch (error) { + terminal.writeErrorLine(`Failed to terminate the operation process tree: ${error}`); + } + }; + if (abortSignal?.aborted) { + terminateProcessTree(); + } else { + abortSignal?.addEventListener('abort', terminateProcessTree, { once: true }); + } let reporterError: Error | undefined; const reporterDrainPromise: Promise = childProcessReporter ? childProcessReporter @@ -164,9 +184,14 @@ export class ShellOperationRunner implements IOperationRunner { const [{ exitCode, signal }]: [ { readonly exitCode: number | null; readonly signal: NodeJS.Signals | null }, void - ] = await Promise.all([closePromise, reporterDrainPromise]); + ] = await Promise.all([closePromise, reporterDrainPromise]).finally(() => { + abortSignal?.removeEventListener('abort', terminateProcessTree); + }); - if (signal) { + if (abortSignal?.aborted) { + terminal.writeLine('Terminated because the operation was aborted.'); + return OperationStatus.Aborted; + } else if (signal) { // eslint-disable-next-line require-atomic-updates -- This operation context has one active runner. context.error = new OperationError('error', `Terminated by signal: ${signal}`); return OperationStatus.Failure; @@ -195,6 +220,17 @@ export class ShellOperationRunner implements IOperationRunner { } } +function killExitedProcessGroup(pid: number): void { + try { + // The process group ID cannot be reused while any member of the group is still alive. + process.kill(-pid, 'SIGKILL'); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ESRCH') { + throw error; + } + } +} + /** * Returns whether a lifecycle command directly launches Heft. * diff --git a/libraries/rush-lib/src/logic/operations/test/ShellOperationRunnerAbort.test.ts b/libraries/rush-lib/src/logic/operations/test/ShellOperationRunnerAbort.test.ts new file mode 100644 index 0000000000..930f63170d --- /dev/null +++ b/libraries/rush-lib/src/logic/operations/test/ShellOperationRunnerAbort.test.ts @@ -0,0 +1,199 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. Licensed under the MIT license. +// See LICENSE in the project root for license information. + +jest.mock('../OperationStateFile'); +jest.mock('../ProjectLogWritable', () => { + const actual = jest.requireActual('../ProjectLogWritable'); + const { MockWritable } = jest.requireActual('@rushstack/terminal'); + return { ...actual, initializeProjectLogFilesAsync: jest.fn(async () => new MockWritable()) }; +}); + +import { spawn, type ChildProcess } from 'node:child_process'; +import { once } from 'node:events'; +import { setTimeout as delayAsync } from 'node:timers/promises'; + +import { SubprocessTerminator } from '@rushstack/node-core-library'; +import { MockWritable } from '@rushstack/terminal'; + +import type { IPhase } from '../../../api/CommandLineConfiguration'; +import type { RushConfigurationProject } from '../../../api/RushConfigurationProject'; +import { type ILifecycleCommandOptions, Utilities } from '../../../utilities/Utilities'; +import type { IExecutionResult } from '../IOperationExecutionResult'; +import { Operation } from '../Operation'; +import { OperationGraph } from '../OperationGraph'; +import { OperationStatus } from '../OperationStatus'; +import { ShellOperationRunner } from '../ShellOperationRunner'; + +// Spawns a grandchild that never exits, reports its PID, then waits forever itself. +const NEVER_ENDING_TREE_SCRIPT: string = ` + const grandchild = require('node:child_process').spawn(process.execPath, ['-e', 'setInterval(() => {}, 1000)'], { stdio: 'ignore' }); + process.stdout.write('grandchild=' + grandchild.pid + '\\n'); + setInterval(() => {}, 1000); +`; + +// Spawns a grandchild that inherits stdout and never exits, reports its PID, then exits itself. +const EXITED_PARENT_TREE_SCRIPT: string = ` + const grandchild = require('node:child_process').spawn(process.execPath, ['-e', 'setInterval(() => {}, 1000)'], { stdio: ['ignore', 'inherit', 'ignore'] }); + process.stdout.write('grandchild=' + grandchild.pid + '\\n', () => process.exit(0)); +`; + +function isAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + +async function waitForExitAsync(pid: number, timeoutMs: number): Promise { + const deadline: number = Date.now() + timeoutMs; + while (isAlive(pid)) { + if (Date.now() >= deadline) return false; + await delayAsync(20); + } + return true; +} + +describe('ShellOperationRunner hard abort', () => { + let child: ChildProcess | undefined; + let grandchildPid: number | undefined; + let spawnOptions: ILifecycleCommandOptions | undefined; + + function createGraph( + supportsTerminateRunning: boolean, + script: string = NEVER_ENDING_TREE_SCRIPT + ): { + graph: OperationGraph; + grandchildStarted: Promise; + } { + let onGrandchild: (pid: number) => void = () => undefined; + const grandchildStarted: Promise = new Promise((resolve) => (onGrandchild = resolve)); + jest.spyOn(Utilities, 'executeLifecycleCommandAsync').mockImplementation((command, options) => { + spawnOptions = options; + child = spawn(process.execPath, ['-e', script], { + stdio: ['ignore', 'pipe', 'pipe'], + detached: !!options.connectSubprocessTerminator && SubprocessTerminator.RECOMMENDED_OPTIONS.detached + }); + child.stdout!.on('data', (chunk: Buffer) => { + const match: RegExpMatchArray | null = chunk.toString().match(/grandchild=(\d+)/); + if (match) { + grandchildPid = Number(match[1]); + onGrandchild(grandchildPid); + } + }); + return child; + }); + const phase: IPhase = { + name: 'build', + allowWarningsOnSuccess: false, + associatedParameters: new Set(), + dependencies: { self: new Set(), upstream: new Set() }, + isSynthetic: false, + logFilenameIdentifier: 'build', + missingScriptBehavior: 'silent' + }; + const project: RushConfigurationProject = { + packageName: 'sleeper', + projectFolder: __dirname, + rushConfiguration: { commonTempFolder: __dirname } + } as RushConfigurationProject; + const runner: ShellOperationRunner = new ShellOperationRunner({ + phase, + rushProject: project, + displayName: 'sleeper', + initialCommand: 'node sleeper.js', + incrementalCommand: undefined, + commandForHash: 'node sleeper.js', + ignoredParameterValues: [] + }); + const operation: Operation = new Operation({ phase, project, runner, logFilenameIdentifier: 'sleeper' }); + const graph: OperationGraph = new OperationGraph(new Set([operation]), { + quietMode: true, + debugMode: false, + parallelism: 1, + allowOversubscription: true, + destinations: [new MockWritable()], + abortController: new AbortController(), + supportsTerminateRunning + }); + return { graph, grandchildStarted }; + } + + afterEach(async () => { + jest.restoreAllMocks(); + const lastChild: ChildProcess | undefined = child; + const lastGrandchildPid: number | undefined = grandchildPid; + child = undefined; + grandchildPid = undefined; + spawnOptions = undefined; + if (lastGrandchildPid !== undefined && isAlive(lastGrandchildPid)) { + process.kill(lastGrandchildPid, 'SIGKILL'); + } + if (lastChild && lastChild.exitCode === null && lastChild.signalCode === null) { + const closed: Promise = once(lastChild, 'close'); + lastChild.kill('SIGKILL'); + await closed; + } + }); + + it('kills the running process tree and reports Aborted without waiting for it to finish', async () => { + const { graph, grandchildStarted } = createGraph(true); + const execution: Promise = graph.executeAsync({}); + const pid: number = await grandchildStarted; + expect(spawnOptions?.connectSubprocessTerminator).toBe(true); + + const abortStart: number = Date.now(); + await graph.abortCurrentIterationAsync({ terminateRunning: true }); + const result: IExecutionResult = await execution; + + expect(Date.now() - abortStart).toBeLessThan(5000); + expect(result.status).toBe(OperationStatus.Aborted); + const [record] = [...result.operationResults.values()]; + expect(record.status).toBe(OperationStatus.Aborted); + expect(record.error).toBeUndefined(); + expect(child!.exitCode !== null || child!.signalCode !== null).toBe(true); + expect(await waitForExitAsync(pid, 5000)).toBe(true); + // Aborted operations are not retained as the last execution result. + expect(graph.resultByOperation.size).toBe(0); + }, 20000); + + // Process groups are POSIX-only; Windows terminates the tree via TaskKill while the parent is alive. + (process.platform === 'win32' ? it.skip : it)( + 'kills descendants that keep the output open after the shell itself has exited', + async () => { + const { graph, grandchildStarted } = createGraph(true, EXITED_PARENT_TREE_SCRIPT); + const execution: Promise = graph.executeAsync({}); + const pid: number = await grandchildStarted; + const parent: ChildProcess = child!; + if (parent.exitCode === null) { + await once(parent, 'exit'); + } + expect(isAlive(pid)).toBe(true); + + await graph.abortCurrentIterationAsync({ terminateRunning: true }); + const result: IExecutionResult = await execution; + + expect(result.status).toBe(OperationStatus.Aborted); + expect(await waitForExitAsync(pid, 5000)).toBe(true); + }, + 20000 + ); + + it('does not isolate or terminate processes when the graph does not support it', async () => { + const { graph, grandchildStarted } = createGraph(false); + const execution: Promise = graph.executeAsync({}); + await grandchildStarted; + expect(spawnOptions?.connectSubprocessTerminator).toBe(false); + + // A soft abort only prevents unstarted work; the running process keeps going. + const abortPromise: Promise = graph.abortCurrentIterationAsync({ terminateRunning: true }); + await delayAsync(300); + expect(child!.exitCode).toBeNull(); + expect(child!.signalCode).toBeNull(); + + child!.kill('SIGKILL'); + await abortPromise; + expect((await execution).status).toBe(OperationStatus.Failure); + }, 20000); +});