diff --git a/common/changes/@rushstack/rush-client-core/fix-restart-arbitration_2026-09-24-01-45.json b/common/changes/@rushstack/rush-client-core/fix-restart-arbitration_2026-09-24-01-45.json new file mode 100644 index 0000000000..06c96fc736 --- /dev/null +++ b/common/changes/@rushstack/rush-client-core/fix-restart-arbitration_2026-09-24-01-45.json @@ -0,0 +1,10 @@ +{ + "changes": [ + { + "packageName": "@rushstack/rush-client-core", + "comment": "Retry daemon restarts with jittered backoff inside the admission deadline and return a `restartRetriesExhausted` fallback outcome instead of failing when successors keep restarting for other environments.", + "type": "patch" + } + ], + "packageName": "@rushstack/rush-client-core" +} diff --git a/common/changes/@rushstack/rush-daemon-protocol/fix-restart-arbitration_2026-09-24-18-10.json b/common/changes/@rushstack/rush-daemon-protocol/fix-restart-arbitration_2026-09-24-18-10.json new file mode 100644 index 0000000000..a61beef311 --- /dev/null +++ b/common/changes/@rushstack/rush-daemon-protocol/fix-restart-arbitration_2026-09-24-18-10.json @@ -0,0 +1,10 @@ +{ + "changes": [ + { + "packageName": "@rushstack/rush-daemon-protocol", + "comment": "Document that clients may retry `retryAfterRestart` a bounded number of times within the admission deadline.", + "type": "none" + } + ], + "packageName": "@rushstack/rush-daemon-protocol" +} diff --git a/common/changes/@rushstack/rush-daemon/fix-restart-arbitration_2026-09-24-01-45.json b/common/changes/@rushstack/rush-daemon/fix-restart-arbitration_2026-09-24-01-45.json new file mode 100644 index 0000000000..69802e647c --- /dev/null +++ b/common/changes/@rushstack/rush-daemon/fix-restart-arbitration_2026-09-24-01-45.json @@ -0,0 +1,10 @@ +{ + "changes": [ + { + "packageName": "@rushstack/rush-daemon", + "comment": "Queue a request whose environment requires a process restart until every queued or in-flight request that matches the running process has drained, so concurrent clients with different environments no longer preempt each other.", + "type": "patch" + } + ], + "packageName": "@rushstack/rush-daemon" +} diff --git a/common/reviews/api/rush-client-core.api.md b/common/reviews/api/rush-client-core.api.md index 103df5d039..d24d6380cb 100644 --- a/common/reviews/api/rush-client-core.api.md +++ b/common/reviews/api/rush-client-core.api.md @@ -49,7 +49,7 @@ export type DaemonClientOutcome = { readonly result: IDaemonCommandResult; } | { readonly kind: 'fallback'; - readonly reason: 'unsupported' | 'controllingTerminalRequired' | 'stdinEndUnsupported'; + readonly reason: 'unsupported' | 'controllingTerminalRequired' | 'stdinEndUnsupported' | 'restartRetriesExhausted'; readonly message?: string; } | { readonly kind: 'rejected'; diff --git a/libraries/rush-client-core/README.md b/libraries/rush-client-core/README.md index 6c1afe75e6..19ea090f36 100644 --- a/libraries/rush-client-core/README.md +++ b/libraries/rush-client-core/README.md @@ -24,19 +24,22 @@ or admitting stdin: that is a protocol error, not permission to replay the comma Abort signals send `requestCancel`, then wait for the result; cancellation has a bounded grace period. Disconnects, protocol errors and sink failures are errors, never reasons to replay possibly executed work. Only pre-execution `unsupported`, -`controllingTerminalRequired`, and `stdinEndUnsupported` outcomes permit fallback. Raw-mode changes are +`controllingTerminalRequired`, `stdinEndUnsupported`, and `restartRetriesExhausted` outcomes permit fallback. Raw-mode changes are acknowledged only after applying them. Input listeners and raw state are restored on success, cancellation, disconnect and failure. No resize messages are sent. `executeWithDaemonRestartAsync(readyClient, connectionOptions, executionOptions)` -adds one bounded retry for an explicit `retryAfterRestart: true` result. It captures -the endpoint's PID/start identity before sending, requires protocol 0.10, waits for -that ownership to be released, and reconnects through the same startup mutex. +retries an explicit `retryAfterRestart: true` result a bounded number of times. Before +each hand-off it captures the endpoint's PID/start identity, requires protocol 0.10, +waits for that ownership to be released, and reconnects through the same startup mutex. +Retries after the first use jittered backoff, and the backoff, the successor hand-off +and the resubmitted request all share the request's admission deadline. The original immutable request and unread input are preserved. Output, events, terminal control, or stdin admission forbid retry, as do connection loss and plain -error messages. A second restart result fails explicitly. Cancellation stops waiting -without killing a daemon. Disabling auto-start still permits waiting for a -host-started successor, but never lets the client spawn one. +error messages. When the retry bound or the admission deadline is exhausted, it +returns a `restartRetriesExhausted` fallback outcome so the caller can run in-process. +Cancellation stops waiting without killing a daemon. Disabling auto-start still +permits waiting for a host-started successor, but never lets the client spawn one. `connectOrStartDaemonAsync()` accepts an **explicit, version-selected** executable, arguments, environment and cwd. It does not discover or install a Rush version. diff --git a/libraries/rush-client-core/src/DaemonClient.ts b/libraries/rush-client-core/src/DaemonClient.ts index 012b215cf3..911021c50c 100644 --- a/libraries/rush-client-core/src/DaemonClient.ts +++ b/libraries/rush-client-core/src/DaemonClient.ts @@ -67,7 +67,11 @@ export type DaemonClientOutcome = | { readonly kind: 'result'; readonly result: IDaemonCommandResult } | { readonly kind: 'fallback'; - readonly reason: 'unsupported' | 'controllingTerminalRequired' | 'stdinEndUnsupported'; + readonly reason: + | 'unsupported' + | 'controllingTerminalRequired' + | 'stdinEndUnsupported' + | 'restartRetriesExhausted'; readonly message?: string; } | { readonly kind: 'rejected'; readonly rejection: IDaemonRequestRejectedMessage['payload'] }; diff --git a/libraries/rush-client-core/src/executeWithDaemonRestart.ts b/libraries/rush-client-core/src/executeWithDaemonRestart.ts index 2299d55e27..e0d8563648 100644 --- a/libraries/rush-client-core/src/executeWithDaemonRestart.ts +++ b/libraries/rush-client-core/src/executeWithDaemonRestart.ts @@ -1,6 +1,8 @@ // Copyright (c) Microsoft Corporation. All rights reserved. Licensed under the MIT license. // See LICENSE in the project root for license information. +import { setTimeout as delayAsync } from 'node:timers/promises'; + import { DAEMON_WORKSPACE_RESTART_PROTOCOL_MINOR } from '@rushstack/rush-daemon-protocol'; import { readDaemonLockfile, type IDaemonLockfile } from '@rushstack/rush-daemon-transport'; @@ -10,8 +12,20 @@ import type { DaemonClient, DaemonClientOutcome, IDaemonClientExecuteOptions } f import { DaemonClientError } from './DaemonClientError'; /** - * Executes on a ready client, retrying once only for a typed pre-execution restart. + * The maximum number of successors a single request follows. Each restart serves at least one other + * environment first, so this bounds the wait when several environments share one workspace daemon. + */ +const MAX_RESTART_RETRIES: number = 6; +const RETRY_JITTER_BASE_MS: number = 50; +const RETRY_JITTER_MAX_MS: number = 1000; +/** Matches the default of {@link IConnectOrStartDaemonOptions.startupTimeoutMs}. */ +const DEFAULT_STARTUP_TIMEOUT_MS: number = 15000; + +/** + * Executes on a ready client, retrying only for a typed pre-execution restart. * Preserves the original request, callbacks and unread input; never retries connection loss. + * Restarts are retried with jittered backoff inside the request's admission deadline; once the + * retries or the deadline are exhausted, a `fallback` outcome lets the caller run in-process instead. * The connection options must select the request's expected daemon and startup environment. * @beta */ @@ -25,64 +39,117 @@ export async function executeWithDaemonRestartAsync( execution.abortSignal && connection.abortSignal ? AbortSignal.any([execution.abortSignal, connection.abortSignal]) : (execution.abortSignal ?? connection.abortSignal); + const waitTimeoutMs: number | undefined = execution.request.admission?.waitTimeoutMs; + let owner: IDaemonLockfile | undefined = await attestOwnerAsync(client, connection); + let outcome: DaemonClientOutcome = await client.executeAsync({ ...execution, abortSignal }); + let previous: DaemonClient | undefined; + try { + for (let retry: number = 1; outcome.kind === 'result' && outcome.result.retryAfterRestart; retry++) { + if (!owner) { + throw new DaemonClientError( + 'startupFailed', + 'Cannot attest the restarting daemon ownership; the request was not retried.' + ); + } + if (abortSignal?.aborted) return abortedOutcome(execution); + const getRemainingMs = (): number | undefined => + waitTimeoutMs === undefined ? undefined : waitTimeoutMs - (Date.now() - startedAt); + if (retry > MAX_RESTART_RETRIES || isExpired(getRemainingMs())) { + return restartExhaustedOutcome(retry - 1); + } + let successor: DaemonClient; + let boundedByAdmission: boolean = false; + try { + // The first retry follows the planned successor immediately; later ones back off with jitter so + // clients whose environments differ do not reach each new successor in lockstep. + if (retry > 1) { + await delayAsync(getRetryDelayMs(retry, getRemainingMs()), undefined, { signal: abortSignal }); + } + const remainingMs: number | undefined = getRemainingMs(); + if (isExpired(remainingMs)) return restartExhaustedOutcome(retry - 1); + const startupTimeoutMs: number = connection.startupTimeoutMs ?? DEFAULT_STARTUP_TIMEOUT_MS; + // The successor handoff shares the request's admission deadline rather than starting a fresh one. + boundedByAdmission = remainingMs !== undefined && remainingMs < startupTimeoutMs; + successor = await connectOrStartDaemonAsync({ + ...connection, + startupTimeoutMs: boundedByAdmission ? Math.max(1, Math.ceil(remainingMs!)) : startupTimeoutMs, + previousDaemon: { pid: owner.pid, startedAt: owner.startedAt }, + abortSignal + }); + } catch (error) { + if ( + abortSignal?.aborted && + (error === abortSignal.reason || + (typeof error === 'object' && error !== null && 'code' in error && error.code === 'ABORT_ERR')) + ) { + return abortedOutcome(execution); + } + // A startup error after the admission deadline expired is the deadline, not a new failure mode. + if (boundedByAdmission && error instanceof DaemonClientError && isExpired(getRemainingMs())) { + return restartExhaustedOutcome(retry); + } + throw error; + } finally { + await previous?.closeAsync().catch(() => undefined); + previous = undefined; + } + previous = successor; + const remainingMs: number | undefined = getRemainingMs(); + if (isExpired(remainingMs)) return restartExhaustedOutcome(retry); + owner = await attestOwnerAsync(successor, connection); + outcome = await successor.executeAsync({ + ...execution, + abortSignal, + request: + remainingMs === undefined + ? execution.request + : captureDaemonRequest({ + ...execution.request, + admission: { ...execution.request.admission, waitTimeoutMs: Math.floor(remainingMs) } + }) + }); + } + return outcome; + } finally { + await previous?.closeAsync().catch(() => undefined); + } +} + +function isExpired(remainingMs: number | undefined): boolean { + return remainingMs !== undefined && remainingMs <= 0; +} +/** Returns the published ownership record only when it names the connected, restart-capable process. */ +async function attestOwnerAsync( + client: DaemonClient, + connection: IConnectOrStartDaemonOptions +): Promise { const owner: IDaemonLockfile | undefined = readDaemonLockfile(connection.paths.lockfilePath); const { pid } = await client.status; - const attested: boolean = - client.protocolVersion.minor >= DAEMON_WORKSPACE_RESTART_PROTOCOL_MINOR && + return client.protocolVersion.minor >= DAEMON_WORKSPACE_RESTART_PROTOCOL_MINOR && owner !== undefined && owner.pid === pid && owner.socketPath === connection.paths.socketPath && Number.isSafeInteger(owner.pid) && owner.pid > 0 && - Number.isFinite(Date.parse(owner.startedAt)); - const outcome: DaemonClientOutcome = await client.executeAsync({ ...execution, abortSignal }); - if (outcome.kind !== 'result' || !outcome.result.retryAfterRestart) return outcome; - if (!attested || !owner) { - throw new DaemonClientError( - 'startupFailed', - 'Cannot attest the restarting daemon ownership; the request was not retried.' - ); - } - if (abortSignal?.aborted) return abortedOutcome(execution); - let successor: DaemonClient; - try { - successor = await connectOrStartDaemonAsync({ - ...connection, - previousDaemon: { pid: owner.pid, startedAt: owner.startedAt }, - abortSignal - }); - } catch (error) { - if ( - abortSignal?.aborted && - (error === abortSignal.reason || - (typeof error === 'object' && error !== null && 'code' in error && error.code === 'ABORT_ERR')) - ) { - return abortedOutcome(execution); - } - throw error; - } - const waitTimeoutMs: number | undefined = execution.request.admission?.waitTimeoutMs; - const result: DaemonClientOutcome = await successor.executeAsync({ - ...execution, - abortSignal, - request: - waitTimeoutMs === undefined - ? execution.request - : captureDaemonRequest({ - ...execution.request, - admission: { - ...execution.request.admission, - waitTimeoutMs: Math.max(0, waitTimeoutMs - (Date.now() - startedAt)) - } - }) - }); - if (result.kind === 'result' && result.result.retryAfterRestart) { - throw new DaemonClientError( - 'startupFailed', - 'The successor requested another restart; the single safe retry was exhausted.' - ); - } - return result; + Number.isFinite(Date.parse(owner.startedAt)) + ? owner + : undefined; +} + +function getRetryDelayMs(retry: number, remainingMs: number | undefined): number { + const ceilingMs: number = Math.min(RETRY_JITTER_MAX_MS, RETRY_JITTER_BASE_MS * 2 ** (retry - 1)); + const delayMs: number = Math.floor(ceilingMs / 2 + (Math.random() * ceilingMs) / 2); + return remainingMs === undefined ? delayMs : Math.max(0, Math.min(delayMs, remainingMs - 1)); +} + +function restartExhaustedOutcome(restarts: number): DaemonClientOutcome { + return { + kind: 'fallback', + reason: 'restartRetriesExhausted', + message: `The daemon was still restarting for other environments after ${restarts} ${ + restarts === 1 ? 'restart' : 'restarts' + }; no operation was started by the daemon` + }; } function abortedOutcome(execution: IDaemonClientExecuteOptions): DaemonClientOutcome { @@ -90,4 +157,4 @@ function abortedOutcome(execution: IDaemonClientExecuteOptions): DaemonClientOut kind: 'result', result: { requestId: execution.request.requestId, exitCode: 130, outcome: 'aborted', aborted: true } }; -} +} \ No newline at end of file diff --git a/libraries/rush-client-core/src/test/connectOrStartDaemon.test.ts b/libraries/rush-client-core/src/test/connectOrStartDaemon.test.ts index dc4275cd74..698835c0e3 100644 --- a/libraries/rush-client-core/src/test/connectOrStartDaemon.test.ts +++ b/libraries/rush-client-core/src/test/connectOrStartDaemon.test.ts @@ -380,9 +380,11 @@ describe('detached daemon startup', () => { expect(fs.readFileSync(path.join(folder, 'starts'), 'utf8').trim().split('\n')).toHaveLength(2); }); - it.each(['restart-once', 'restart-always', 'restart-held'])( + it.each(['restart-once', 'restart-twice', 'restart-always', 'restart-held'])( 'retries only the typed pre-execution result for %s after ownership release', async (mode) => { + // restart-always exhausts the deadline; the others need room for loaded CI machines. + const waitTimeoutMs: number = mode === 'restart-always' ? 1000 : 10000; const connection: IConnectOrStartDaemonOptions = { ...options, startCommand: { ...options.startCommand!, args: [...options.startCommand!.args, 'fixture', mode] } @@ -395,7 +397,7 @@ describe('detached daemon startup', () => { cwd: folder, environment: {}, terminal: { isTTY: false, supportsColor: false }, - admission: { waitTimeoutMs: 1000 } + admission: { waitTimeoutMs } }); const pending = executeWithDaemonRestartAsync( client, @@ -405,24 +407,36 @@ describe('detached daemon startup', () => { }, { request } ); - if (mode === 'restart-once') { + if (mode === 'restart-once' || mode === 'restart-twice') { + const restarts: number = mode === 'restart-once' ? 1 : 2; expect(await pending).toMatchObject({ kind: 'result', result: { exitCode: 0 } }); const waits = fs.readFileSync(path.join(folder, 'waits'), 'utf8').trim().split('\n').map(Number); - expect(waits[0]).toBe(1000); - expect(waits[1]).toBeLessThan(1000); - expect(waits[1]).toBeGreaterThanOrEqual(0); - expect(request.admission?.waitTimeoutMs).toBe(1000); - } else { - await expect(pending).rejects.toThrow( - mode === 'restart-held' ? 'previous daemon still owns' : 'single safe retry was exhausted' + expect(waits).toHaveLength(restarts + 1); + expect(waits[0]).toBe(waitTimeoutMs); + for (let index: number = 1; index <= restarts; index++) { + expect(waits[index]).toBeLessThanOrEqual(waits[index - 1]); + expect(waits[index]).toBeLessThan(waitTimeoutMs); + expect(waits[index]).toBeGreaterThanOrEqual(0); + } + expect(request.admission?.waitTimeoutMs).toBe(waitTimeoutMs); + expect(fs.readFileSync(path.join(folder, 'starts'), 'utf8').trim().split('\n')).toHaveLength( + restarts + 1 ); + } else if (mode === 'restart-always') { + // Every successor asks again: bounded retries inside the admission deadline, then a fallback. + expect(await pending).toMatchObject({ kind: 'fallback', reason: 'restartRetriesExhausted' }); + const starts: number = fs.readFileSync(path.join(folder, 'starts'), 'utf8').trim().split('\n').length; + expect(starts).toBeGreaterThanOrEqual(2); + expect(starts).toBeLessThanOrEqual(7); + // The deadline may expire after the last successor started but before the request was resubmitted. + const requests: number = fs.readFileSync(path.join(folder, 'requests'), 'utf8').trim().split('\n').length; + expect(requests).toBeGreaterThanOrEqual(starts - 1); + expect(requests).toBeLessThanOrEqual(starts); + } else { + await expect(pending).rejects.toThrow('previous daemon still owns'); + expect(fs.readFileSync(path.join(folder, 'starts'), 'utf8').trim().split('\n')).toHaveLength(1); + expect(fs.readFileSync(path.join(folder, 'requests'), 'utf8').trim().split('\n')).toHaveLength(1); } - expect(fs.readFileSync(path.join(folder, 'starts'), 'utf8').trim().split('\n')).toHaveLength( - mode === 'restart-held' ? 1 : 2 - ); - expect(fs.readFileSync(path.join(folder, 'requests'), 'utf8').trim().split('\n')).toHaveLength( - mode === 'restart-held' ? 1 : 2 - ); }, 15000 ); @@ -463,6 +477,35 @@ describe('detached daemon startup', () => { } ); + it('bounds the successor hand-off by the admission deadline instead of a fresh startup timeout', async () => { + const connection: IConnectOrStartDaemonOptions = { + ...options, + startupTimeoutMs: 7000, + startCommand: { + ...options.startCommand!, + args: [...options.startCommand!.args, 'fixture', 'restart-held'] + } + }; + const client = await connectOrStartDaemonAsync(connection); + const request = captureDaemonRequest({ + argv: ['test'], + commandName: 'test', + commandOrigin: 'custom', + cwd: folder, + environment: {}, + terminal: { isTTY: false, supportsColor: false }, + admission: { waitTimeoutMs: 500 } + }); + const startedAt: number = Date.now(); + expect(await executeWithDaemonRestartAsync(client, connection, { request })).toMatchObject({ + kind: 'fallback', + reason: 'restartRetriesExhausted' + }); + expect(Date.now() - startedAt).toBeLessThan(5000); + expect(fs.readFileSync(path.join(folder, 'starts'), 'utf8').trim().split('\n')).toHaveLength(1); + expect(fs.readFileSync(path.join(folder, 'requests'), 'utf8').trim().split('\n')).toHaveLength(1); + }); + it('refuses restart retry if ownership was not attested before submitting', async () => { const connection: IConnectOrStartDaemonOptions = { ...options, diff --git a/libraries/rush-client-core/src/test/fixtures/daemon.ts b/libraries/rush-client-core/src/test/fixtures/daemon.ts index 6c00f3b323..e03e1cc12e 100644 --- a/libraries/rush-client-core/src/test/fixtures/daemon.ts +++ b/libraries/rush-client-core/src/test/fixtures/daemon.ts @@ -68,9 +68,13 @@ async function mainAsync(): Promise { if (message.payload.admission?.waitTimeoutMs !== undefined) { fs.appendFileSync(path.join(folder, 'waits'), `${message.payload.admission.waitTimeoutMs}\n`); } + const restartCount: number = fs.existsSync(path.join(folder, 'restarted')) + ? fs.readFileSync(path.join(folder, 'restarted'), 'utf8').length + : 0; const restart: boolean = restartMode !== undefined && - (restartMode !== 'restart-once' || !fs.existsSync(path.join(folder, 'restarted'))); + (restartMode !== 'restart-once' || restartCount < 1) && + (restartMode !== 'restart-twice' || restartCount < 2); await connection.sendFrameAsync({ kind: DaemonFrameType.controlJson, payload: encodeDaemonControlMessage({ @@ -85,7 +89,7 @@ async function mainAsync(): Promise { }) }); if (restart && restartMode !== 'restart-held') { - fs.writeFileSync(path.join(folder, 'restarted'), ''); + fs.appendFileSync(path.join(folder, 'restarted'), 'r'); await stopAsync(); } } else if (message.kind === 'shutdown') { diff --git a/libraries/rush-daemon-protocol/src/DaemonCommandResult.ts b/libraries/rush-daemon-protocol/src/DaemonCommandResult.ts index c3726e2b70..293776da32 100644 --- a/libraries/rush-daemon-protocol/src/DaemonCommandResult.ts +++ b/libraries/rush-daemon-protocol/src/DaemonCommandResult.ts @@ -18,7 +18,8 @@ export type DaemonCommandOutcome = 'success' | 'success-with-warning' | 'failure export interface IDaemonCommandResult { /** * Protocol 0.10: no execution or request IO occurred, and a successor has been selected. - * Retry at most once, after attested predecessor ownership release. Never infer this from an error. + * Retry only after attested predecessor ownership release, within the request's admission deadline and a + * small client-defined retry bound; then fall back instead of retrying. Never infer this from an error. */ readonly retryAfterRestart?: true; /** Whether cancellation or disconnect was observed, even if a cleanup failure determines the outcome. */ diff --git a/libraries/rush-daemon/src/WorkspaceRequestAdmission.ts b/libraries/rush-daemon/src/WorkspaceRequestAdmission.ts index 0863216263..a541976266 100644 --- a/libraries/rush-daemon/src/WorkspaceRequestAdmission.ts +++ b/libraries/rush-daemon/src/WorkspaceRequestAdmission.ts @@ -18,6 +18,7 @@ import { } from './RequestScheduler'; import type { IWorkspaceSession } from './WorkspaceSession'; import { assertWorkspaceRequestResourcesHealthy } from './WorkspaceRequestResources'; +import type { IWorkspaceRestartTicket, WorkspaceRestartArbiter } from './WorkspaceRestartArbiter'; export interface IRequestAdmissionClient { readonly abortSignal: AbortSignal; @@ -191,6 +192,19 @@ export class RequestAdmissionController { } } + /** Waits, within the same admission budget, until a restart would not preempt another request. */ + public async waitForRestartDrainAsync( + arbiter: WorkspaceRestartArbiter, + ticket: IWorkspaceRestartTicket + ): Promise { + // The arbiter reports its own admission errors, so this does not depend on the scheduler error mapping. + await arbiter.waitForDrainAsync(ticket, { + abortSignal: this.#abortController.signal, + noWait: this.#admission?.noWait, + waitTimeoutMs: this.#getRemainingWaitTimeoutMs() + }); + } + public dispose(): void { this.#client.abortSignal.removeEventListener('abort', this.#abortFromClient); } diff --git a/libraries/rush-daemon/src/WorkspaceRequestLifecycle.ts b/libraries/rush-daemon/src/WorkspaceRequestLifecycle.ts index da5ca8230a..adb9753ce6 100644 --- a/libraries/rush-daemon/src/WorkspaceRequestLifecycle.ts +++ b/libraries/rush-daemon/src/WorkspaceRequestLifecycle.ts @@ -49,6 +49,7 @@ import type { IWorkspaceProcessRestartPlan, IWorkspaceSuccessorLaunch } from './WorkspaceProcessRestart'; +import { WorkspaceRestartArbiter, type IWorkspaceRestartTicket } from './WorkspaceRestartArbiter'; interface IExecutionState { began: boolean; @@ -106,6 +107,7 @@ class RestartPendingBeforeExecution extends Error { export class WorkspaceRequestLifecycle implements IDaemonRequestLifecycle { readonly #options: IWorkspaceRequestLifecycleOptions; readonly #gate: RequestScheduler = new RequestScheduler(); + readonly #restartArbiter: WorkspaceRestartArbiter = new WorkspaceRestartArbiter(); readonly #abortController: AbortController = new AbortController(); readonly #observers: Set = new Set(); readonly #terminal: Terminal = new Terminal(new NoOpTerminalProvider()); @@ -195,11 +197,13 @@ export class WorkspaceRequestLifecycle implements IDaemonRequestLifecycle { client, requestId: envelope.requestId }); + // Long-lived observers are cancelled by a transition, so they never delay a restart. + const ticket: IWorkspaceRestartTicket | undefined = observer ? undefined : this.#restartArbiter.enter(); let generation: IPreparedGeneration | undefined; try { for (let attempt: number = 0; ; attempt++) { try { - generation = await this.#prepareAsync(envelope, client, admission); + generation = await this.#prepareAsync(envelope, client, admission, ticket); const requestEnvelope: IDaemonRequestEnvelope = { ...envelope, admission: admission.remainingAdmission @@ -272,6 +276,7 @@ export class WorkspaceRequestLifecycle implements IDaemonRequestLifecycle { } } } finally { + if (ticket) this.#restartArbiter.leave(ticket); admission.dispose(); if (observer) this.#observers.delete(observer); } @@ -281,6 +286,7 @@ export class WorkspaceRequestLifecycle implements IDaemonRequestLifecycle { envelope: IDaemonRequestEnvelope, client: IDaemonRequestDispatchClient, admission: RequestAdmissionController, + ticket: IWorkspaceRestartTicket | undefined, admittedLease?: IRequestLease ): Promise { let lease: IRequestLease = @@ -316,6 +322,11 @@ export class WorkspaceRequestLifecycle implements IDaemonRequestLifecycle { const currentTier: WorkspaceInputChangeTier = this.#classify(current, false); if (currentTier === WorkspaceInputChangeTier.Restart) { lease.release(); + if (ticket) { + // Like build requests, a graph-control restart must not preempt requests this process can serve. + await admission.waitForRestartDrainAsync(this.#restartArbiter, ticket); + if (this.#restartPending) throw new RestartPendingBeforeExecution(); + } this.#cancelObservers(); lease = await admission.acquireAsync(this.#gate, RequestExclusivityClass.Exclusive); session = await this.#options.provider.getSessionAsync(); @@ -408,12 +419,17 @@ export class WorkspaceRequestLifecycle implements IDaemonRequestLifecycle { } lease.release(); + if (tier === WorkspaceInputChangeTier.Restart && ticket) { + // Serve every queued or in-flight request that matches this process before restarting for another one. + await admission.waitForRestartDrainAsync(this.#restartArbiter, ticket); + if (this.#restartPending) throw new RestartPendingBeforeExecution(); + } if (this.#transitioning) { const shared: IRequestLease = await admission.acquireAsync( this.#gate, RequestExclusivityClass.SharedBuild ); - return await this.#prepareAsync(envelope, client, admission, shared); + return await this.#prepareAsync(envelope, client, admission, ticket, shared); } this.#transitioning = ownsTransition = true; this.#cancelObservers(); diff --git a/libraries/rush-daemon/src/WorkspaceRestartArbiter.ts b/libraries/rush-daemon/src/WorkspaceRestartArbiter.ts new file mode 100644 index 0000000000..9fb26f6455 --- /dev/null +++ b/libraries/rush-daemon/src/WorkspaceRestartArbiter.ts @@ -0,0 +1,125 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. Licensed under the MIT license. +// See LICENSE in the project root for license information. + +import { RequestSchedulerError, RequestSchedulerErrorCode } from './RequestScheduler'; + +const MAX_TIMER_DELAY_MS: number = 0x7fffffff; + +/** One dispatched request tracked by a {@link WorkspaceRestartArbiter}. */ +export interface IWorkspaceRestartTicket { + readonly waitingForDrain: boolean; +} + +interface IMutableTicket { + waitingForDrain: boolean; + left: boolean; +} + +/** Options for {@link WorkspaceRestartArbiter.waitForDrainAsync}, supplied by request admission. */ +export interface IWorkspaceRestartDrainOptions { + readonly abortSignal: AbortSignal; + readonly noWait: boolean | undefined; + readonly waitTimeoutMs: number | undefined; +} + +/** + * Arbitrates process restarts between requests whose environments differ from the running daemon. + * A request that needs a restart waits until every other request this process can serve has finished, + * so a mismatched environment never preempts queued or in-flight work that matches the running process. + */ +export class WorkspaceRestartArbiter { + readonly #listeners: Set<() => void> = new Set(); + #servingCount: number = 0; + + /** The number of tracked requests that are not waiting for a restart. */ + public get servingCount(): number { + return this.#servingCount; + } + + public enter(): IWorkspaceRestartTicket { + this.#servingCount++; + const ticket: IMutableTicket = { waitingForDrain: false, left: false }; + return ticket; + } + + public leave(ticket: IWorkspaceRestartTicket): void { + const state: IMutableTicket = ticket as IMutableTicket; + if (state.left) return; + state.left = true; + if (!state.waitingForDrain) this.#decrement(); + } + + /** + * Waits until no other tracked request is still being served by this process, then counts the ticket + * as served again so concurrent restart candidates proceed one at a time. + */ + public async waitForDrainAsync( + ticket: IWorkspaceRestartTicket, + options: IWorkspaceRestartDrainOptions + ): Promise { + const state: IMutableTicket = ticket as IMutableTicket; + if (state.left || state.waitingForDrain) throw new Error('The restart ticket is not being served.'); + state.waitingForDrain = true; + this.#decrement(); + try { + const deadline: number | undefined = + options.waitTimeoutMs === undefined ? undefined : Date.now() + options.waitTimeoutMs; + while (this.#servingCount > 0) { + if (options.noWait) { + throw new RequestSchedulerError( + RequestSchedulerErrorCode.NoWait, + 'Another environment is still being served; the request did not wait for a restart.' + ); + } + await this.#waitForChangeAsync(options.abortSignal, deadline); + } + } finally { + state.waitingForDrain = false; + this.#servingCount++; + } + } + + #decrement(): void { + this.#servingCount--; + if (this.#servingCount === 0) { + for (const listener of Array.from(this.#listeners)) listener(); + } + } + + #waitForChangeAsync(abortSignal: AbortSignal, deadline: number | undefined): Promise { + return new Promise((resolve, reject) => { + let timer: ReturnType | undefined; + const unsubscribe: AbortController = new AbortController(); + const settle = (error?: RequestSchedulerError): void => { + this.#listeners.delete(settle); + unsubscribe.abort(); + if (timer) clearTimeout(timer); + if (error) reject(error); + else resolve(); + }; + const settleAborted = (): void => + settle( + new RequestSchedulerError(RequestSchedulerErrorCode.Aborted, 'The request was aborted before execution.') + ); + if (abortSignal.aborted) { + settleAborted(); + return; + } + this.#listeners.add(settle); + abortSignal.addEventListener('abort', settleAborted, { once: true, signal: unsubscribe.signal }); + if (deadline !== undefined) { + timer = setTimeout( + () => + settle( + new RequestSchedulerError( + RequestSchedulerErrorCode.WaitTimeout, + 'The request was not admitted before the daemon could restart for its environment. ' + + 'Use --wait-timeout or RUSH_DAEMON_QUEUE_TIMEOUT_SECONDS to wait longer.' + ) + ), + Math.min(MAX_TIMER_DELAY_MS, Math.max(0, deadline - Date.now())) + ); + } + }); + } +} diff --git a/libraries/rush-daemon/src/test/WorkspaceRestartArbiter.test.ts b/libraries/rush-daemon/src/test/WorkspaceRestartArbiter.test.ts new file mode 100644 index 0000000000..8ee1001739 --- /dev/null +++ b/libraries/rush-daemon/src/test/WorkspaceRestartArbiter.test.ts @@ -0,0 +1,75 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. Licensed under the MIT license. +// See LICENSE in the project root for license information. + +import { RequestSchedulerError, RequestSchedulerErrorCode } from '../RequestScheduler'; +import { WorkspaceRestartArbiter, type IWorkspaceRestartTicket } from '../WorkspaceRestartArbiter'; + +const WAIT: { abortSignal: AbortSignal; noWait: undefined; waitTimeoutMs: undefined } = { + abortSignal: new AbortController().signal, + noWait: undefined, + waitTimeoutMs: undefined +}; + +async function isSettledAsync(promise: Promise): Promise { + let settled: boolean = false; + void promise.then( + () => (settled = true), + () => (settled = true) + ); + await new Promise((resolve) => setImmediate(resolve)); + return settled; +} + +describe(WorkspaceRestartArbiter.name, () => { + it('proceeds immediately when no other request is being served', async () => { + const arbiter: WorkspaceRestartArbiter = new WorkspaceRestartArbiter(); + const ticket: IWorkspaceRestartTicket = arbiter.enter(); + await arbiter.waitForDrainAsync(ticket, WAIT); + expect(arbiter.servingCount).toBe(1); + arbiter.leave(ticket); + expect(arbiter.servingCount).toBe(0); + }); + + it('waits for served requests, including ones that arrive later, and admits candidates one at a time', async () => { + const arbiter: WorkspaceRestartArbiter = new WorkspaceRestartArbiter(); + const serving: IWorkspaceRestartTicket = arbiter.enter(); + const first: IWorkspaceRestartTicket = arbiter.enter(); + const second: IWorkspaceRestartTicket = arbiter.enter(); + const firstWait: Promise = arbiter.waitForDrainAsync(first, WAIT); + const secondWait: Promise = arbiter.waitForDrainAsync(second, WAIT); + const late: IWorkspaceRestartTicket = arbiter.enter(); + arbiter.leave(serving); + expect(await isSettledAsync(firstWait)).toBe(false); + arbiter.leave(late); + await firstWait; + expect(await isSettledAsync(secondWait)).toBe(false); + arbiter.leave(first); + await secondWait; + arbiter.leave(second); + expect(arbiter.servingCount).toBe(0); + }); + + it.each([ + ['no-wait', RequestSchedulerErrorCode.NoWait], + ['timeout', RequestSchedulerErrorCode.WaitTimeout], + ['abort', RequestSchedulerErrorCode.Aborted] + ])('reports %s as an admission failure and restores its accounting', async (mode, code) => { + const arbiter: WorkspaceRestartArbiter = new WorkspaceRestartArbiter(); + const serving: IWorkspaceRestartTicket = arbiter.enter(); + const ticket: IWorkspaceRestartTicket = arbiter.enter(); + const abort: AbortController = new AbortController(); + const waiting: Promise = arbiter.waitForDrainAsync(ticket, { + abortSignal: abort.signal, + noWait: mode === 'no-wait' ? true : undefined, + waitTimeoutMs: mode === 'timeout' ? 10 : undefined + }); + if (mode === 'abort') abort.abort(); + const error: unknown = await waiting.catch((caught: unknown) => caught); + expect(error).toBeInstanceOf(RequestSchedulerError); + expect((error as RequestSchedulerError).code).toBe(code); + expect(arbiter.servingCount).toBe(2); + arbiter.leave(ticket); + arbiter.leave(serving); + expect(arbiter.servingCount).toBe(0); + }); +}); \ No newline at end of file diff --git a/libraries/rush-daemon/src/test/WorkspaceRestartArbitration.test.ts b/libraries/rush-daemon/src/test/WorkspaceRestartArbitration.test.ts new file mode 100644 index 0000000000..557c54ddc5 --- /dev/null +++ b/libraries/rush-daemon/src/test/WorkspaceRestartArbitration.test.ts @@ -0,0 +1,82 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. Licensed under the MIT license. +// See LICENSE in the project root for license information. + +import * as fs from 'node:fs'; +import * as path from 'node:path'; +import { setTimeout as delayAsync } from 'node:timers/promises'; + +import { WorkspaceInputChangeTier } from '@microsoft/rush-lib'; + +import { getInstalledWorkspaceSuccessorLaunchAsync } from '../WorkspaceProcessRestart'; +import { DaemonGraphTestFixture } from './DaemonGraphTestFixture'; +import type { ITerminalExchange } from './DaemonRequestWireTestUtilities'; +import { pongAsync, setDaemonPolicy } from './WarmGenerationTestUtilities'; +import { stopSuccessorAsync } from './WorkspaceLifecycleTestProcess'; + +jest.setTimeout(60_000); + +it('queues a mismatched-environment restart until matching queued and in-flight requests drain', async () => { + const fixture = await DaemonGraphTestFixture.createAsync((created) => { + setDaemonPolicy(created, {}); + created.getSuccessorLaunchAsync = getInstalledWorkspaceSuccessorLaunchAsync; + // Project c holds its build open until the test removes the marker. + created.write('hold', ''); + created.write( + 'c/build.cjs', + "const fs=require('node:fs');fs.appendFileSync('../runs.txt','c\\n');" + + "const t=setInterval(()=>{if(!fs.existsSync('../hold')){clearInterval(t);console.log('finished-c');}},20);" + ); + }); + const order: string[] = []; + const track = (name: string, exchange: Promise): Promise => + exchange.then((result) => { + order.push(name); + return result; + }); + try { + const before = await pongAsync(fixture); + const matching = track('matching', fixture.runAsync(['build', '--to', 'c', '--parallelism', '3'])); + const deadline: number = Date.now() + 30_000; + while (!fixture.runs().includes('c') && Date.now() < deadline) await delayAsync(20); + expect(fixture.runs()).toContain('c'); + + const mismatched = track( + 'mismatched', + fixture.runAsync(['build', '--to', 'b', '--parallelism', '3'], { + environment: { ...fixture.environment, RUSHD_RELOAD_TIER_TEST: 'changed' } + }) + ); + await delayAsync(1000); + // Arrives after the mismatched request; it must still be served by this process. + const lateMatching = track('late', fixture.runAsync(['build', '--to', 'c', '--parallelism', '3'])); + await delayAsync(1000); + expect(order).toEqual([]); + expect(fixture.host.workspaceStatus.lastReloadTier).not.toBe(WorkspaceInputChangeTier.Restart); + + fs.rmSync(path.join(fixture.folder, 'hold')); + const [held, late, restart] = await Promise.all([matching, lateMatching, mismatched]); + expect(held.terminal).toMatchObject({ kind: 'requestResult', payload: { exitCode: 0 } }); + expect(held.terminal.payload).not.toHaveProperty('retryAfterRestart'); + expect(late.terminal).toMatchObject({ kind: 'requestResult', payload: { exitCode: 0 } }); + expect(late.terminal.payload).not.toHaveProperty('retryAfterRestart'); + expect(restart.terminal).toMatchObject({ + kind: 'requestResult', + payload: { exitCode: 1, retryAfterRestart: true } + }); + expect(order[order.length - 1]).toBe('mismatched'); + + const restarted = await fixture.host.restartCompleted; + expect(restarted?.pid).not.toBe(before.pid); + expect((await pongAsync(fixture)).pid).toBe(restarted?.pid); + expect(fixture.runs()).not.toContain('b'); + } finally { + fs.rmSync(path.join(fixture.folder, 'hold'), { force: true }); + try { + await fixture.host.closeAsync(); + await fixture.host.restartCompleted; + } finally { + await stopSuccessorAsync(fixture.host.paths); + await fixture[Symbol.asyncDispose](); + } + } +}); \ No newline at end of file