From f520a49bdb02c48055feacf2281498d8c478036b Mon Sep 17 00:00:00 2001 From: selarkin Date: Wed, 23 Sep 2026 20:03:58 -0700 Subject: [PATCH 1/4] [rush-daemon] Run rushx scripts without exclusive admission and complete on child exit Fixes #6085 Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../rushx-admission-6085_2026-09-24.json | 10 ++ common/reviews/api/rush-daemon.api.md | 4 + .../src/DaemonRequestDispatcher.ts | 1 + .../src/GlobalCommandExecutionContext.ts | 121 +++++++++++++--- .../rush-daemon/src/GlobalCommandRequest.ts | 5 + .../src/GlobalCommandRequestRouter.ts | 35 +++-- .../test/GlobalCommandRequestRouter.test.ts | 130 ++++++++++++++++++ 7 files changed, 273 insertions(+), 33 deletions(-) create mode 100644 common/changes/@rushstack/rush-daemon/rushx-admission-6085_2026-09-24.json diff --git a/common/changes/@rushstack/rush-daemon/rushx-admission-6085_2026-09-24.json b/common/changes/@rushstack/rush-daemon/rushx-admission-6085_2026-09-24.json new file mode 100644 index 0000000000..413e48926e --- /dev/null +++ b/common/changes/@rushstack/rush-daemon/rushx-admission-6085_2026-09-24.json @@ -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" +} diff --git a/common/reviews/api/rush-daemon.api.md b/common/reviews/api/rush-daemon.api.md index bbad30a4b9..9ab111f8a7 100644 --- a/common/reviews/api/rush-daemon.api.md +++ b/common/reviews/api/rush-daemon.api.md @@ -7,6 +7,7 @@ /// 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'; @@ -414,6 +415,8 @@ export interface IResolvedGlobalCommandRequest { // (undocumented) readonly environment: IGlobalCommandEnvironment; // (undocumented) + readonly invocationKind?: DaemonInvocationKind; + // (undocumented) readonly requestId: string; // (undocumented) readonly terminal: IGlobalCommandTerminalProperties; @@ -431,6 +434,7 @@ export interface IResolveGlobalCommandRequestOptions { readonly cwd: string; // (undocumented) readonly environment: Readonly; + readonly invocationKind?: DaemonInvocationKind; // (undocumented) readonly requestId: string; // (undocumented) diff --git a/libraries/rush-daemon/src/DaemonRequestDispatcher.ts b/libraries/rush-daemon/src/DaemonRequestDispatcher.ts index 38f8d53e7f..708821f6b9 100644 --- a/libraries/rush-daemon/src/DaemonRequestDispatcher.ts +++ b/libraries/rush-daemon/src/DaemonRequestDispatcher.ts @@ -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, diff --git a/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts b/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts index bbdfa7ef93..8b4d484d76 100644 --- a/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts +++ b/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts @@ -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. @@ -89,6 +94,10 @@ class OrderedTerminalWriter { this.#onFailure = onFailure; } + public get hasPendingWrites(): boolean { + return this.#pendingByteCount > 0; + } + public write(stream: 'stdout' | 'stderr', chunk: Uint8Array): void { void this.writeAsync(stream, chunk); } @@ -172,6 +181,7 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon readonly #writer: OrderedTerminalWriter; #disposePromise: Promise | undefined; #closed: boolean = false; + #outputProgress: number = 0; #requestAborted: boolean = false; public readonly terminal: ITerminal; @@ -260,16 +270,19 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon windowsHide: options.windowsHide }); SubprocessTerminator.killProcessTreeOnExit(child, SubprocessTerminator.RECOMMENDED_OPTIONS); - const completion: Promise = this.#trackChildAsync(child).catch((error: unknown) => { - this.#childCompletionErrors.push(error); - }); + const outputForwarded: boolean = options.forwardOutput !== false; + const completion: Promise = this.#trackChildAsync(child, outputForwarded).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) { + if (outputForwarded) { this.#forwardChildOutput(child.stdout, 'stdout'); this.#forwardChildOutput(child.stderr, 'stderr'); } @@ -339,33 +352,30 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon this.#abortController.abort(reason); } - async #trackChildAsync(child: childProcess.ChildProcessWithoutNullStreams): Promise { - 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, + outputForwarded: boolean + ): Promise { + const terminateChild = (): void => this.#terminateChild(child); this.abortSignal.addEventListener('abort', terminateChild, { once: true }); let childError: Error | undefined; try { + // Complete on 'exit' rather than 'close': a background descendant may inherit and hold the output pipes open. await new Promise((resolve) => { child.once('error', (error: Error) => { childError = error; }); + child.once('exit', () => resolve()); child.once('close', () => resolve()); }); + terminateExitedChildProcessGroup(child); + await this.#drainChildOutputAsync(child, outputForwarded); + } 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); } @@ -375,6 +385,75 @@ 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 while + * output is still being forwarded, so a slow client does not truncate output. 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, + outputForwarded: boolean + ): Promise { + if (child.stdout.closed && child.stderr.closed) { + return; + } + let lastProgress: number = this.#outputProgress; + await new Promise((resolve) => { + let timer: NodeJS.Timeout | undefined; + const finish = (): void => { + clearTimeout(timer); + this.abortSignal.removeEventListener('abort', finish); + resolve(); + }; + const check = (): void => { + const progressed: boolean = this.#outputProgress !== lastProgress || this.#writer.hasPendingWrites; + if (!progressed || this.abortSignal.aborted) { + finish(); + } else { + lastProgress = this.#outputProgress; + 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 (outputForwarded) { + 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); } @@ -386,7 +465,9 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon source.on('data', (chunk: Buffer | string) => { source.pause(); const bytes: Uint8Array = typeof chunk === 'string' ? Buffer.from(chunk) : chunk; + this.#outputProgress++; void this.#writer.writeAsync(stream, bytes).then(() => { + this.#outputProgress++; if (!this.abortSignal.aborted) { source.resume(); } @@ -402,6 +483,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; diff --git a/libraries/rush-daemon/src/GlobalCommandRequest.ts b/libraries/rush-daemon/src/GlobalCommandRequest.ts index ed396b7dfd..0b41a05f27 100644 --- a/libraries/rush-daemon/src/GlobalCommandRequest.ts +++ b/libraries/rush-daemon/src/GlobalCommandRequest.ts @@ -7,6 +7,7 @@ import * as path from 'node:path'; import { EnvironmentMap } from '@rushstack/node-core-library'; import { validateDaemonRequestAdmissionOptions } from '@rushstack/rush-daemon-protocol'; import type { + DaemonInvocationKind, DaemonRushCommandOrigin, DaemonTerminalRequirement, IDaemonRequestAdmissionOptions @@ -48,6 +49,8 @@ export interface IResolveGlobalCommandRequestOptions { readonly commandOrigin: DaemonRushCommandOrigin; readonly cwd: string; readonly environment: Readonly; + /** Whether the request runs a Rush command or a Rushx package script. Defaults to `rush`. */ + readonly invocationKind?: DaemonInvocationKind; readonly requestId: string; readonly terminal: IGlobalCommandTerminalProperties; } @@ -63,6 +66,7 @@ export interface IResolvedGlobalCommandRequest { readonly commandOrigin: DaemonRushCommandOrigin; readonly cwd: string; readonly environment: IGlobalCommandEnvironment; + readonly invocationKind?: DaemonInvocationKind; readonly requestId: string; readonly terminal: IGlobalCommandTerminalProperties; } @@ -101,6 +105,7 @@ export function resolveGlobalCommandRequest( commandOrigin: options.commandOrigin, cwd, environment: new GlobalCommandEnvironment(options.environment), + invocationKind: options.invocationKind === 'rushx' ? 'rushx' : 'rush', requestId: options.requestId, terminal: resolveTerminalProperties(options.terminal) }); diff --git a/libraries/rush-daemon/src/GlobalCommandRequestRouter.ts b/libraries/rush-daemon/src/GlobalCommandRequestRouter.ts index 256564c846..25af327f13 100644 --- a/libraries/rush-daemon/src/GlobalCommandRequestRouter.ts +++ b/libraries/rush-daemon/src/GlobalCommandRequestRouter.ts @@ -96,20 +96,25 @@ export class GlobalCommandRequestRouter { throw new DaemonRequiresInProcessError(policy); } let admissionController: RequestAdmissionController | undefined; - let lease: IRequestLease; + let lease: IRequestLease | undefined; try { - admissionController = new RequestAdmissionController({ - admission: request.admission, - client, - requestId: request.requestId - }); - lease = await admissionController.acquireAsync( - getWorkspaceRequestScheduler(this.#workspaceSession), - classifyRushCommand({ - commandName: request.commandName, - commandOrigin: request.commandOrigin - }) - ); + // Rushx package scripts do not read or mutate daemon-owned workspace state after resolution, so they bypass + // workspace admission. Native rushx takes no workspace lock either, and holding a scheduler lease for the + // script's whole lifetime would serialize concurrent scripts and block unrelated builds (issue #6085). + if (request.invocationKind !== 'rushx') { + admissionController = new RequestAdmissionController({ + admission: request.admission, + client, + requestId: request.requestId + }); + lease = await admissionController.acquireAsync( + getWorkspaceRequestScheduler(this.#workspaceSession), + classifyRushCommand({ + commandName: request.commandName, + commandOrigin: request.commandOrigin + }) + ); + } } catch (error) { admissionController?.dispose(); return await finishAfterAdmissionErrorAsync(request.requestId, client, interactiveSession, error); @@ -125,10 +130,10 @@ export class GlobalCommandRequestRouter { this.#workspaceSession ); } finally { - lease.release(); + lease?.release(); } } finally { - admissionController.dispose(); + admissionController?.dispose(); } } } diff --git a/libraries/rush-daemon/src/test/GlobalCommandRequestRouter.test.ts b/libraries/rush-daemon/src/test/GlobalCommandRequestRouter.test.ts index dc549b31f1..842c12720c 100644 --- a/libraries/rush-daemon/src/test/GlobalCommandRequestRouter.test.ts +++ b/libraries/rush-daemon/src/test/GlobalCommandRequestRouter.test.ts @@ -466,6 +466,136 @@ describe(GlobalCommandRequestRouter.name, () => { } ); + it('does not hold workspace admission for Rushx package scripts', async () => { + const session: TestWorkspaceSession = new TestWorkspaceSession(TEST_REPO_ROOT); + const router: GlobalCommandRequestRouter = new GlobalCommandRequestRouter(session); + const exclusiveLease = await getWorkspaceRequestScheduler(session).acquireAsync({ + exclusivityClass: RequestExclusivityClass.Exclusive, + noWait: true + }); + const bothStarted = createDeferred(); + let startedCount: number = 0; + const runScriptAsync = (requestId: string): Promise => + router.executeAsync( + router.resolveRequest({ + ...createRequestOptions(requestId, FIRST_CWD, {}, 80), + invocationKind: 'rushx' + }), + async () => { + if (++startedCount === 2) bothStarted.resolve(); + await bothStarted.promise; + return { exitCode: 0 }; + }, + new TestGlobalCommandClient() + ); + try { + const results: IGlobalCommandRequestResult[] = await Promise.all([ + runScriptAsync('script-1'), + runScriptAsync('script-2') + ]); + expect(results.map(({ outcome }) => outcome)).toEqual(['success', 'success']); + await expect( + router.executeAsync( + router.resolveRequest({ + ...createRequestOptions('rush-custom', FIRST_CWD, {}, 80), + admission: { noWait: true } + }), + async () => ({ exitCode: 0 }), + new TestGlobalCommandClient() + ) + ).resolves.toMatchObject({ outcome: 'failure' }); + } finally { + exclusiveLease.release(); + } + }); + + (process.platform === 'win32' ? it.skip : it)( + 'completes on child exit when a background descendant holds the output pipes', + async () => { + const session: TestWorkspaceSession = new TestWorkspaceSession(TEST_REPO_ROOT); + const router: GlobalCommandRequestRouter = new GlobalCommandRequestRouter(session); + const client: TestGlobalCommandClient = new TestGlobalCommandClient(); + // The detached grandchild leaves the child's process group, so only the bounded pipe drain can release it. + const script: string = [ + "const { spawn } = require('node:child_process');", + "const grandchild = spawn(process.execPath, ['-e', 'setTimeout(() => {}, 60000)'],", + " { detached: true, stdio: 'inherit' });", + 'grandchild.unref();', + "process.stdout.write('grandchild=' + grandchild.pid + '\\n');" + ].join('\n'); + const startTime: number = Date.now(); + const result: IGlobalCommandRequestResult = await router.executeAsync( + router.resolveRequest(createRequestOptions('background-descendant', FIRST_CWD, {}, 80)), + async (context) => { + const child = context.spawnChild(process.execPath, ['-e', script]); + const exitCode: number | null = await new Promise((resolve) => + child.once('close', (code: number | null) => resolve(code)) + ); + return { exitCode: exitCode ?? 1 }; + }, + client + ); + const output: string = client.chunks.map(({ text }) => text).join(''); + const grandchildPid: number = Number(/grandchild=(\d+)/.exec(output)?.[1]); + try { + expect(result).toMatchObject({ exitCode: 0, outcome: 'success' }); + expect(Date.now() - startTime).toBeLessThan(30000); + } finally { + if (grandchildPid > 0) { + process.kill(grandchildPid, 'SIGKILL'); + } + } + } + ); + + (process.platform === 'win32' ? it.skip : it)( + 'terminates the process group on cancellation after the direct child exited', + async () => { + const session: TestWorkspaceSession = new TestWorkspaceSession(TEST_REPO_ROOT); + const router: GlobalCommandRequestRouter = new GlobalCommandRequestRouter(session); + const client: TestGlobalCommandClient = new TestGlobalCommandClient(); + const killProcessTreeSpy: jest.SpyInstance = jest.spyOn(SubprocessTerminator, 'killProcessTree'); + const processKillSpy: jest.SpyInstance = jest.spyOn(process, 'kill'); + const script: string = [ + "const { spawn } = require('node:child_process');", + "const grandchild = spawn(process.execPath, ['-e', 'setTimeout(() => {}, 60000)'],", + " { detached: true, stdio: 'inherit' });", + 'grandchild.unref();', + "process.stdout.write('grandchild=' + grandchild.pid + '\\n');" + ].join('\n'); + let childPid: number | undefined; + let output: string = ''; + try { + const result: IGlobalCommandRequestResult = await router.executeAsync( + router.resolveRequest(createRequestOptions('cancel-after-exit', FIRST_CWD, {}, 80)), + async (context) => { + const child = context.spawnChild(process.execPath, ['-e', script], { forwardOutput: false }); + childPid = child.pid; + child.stdout.on('data', (chunk: Buffer) => { + output += chunk.toString(); + }); + await new Promise((resolve) => child.once('exit', () => resolve())); + processKillSpy.mockClear(); + client.abortController.abort(new Error('client cancelled')); + await new Promise((resolve) => child.once('close', () => resolve())); + return { exitCode: 0 }; + }, + client + ); + expect(result).toMatchObject({ aborted: true, outcome: 'aborted' }); + expect(killProcessTreeSpy).not.toHaveBeenCalled(); + expect(processKillSpy).toHaveBeenCalledWith(-(childPid ?? 0), 'SIGKILL'); + } finally { + killProcessTreeSpy.mockRestore(); + processKillSpy.mockRestore(); + const grandchildPid: number = Number(/grandchild=(\d+)/.exec(output)?.[1]); + if (grandchildPid > 0) { + process.kill(grandchildPid, 'SIGKILL'); + } + } + } + ); + it('reports child spawn failures during request cleanup', async () => { const session: TestWorkspaceSession = new TestWorkspaceSession(TEST_REPO_ROOT); const router: GlobalCommandRequestRouter = new GlobalCommandRequestRouter(session); From efc5f3ac030117ad71fb110773732ab8b12cd380 Mon Sep 17 00:00:00 2001 From: selarkin Date: Wed, 23 Sep 2026 22:02:21 -0700 Subject: [PATCH 2/4] [rush-cli-client] Update rushx tests for admission-free script execution Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../src/test/RushXDaemonAlias.test.ts | 29 ++++--------------- .../src/test/RushXDaemonBoundaries.test.ts | 25 +++++----------- .../src/test/pipedInput.test.ts | 14 ++++----- 3 files changed, 20 insertions(+), 48 deletions(-) diff --git a/apps/rush-cli-client/src/test/RushXDaemonAlias.test.ts b/apps/rush-cli-client/src/test/RushXDaemonAlias.test.ts index 7c6a4361be..b8f81115f6 100644 --- a/apps/rush-cli-client/src/test/RushXDaemonAlias.test.ts +++ b/apps/rush-cli-client/src/test/RushXDaemonAlias.test.ts @@ -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({ @@ -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); diff --git a/apps/rush-cli-client/src/test/RushXDaemonBoundaries.test.ts b/apps/rush-cli-client/src/test/RushXDaemonBoundaries.test.ts index 451373461d..ecb58c4df2 100644 --- a/apps/rush-cli-client/src/test/RushXDaemonBoundaries.test.ts +++ b/apps/rush-cli-client/src/test/RushXDaemonBoundaries.test.ts @@ -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 = () => {}; @@ -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; diff --git a/apps/rush-cli-client/src/test/pipedInput.test.ts b/apps/rush-cli-client/src/test/pipedInput.test.ts index aac77051f4..0de06c1cc3 100644 --- a/apps/rush-cli-client/src/test/pipedInput.test.ts +++ b/apps/rush-cli-client/src/test/pipedInput.test.ts @@ -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 = new Promise((resolve) => { @@ -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; From 3500d1e70565925fd152f4111cfe37818ddb310f Mon Sep 17 00:00:00 2001 From: selarkin Date: Wed, 23 Sep 2026 22:31:57 -0700 Subject: [PATCH 3/4] [rush-daemon] Keep close-based child completion on Windows Windows cannot recover a child's descendants after it exits, so closing pipes remain the only signal that native install/update workers have finished. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../rush-daemon/src/GlobalCommandExecutionContext.ts | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts b/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts index 8b4d484d76..51bfc273e2 100644 --- a/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts +++ b/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts @@ -360,12 +360,16 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon this.abortSignal.addEventListener('abort', terminateChild, { once: true }); let childError: Error | undefined; try { - // Complete on 'exit' rather than 'close': a background descendant may inherit and hold the output pipes open. + // 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((resolve) => { child.once('error', (error: Error) => { childError = error; }); - child.once('exit', () => resolve()); + if (process.platform !== 'win32') { + child.once('exit', () => resolve()); + } child.once('close', () => resolve()); }); terminateExitedChildProcessGroup(child); From baa4b06687d386d4daebd33ad55238299acf1556 Mon Sep 17 00:00:00 2001 From: selarkin Date: Thu, 24 Sep 2026 11:05:49 -0700 Subject: [PATCH 4/4] [rush-daemon] Scope post-exit output drain progress to each child Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../src/GlobalCommandExecutionContext.ts | 55 +++++++++++-------- .../test/GlobalCommandRequestRouter.test.ts | 47 ++++++++++++++++ 2 files changed, 78 insertions(+), 24 deletions(-) diff --git a/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts b/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts index 51bfc273e2..c2cffbb35a 100644 --- a/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts +++ b/libraries/rush-daemon/src/GlobalCommandExecutionContext.ts @@ -94,10 +94,6 @@ class OrderedTerminalWriter { this.#onFailure = onFailure; } - public get hasPendingWrites(): boolean { - return this.#pendingByteCount > 0; - } - public write(stream: 'stdout' | 'stderr', chunk: Uint8Array): void { void this.writeAsync(stream, chunk); } @@ -169,6 +165,12 @@ interface ITrackedChild { readonly completion: Promise; } +/** 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; @@ -181,7 +183,6 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon readonly #writer: OrderedTerminalWriter; #disposePromise: Promise | undefined; #closed: boolean = false; - #outputProgress: number = 0; #requestAborted: boolean = false; public readonly terminal: ITerminal; @@ -270,8 +271,9 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon windowsHide: options.windowsHide }); SubprocessTerminator.killProcessTreeOnExit(child, SubprocessTerminator.RECOMMENDED_OPTIONS); - const outputForwarded: boolean = options.forwardOutput !== false; - const completion: Promise = this.#trackChildAsync(child, outputForwarded).catch( + const outputProgress: IChildOutputProgress | undefined = + options.forwardOutput !== false ? { events: 0, pendingWrites: 0 } : undefined; + const completion: Promise = this.#trackChildAsync(child, outputProgress).catch( (error: unknown) => { this.#childCompletionErrors.push(error); } @@ -282,9 +284,9 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon if (options.forwardInput === true) { this.#attachChildInput(child); } - if (outputForwarded) { - 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; } @@ -354,7 +356,7 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon async #trackChildAsync( child: childProcess.ChildProcessWithoutNullStreams, - outputForwarded: boolean + outputProgress: IChildOutputProgress | undefined ): Promise { const terminateChild = (): void => this.#terminateChild(child); this.abortSignal.addEventListener('abort', terminateChild, { once: true }); @@ -373,7 +375,7 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon child.once('close', () => resolve()); }); terminateExitedChildProcessGroup(child); - await this.#drainChildOutputAsync(child, outputForwarded); + await this.#drainChildOutputAsync(child, outputProgress); } catch (error) { throw this.#recordResourceCleanupFailure(error); } finally { @@ -417,19 +419,20 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon * 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 while - * output is still being forwarded, so a slow client does not truncate output. 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. + * 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, - outputForwarded: boolean + progress: IChildOutputProgress | undefined ): Promise { if (child.stdout.closed && child.stderr.closed) { return; } - let lastProgress: number = this.#outputProgress; + let lastEvents: number | undefined = progress?.events; await new Promise((resolve) => { let timer: NodeJS.Timeout | undefined; const finish = (): void => { @@ -438,11 +441,12 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon resolve(); }; const check = (): void => { - const progressed: boolean = this.#outputProgress !== lastProgress || this.#writer.hasPendingWrites; + const events: number | undefined = progress?.events; + const progressed: boolean = events !== lastEvents || (progress?.pendingWrites ?? 0) > 0; if (!progressed || this.abortSignal.aborted) { finish(); } else { - lastProgress = this.#outputProgress; + lastEvents = events; timer = setTimeout(check, CHILD_OUTPUT_DRAIN_IDLE_TIMEOUT_MS); } }; @@ -450,7 +454,7 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon this.abortSignal.addEventListener('abort', finish, { once: true }); if (this.abortSignal.aborted) { finish(); - } else if (outputForwarded) { + } else if (progress) { timer = setTimeout(check, CHILD_OUTPUT_DRAIN_IDLE_TIMEOUT_MS); } }); @@ -464,14 +468,17 @@ export class GlobalCommandExecutionContext implements IGlobalCommandExecutionCon #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; - this.#outputProgress++; + progress.events++; + progress.pendingWrites++; void this.#writer.writeAsync(stream, bytes).then(() => { - this.#outputProgress++; + progress.events++; + progress.pendingWrites--; if (!this.abortSignal.aborted) { source.resume(); } diff --git a/libraries/rush-daemon/src/test/GlobalCommandRequestRouter.test.ts b/libraries/rush-daemon/src/test/GlobalCommandRequestRouter.test.ts index 842c12720c..4fe768425a 100644 --- a/libraries/rush-daemon/src/test/GlobalCommandRequestRouter.test.ts +++ b/libraries/rush-daemon/src/test/GlobalCommandRequestRouter.test.ts @@ -538,6 +538,7 @@ describe(GlobalCommandRequestRouter.name, () => { const output: string = client.chunks.map(({ text }) => text).join(''); const grandchildPid: number = Number(/grandchild=(\d+)/.exec(output)?.[1]); try { + expect(grandchildPid).toBeGreaterThan(0); expect(result).toMatchObject({ exitCode: 0, outcome: 'success' }); expect(Date.now() - startTime).toBeLessThan(30000); } finally { @@ -548,6 +549,52 @@ describe(GlobalCommandRequestRouter.name, () => { } ); + (process.platform === 'win32' ? it.skip : it)( + "does not extend an exited child's drain with another child's output", + async () => { + const session: TestWorkspaceSession = new TestWorkspaceSession(TEST_REPO_ROOT); + const router: GlobalCommandRequestRouter = new GlobalCommandRequestRouter(session); + const client: TestGlobalCommandClient = new TestGlobalCommandClient(); + const heldScript: string = [ + "const { spawn } = require('node:child_process');", + "const grandchild = spawn(process.execPath, ['-e', 'setTimeout(() => {}, 60000)'],", + " { detached: true, stdio: 'inherit' });", + 'grandchild.unref();', + "process.stdout.write('grandchild=' + grandchild.pid + '\\n');" + ].join('\n'); + let heldCloseMs: number | undefined; + const result: IGlobalCommandRequestResult = await router.executeAsync( + router.resolveRequest(createRequestOptions('unrelated-output', FIRST_CWD, {}, 80)), + async (context) => { + const chatty = context.spawnChild(process.execPath, [ + '-e', + "setInterval(() => process.stdout.write('.'), 20)" + ]); + const startTime: number = Date.now(); + const held = context.spawnChild(process.execPath, ['-e', heldScript]); + await new Promise((resolve) => held.once('close', () => resolve())); + heldCloseMs = Date.now() - startTime; + chatty.kill('SIGKILL'); + await new Promise((resolve) => chatty.once('close', () => resolve())); + return { exitCode: 0 }; + }, + client + ); + const output: string = client.chunks.map(({ text }) => text).join(''); + const grandchildPid: number = Number(/grandchild=(\d+)/.exec(output)?.[1]); + try { + expect(grandchildPid).toBeGreaterThan(0); + expect(result).toMatchObject({ outcome: 'success' }); + expect(heldCloseMs).toBeLessThan(10000); + } finally { + if (grandchildPid > 0) { + process.kill(grandchildPid, 'SIGKILL'); + } + } + }, + 30000 + ); + (process.platform === 'win32' ? it.skip : it)( 'terminates the process group on cancellation after the direct child exited', async () => {