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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 40 additions & 0 deletions apps/rush-cli-client/src/clientCancellation.ts
Original file line number Diff line number Diff line change
@@ -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<NodeJS.Signals> = ['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;
}
}
38 changes: 28 additions & 10 deletions apps/rush-cli-client/src/launchClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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,
Expand All @@ -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<void> => {
if (discoveryLines.length > 0) {
Expand Down Expand Up @@ -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);
Expand All @@ -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`);
Expand Down
51 changes: 51 additions & 0 deletions apps/rush-cli-client/src/test/launchClient.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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;
Expand Down Expand Up @@ -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);
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
@@ -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"
}
Original file line number Diff line number Diff line change
@@ -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"
}
Original file line number Diff line number Diff line change
@@ -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"
}
5 changes: 4 additions & 1 deletion common/reviews/api/rush-lib.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -721,7 +721,9 @@ export interface IOperationExecutionResult extends IBaseOperationExecutionResult
// @alpha
export interface IOperationGraph {
readonly abortController: AbortController;
abortCurrentIterationAsync(): Promise<void>;
abortCurrentIterationAsync(options?: {
terminateRunning?: boolean;
}): Promise<void>;
addTerminalDestination(destination: TerminalWritable): void;
allowOversubscription: boolean;
closeRunnersAsync(operations?: Iterable<Operation>): Promise<void>;
Expand Down Expand Up @@ -823,6 +825,7 @@ export interface IOperationRunner {

// @beta
export interface IOperationRunnerContext {
readonly abortSignal?: AbortSignal;
collatedWriter: CollatedWriter;
// @internal
createChildProcessReporter(): _IOperationChildProcessReporter | undefined;
Expand Down
74 changes: 68 additions & 6 deletions libraries/rush-daemon/src/PhasedRequestRouter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,12 @@ const OBSERVED_STATUS_OVERRIDES_RETAINED: ReadonlySet<OperationStatus> = new Set
OperationStatus.Blocked,
OperationStatus.Skipped
]);
const IN_PROGRESS_STATUSES: ReadonlySet<string> = new Set<string>([
OperationStatus.Waiting,
OperationStatus.Ready,
OperationStatus.Queued,
OperationStatus.Executing
]);

/**
* Routes one caller-resolved phased request through a real warm workspace operation graph.
Expand Down Expand Up @@ -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<string> = new Set(
Expand Down Expand Up @@ -531,6 +543,16 @@ class PhasedRequestBatchCoordinator {
}
}

if (
entry.executionStarted &&
this.#currentBatch?.includes(entry) &&
this.#hasLiveBatchParticipant()
) {
Comment thread
TheLarkInn marked this conversation as resolved.
// Other live participants still need the shared work: detach this client and answer it now.
this.#finishDetachedEntry(entry);
return;
}

if (
entry.executionStarted &&
this.#currentBatch &&
Expand All @@ -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);
Expand Down Expand Up @@ -590,7 +645,8 @@ class PhasedRequestBatchCoordinator {
}

#requestIterationAbort(): void {
const abortPromise: Promise<void> = this.#graph.abortCurrentIterationAsync();
// Nobody needs the running work any more, so terminate in-flight operations instead of awaiting them.
const abortPromise: Promise<void> = this.#graph.abortCurrentIterationAsync({ terminateRunning: true });
this.#abortTail = Promise.all([this.#abortTail, abortPromise])
.then(() => undefined)
.catch((error: unknown) => {
Expand Down Expand Up @@ -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;
Expand All @@ -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;
}
Expand Down
2 changes: 2 additions & 0 deletions libraries/rush-daemon/src/test/NativeEngineTestCommands.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>((resolve, reject) => {
Expand Down
Loading
Loading