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
29 changes: 6 additions & 23 deletions apps/rush-cli-client/src/test/RushXDaemonAlias.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,7 @@ console.log(JSON.stringify({
expect(result.stdout.length).toBe(0);
});

it('pins physical request identity when an alias changes while waiting for admission', async () => {
it('runs through an alias without waiting for admission behind an exclusive request', async () => {
fixture.write(
'retargeted/projects/a/package.json',
JSON.stringify({
Expand Down Expand Up @@ -163,33 +163,16 @@ console.log(JSON.stringify({
invocationKind: 'rush'
});
await holding;
const input: PassThrough = new PassThrough();
input.end('not consumed before execution');
const onQueuePositionAsync = jest.fn(async () => undefined);
try {
const result = await fixture.runAsync(
fixture.request(['-q', 'build'], path.join(aliasRoot, 'projects/a')),
undefined,
{
stdin: input,
onQueuePositionAsync: async () => {
fs.unlinkSync(aliasRoot);
fs.symlinkSync(
path.join(physicalRoot, 'retargeted'),
aliasRoot,
process.platform === 'win32' ? 'junction' : 'dir'
);
release();
}
}
{ onQueuePositionAsync }
);
if (process.platform === 'win32') {
expect(result.outcome).toMatchObject({
kind: 'result',
result: { exitCode: 1, errorMessage: expect.stringContaining('invocation directory changed') }
});
expect(input.read().toString()).toBe('not consumed before execution');
} else {
expect(result.exitCode).toBe(0);
expect(result.exitCode).toBe(0);
expect(onQueuePositionAsync).not.toHaveBeenCalled();
if (process.platform !== 'win32') {
expect(JSON.parse(result.stdout.toString()).cwd).toBe(path.join(physicalRoot, 'projects/a'));
}
expect(fs.existsSync(path.join(physicalRoot, 'retargeted/projects/a/executed'))).toBe(false);
Expand Down
25 changes: 7 additions & 18 deletions apps/rush-cli-client/src/test/RushXDaemonBoundaries.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,7 @@ describe('native Rushx execution boundaries', () => {
}
);

it('fails visibly without running or reading when hooks change during queue admission', async () => {
it('runs a script without queueing behind an exclusive workspace request', async () => {
const cwd: string = await startAsync();
let release: () => void = () => {};
let started: () => void = () => {};
Expand All @@ -177,25 +177,14 @@ describe('native Rushx execution boundaries', () => {
});
const holder = fixture.runAsync({ ...fixture.request(['hold'], cwd), invocationKind: 'rush' });
await holding;
const stdin: PassThrough = new PassThrough();
stdin.end('not-consumed');
const onQueuePositionAsync = jest.fn(async () => undefined);
try {
const result = await fixture.runAsync(fixture.request(['build'], cwd), undefined, {
stdin,
onQueuePositionAsync: async () => {
const file: string = path.join(fixture.folder, 'rush.json');
const config = JSON.parse(fs.readFileSync(file, 'utf8'));
fixture.write(
'rush.json',
JSON.stringify({ ...config, eventHooks: { preRushx: ['node hook.cjs'] } })
);
release();
}
const result = await fixture.runAsync(fixture.request(['-q', 'build'], cwd), undefined, {
onQueuePositionAsync
});
expect(result.outcome).toMatchObject({ kind: 'result', result: { exitCode: 1, outcome: 'failure' } });
expect(result.stderr.toString()).toContain('Rush configuration changed');
expect(stdin.read().toString()).toBe('not-consumed');
expect(fs.existsSync(path.join(cwd, 'runs.txt'))).toBe(false);
expect(result.outcome).toMatchObject({ kind: 'result', result: { exitCode: 0, outcome: 'success' } });
expect(onQueuePositionAsync).not.toHaveBeenCalled();
expect(fs.readFileSync(path.join(cwd, 'runs.txt'), 'utf8')).toBe('ran\n');
} finally {
release();
await holder;
Expand Down
14 changes: 7 additions & 7 deletions apps/rush-cli-client/src/test/pipedInput.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -132,11 +132,11 @@ describe('standalone client piped input', () => {
);

it.each([
{ args: ['--no-wait'], admission: { noWait: true }, reason: 'no-wait' },
{ args: ['--wait-timeout=0.01'], admission: { waitTimeoutMs: 10 }, reason: 'wait-timeout' }
{ args: ['--no-wait'], admission: { noWait: true } },
{ args: ['--wait-timeout=0.01'], admission: { waitTimeoutMs: 10 } }
])(
'forwards $args without leaking queue flags to scripts',
async ({ args, admission, reason }) => {
'forwards $args without leaking queue flags to scripts, which do not queue for admission',
async ({ args, admission }) => {
let started: () => void = () => {};
let release: () => void = () => {};
const running: Promise<void> = new Promise((resolve) => {
Expand Down Expand Up @@ -178,9 +178,9 @@ describe('standalone client piped input', () => {
})
]);
const result: IPipedResult = await invokeAsync(Buffer.alloc(0), true, args);
expect(result.code).toBe(1);
expect(result.stderr.toString()).toContain(`daemon admission failed (${reason})`);
expect(runScriptAsync).not.toHaveBeenCalled();
expect(result.code).toBe(0);
expect(result.stderr.toString()).not.toContain('daemon admission failed');
expect(runScriptAsync).toHaveBeenCalledTimes(1);
} finally {
release();
await holding;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
{
"changes": [
{
"packageName": "@rushstack/rush-daemon",
"comment": "Run Rushx package scripts without exclusive workspace admission, complete global command requests on child exit plus a bounded output drain, and terminate an exited child's process group on cancellation.",
"type": "none"
}
],
"packageName": "@rushstack/rush-daemon"
}
4 changes: 4 additions & 0 deletions common/reviews/api/rush-daemon.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
/// <reference types="node" />

import * as childProcess from 'node:child_process';
import type { DaemonInvocationKind } from '@rushstack/rush-daemon-protocol';
import type { DaemonRushCommandOrigin } from '@rushstack/rush-daemon-protocol';
import type { DaemonTerminalRequirement } from '@rushstack/rush-daemon-protocol';
import * as fs from 'node:fs';
Expand Down Expand Up @@ -414,6 +415,8 @@ export interface IResolvedGlobalCommandRequest {
// (undocumented)
readonly environment: IGlobalCommandEnvironment;
// (undocumented)
readonly invocationKind?: DaemonInvocationKind;
// (undocumented)
readonly requestId: string;
// (undocumented)
readonly terminal: IGlobalCommandTerminalProperties;
Expand All @@ -431,6 +434,7 @@ export interface IResolveGlobalCommandRequestOptions {
readonly cwd: string;
// (undocumented)
readonly environment: Readonly<NodeJS.ProcessEnv>;
readonly invocationKind?: DaemonInvocationKind;
// (undocumented)
readonly requestId: string;
// (undocumented)
Expand Down
1 change: 1 addition & 0 deletions libraries/rush-daemon/src/DaemonRequestDispatcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,7 @@ async function dispatchWorkspaceRequestAsync(
commandOrigin: isRushxInvocation(envelope) ? 'custom' : envelope.commandOrigin,
cwd: envelope.cwd,
environment: envelope.environment,
invocationKind: isRushxInvocation(envelope) ? 'rushx' : 'rush',
requestId: envelope.requestId,
terminal: {
...envelope.terminal,
Expand Down
138 changes: 117 additions & 21 deletions libraries/rush-daemon/src/GlobalCommandExecutionContext.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,11 @@ import { waitForLinuxProcessGroupExitAsync } from './LinuxProcessGroupExit';
import { recordWorkspaceRequestCleanupFailure } from './WorkspaceRequestResources';

const MAX_PENDING_TERMINAL_BYTES: number = 1024 * 1024;
/**
* How long an exited child's output pipes may stay idle before they are destroyed. A pipe that is still held open by
* a descendant must not keep the request (and its client) waiting indefinitely.
*/
const CHILD_OUTPUT_DRAIN_IDLE_TIMEOUT_MS: number = 250;

/**
* Options for a request-scoped child process.
Expand Down Expand Up @@ -160,6 +165,12 @@ interface ITrackedChild {
readonly completion: Promise<void>;
}

/** Forwarding activity for one child's output, used to bound its post-exit drain. */
interface IChildOutputProgress {
events: number;
pendingWrites: number;
}

export class GlobalCommandExecutionContext implements IGlobalCommandExecutionContext, AsyncDisposable {
readonly #abortController: AbortController = new AbortController();
readonly #client: IGlobalCommandRequestClient;
Expand Down Expand Up @@ -260,18 +271,22 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon
windowsHide: options.windowsHide
});
SubprocessTerminator.killProcessTreeOnExit(child, SubprocessTerminator.RECOMMENDED_OPTIONS);
const completion: Promise<void> = this.#trackChildAsync(child).catch((error: unknown) => {
this.#childCompletionErrors.push(error);
});
const outputProgress: IChildOutputProgress | undefined =
options.forwardOutput !== false ? { events: 0, pendingWrites: 0 } : undefined;
const completion: Promise<void> = this.#trackChildAsync(child, outputProgress).catch(
(error: unknown) => {
this.#childCompletionErrors.push(error);
}
);
const trackedChild: ITrackedChild = { completion };
this.#trackedChildren.add(trackedChild);
void completion.then(() => this.#trackedChildren.delete(trackedChild));
if (options.forwardInput === true) {
this.#attachChildInput(child);
}
if (options.forwardOutput !== false) {
this.#forwardChildOutput(child.stdout, 'stdout');
this.#forwardChildOutput(child.stderr, 'stderr');
if (outputProgress) {
this.#forwardChildOutput(child.stdout, 'stdout', outputProgress);
this.#forwardChildOutput(child.stderr, 'stderr', outputProgress);
}
return child;
}
Expand Down Expand Up @@ -339,33 +354,34 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon
this.#abortController.abort(reason);
}

async #trackChildAsync(child: childProcess.ChildProcessWithoutNullStreams): Promise<void> {
const terminateChild = (): void => {
try {
SubprocessTerminator.killProcessTree(child, SubprocessTerminator.RECOMMENDED_OPTIONS);
} catch (error) {
this.#childTerminationErrors.push(this.#recordResourceCleanupFailure(error));
try {
child.kill('SIGKILL');
} catch (fallbackError) {
this.#childTerminationErrors.push(this.#recordResourceCleanupFailure(fallbackError));
}
}
};
async #trackChildAsync(
child: childProcess.ChildProcessWithoutNullStreams,
outputProgress: IChildOutputProgress | undefined
): Promise<void> {
const terminateChild = (): void => this.#terminateChild(child);
this.abortSignal.addEventListener('abort', terminateChild, { once: true });
let childError: Error | undefined;
try {
// On POSIX, complete on 'exit' rather than 'close': a background descendant may inherit and hold the output
// pipes open, and the child's process group is killed below. Windows cannot recover descendants after the child
// exits, so the pipes closing remains the only signal that the child's descendants have exited.
await new Promise<void>((resolve) => {
child.once('error', (error: Error) => {
childError = error;
});
if (process.platform !== 'win32') {
child.once('exit', () => resolve());
}
child.once('close', () => resolve());
});
terminateExitedChildProcessGroup(child);
await this.#drainChildOutputAsync(child, outputProgress);
} catch (error) {
throw this.#recordResourceCleanupFailure(error);
} finally {
this.abortSignal.removeEventListener('abort', terminateChild);
}
try {
terminateExitedChildProcessGroup(child);
if (process.platform === 'linux' && child.pid !== undefined) {
await waitForLinuxProcessGroupExitAsync(child.pid);
}
Expand All @@ -375,18 +391,94 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon
if (childError) throw childError;
}

/**
* Terminates the child's whole process tree.
*
* @remarks
* `SubprocessTerminator.killProcessTree` is a no-op once the direct child has exited, which would leave its
* surviving process group running, so an exited child's group is signalled directly.
*/
#terminateChild(child: childProcess.ChildProcessWithoutNullStreams): void {
try {
if (hasChildExited(child)) {
terminateExitedChildProcessGroup(child);
} else {
SubprocessTerminator.killProcessTree(child, SubprocessTerminator.RECOMMENDED_OPTIONS);
}
} catch (error) {
this.#childTerminationErrors.push(this.#recordResourceCleanupFailure(error));
try {
child.kill('SIGKILL');
} catch (fallbackError) {
this.#childTerminationErrors.push(this.#recordResourceCleanupFailure(fallbackError));
}
}
}

/**
* Waits for output that the exited child left in its pipes.
*
* @remarks
* Forwarded output is drained until its pipes close or stay idle for a bounded time; the deadline is extended only
* while this child's own output is still being forwarded, so a slow client does not truncate it and unrelated output
* cannot keep the drain open. An idle pipe, for example one held by a descendant outside the child's process group,
* is then destroyed. Output that the caller consumes itself is not forwarded, so its progress cannot be observed;
* those pipes are awaited until they close or the request aborts.
*/
async #drainChildOutputAsync(
child: childProcess.ChildProcessWithoutNullStreams,
progress: IChildOutputProgress | undefined
): Promise<void> {
if (child.stdout.closed && child.stderr.closed) {
return;
}
let lastEvents: number | undefined = progress?.events;
await new Promise<void>((resolve) => {
let timer: NodeJS.Timeout | undefined;
const finish = (): void => {
clearTimeout(timer);
this.abortSignal.removeEventListener('abort', finish);
resolve();
};
const check = (): void => {
const events: number | undefined = progress?.events;
const progressed: boolean = events !== lastEvents || (progress?.pendingWrites ?? 0) > 0;
if (!progressed || this.abortSignal.aborted) {
finish();
} else {
lastEvents = events;
timer = setTimeout(check, CHILD_OUTPUT_DRAIN_IDLE_TIMEOUT_MS);
}
};
child.once('close', finish);
this.abortSignal.addEventListener('abort', finish, { once: true });
if (this.abortSignal.aborted) {
finish();
} else if (progress) {
timer = setTimeout(check, CHILD_OUTPUT_DRAIN_IDLE_TIMEOUT_MS);
}
});
child.stdout.destroy();
child.stderr.destroy();
}

#recordResourceCleanupFailure(error: unknown): Error {
return recordWorkspaceRequestCleanupFailure(this.workspaceSession, this.#request.requestId, error);
}

#forwardChildOutput(
source: NodeJS.ReadableStream & { pause(): unknown; resume(): unknown },
stream: 'stdout' | 'stderr'
stream: 'stdout' | 'stderr',
progress: IChildOutputProgress
): void {
source.on('data', (chunk: Buffer | string) => {
source.pause();
const bytes: Uint8Array = typeof chunk === 'string' ? Buffer.from(chunk) : chunk;
progress.events++;
progress.pendingWrites++;
void this.#writer.writeAsync(stream, bytes).then(() => {
progress.events++;
progress.pendingWrites--;
if (!this.abortSignal.aborted) {
source.resume();
}
Expand All @@ -402,6 +494,10 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon
}
}

function hasChildExited(child: childProcess.ChildProcess): boolean {
return child.exitCode !== null || child.signalCode !== null;
}

function terminateExitedChildProcessGroup(child: childProcess.ChildProcessWithoutNullStreams): void {
if (process.platform === 'win32' || child.pid === undefined) {
return;
Expand Down
Loading
Loading