diff --git a/common/changes/@rushstack/rush-daemon/rushd-coalesced-early-result_2026-09-24.json b/common/changes/@rushstack/rush-daemon/rushd-coalesced-early-result_2026-09-24.json new file mode 100644 index 0000000000..feecf49b39 --- /dev/null +++ b/common/changes/@rushstack/rush-daemon/rushd-coalesced-early-result_2026-09-24.json @@ -0,0 +1,10 @@ +{ + "changes": [ + { + "packageName": "@rushstack/rush-daemon", + "comment": "Publish a coalesced phased request's result as soon as every operation in its own selection has completed and its output has drained, instead of holding it until the shared iteration finishes the other clients' larger selections.", + "type": "patch" + } + ], + "packageName": "@rushstack/rush-daemon" +} diff --git a/libraries/rush-daemon/README.md b/libraries/rush-daemon/README.md index 713b69bec1..5e65060db2 100644 --- a/libraries/rush-daemon/README.md +++ b/libraries/rush-daemon/README.md @@ -119,8 +119,10 @@ The native Rush lock is held only during graph preparation and each coalesced it is idle. `acquireExecutionLeaseAsync` is an optional engine/session hook invoked once by the batch coordinator, before input reconciliation. Compatible clients share that lease rather than contending independently. It remains held through operation execution, runner cleanup, and every participant's output/input cleanup; the batch barrier -releases it before any final command result is published. Thus ordinary native actions and permanent `--no-daemon` -fallback can run immediately after a completed warm request without stopping the daemon. +releases it before the final command result of the batch's last participant is published. A coalesced participant +whose own operations all completed earlier may receive its result while the iteration still runs for the others (see +below); it must not assume the lock is already released. Thus ordinary native actions and permanent `--no-daemon` +fallback can run immediately after a completed single-client warm request without stopping the daemon. A real native command holding the lock causes preparation or execution to be refused; there is no lock bypass or automatic retry. A later explicit request can retry after contention ends, including contention during the first @@ -442,8 +444,13 @@ request's immutable `RUSH_ALLOW_WARNINGS_IN_SUCCESSFUL_BUILD` environment overri Compatible phased `SHARED-BUILD` requests admitted before the next graph iteration starts are coalesced at a deterministic event-loop-turn boundary. The router reconciles retained invalidations once, unions the clients' enabled dependency closures, and schedules one iteration. Shared operations execute once, while each client subscribes only -to its own closure and derives its final result only from that subset. Requests admitted after scheduling starts form -a later batch. Cancelling or disconnecting one client removes its subscription without aborting work needed by other +to its own closure and derives its final result only from that subset. A client does not wait for the other clients' +larger selections: once every operation of its own closure that the iteration scheduled has completed and its output +has drained, its result is published while the iteration, graph lease, and native execution lease continue for the +remaining clients. The last client that still needs the iteration receives its result after iteration end and lease +release, as for a single client. An early result is not published when any of the client's operations was aborted; +iteration-wide failures that occur after an early result are reported only to the remaining clients. Requests +admitted after scheduling starts form a later batch. Cancelling or disconnecting one client removes its subscription without aborting work needed by other clients; the graph iteration is aborted only after every client in that batch has stopped needing it. The typed phased router remains separate from native initialization. `ProductionDaemonRequestResolver` supplies diff --git a/libraries/rush-daemon/src/PhasedRequestEventSink.ts b/libraries/rush-daemon/src/PhasedRequestEventSink.ts index d25b936da4..d4a315526f 100644 --- a/libraries/rush-daemon/src/PhasedRequestEventSink.ts +++ b/libraries/rush-daemon/src/PhasedRequestEventSink.ts @@ -6,10 +6,10 @@ import { randomUUID } from 'node:crypto'; import type { IOperationExecutionResult, Operation, - OperationStatus, _IOperationActivityOptions, _IOperationGraphEventSink } from '@microsoft/rush-lib'; +import { OperationStatus } from '@microsoft/rush-lib'; import { DAEMON_PROTOCOL_VERSION, RUSHD_OPERATION_HEADER, @@ -28,6 +28,17 @@ import type { IPhasedRequestClient } from './PhasedRequestClient'; const EVENT_SOURCE_PACKAGE: string = '@microsoft/rush-lib'; const EVENT_SOURCE_COMPONENT: string = 'OperationGraph'; const TEXT_ENCODER: InstanceType = new TextEncoder(); +// Mirrors rush-lib's TERMINAL_STATUSES, which is not part of its public API. +const TERMINAL_OPERATION_STATUSES: ReadonlySet = new Set([ + OperationStatus.Success, + OperationStatus.SuccessWithWarning, + OperationStatus.Skipped, + OperationStatus.Blocked, + OperationStatus.FromCache, + OperationStatus.Failure, + OperationStatus.NoOp, + OperationStatus.Aborted +]); interface IObservedOperationResult { readonly executionResult: IOperationExecutionResult; @@ -91,7 +102,10 @@ export class PhasedRequestEventSink implements _IOperationGraphEventSink { readonly #observedResults: Map = new Map(); readonly #rushVersion: string; readonly #writer: OrderedClientWriter; + readonly #onActiveOperationsSettled: (() => void) | undefined; + readonly #pendingOperationIds: Set = new Set(); #completedOperations: number = 0; + #settled: boolean = false; #totalOperations: number = 0; public constructor(options: { @@ -100,10 +114,17 @@ export class PhasedRequestEventSink implements _IOperationGraphEventSink { getNextSequence: () => number; onWriteFailure: (error: Error) => void; rushVersion: string; + /** + * Called at most once per iteration, when every operation of this client's selection that the iteration + * scheduled has emitted its terminal completion event. All of those operations' events and log chunks are + * enqueued on this sink's writer before the callback runs. + */ + onActiveOperationsSettled?: () => void; }) { this.#activeOperationIds = options.activeOperationIds; this.#client = options.client; this.#getNextSequence = options.getNextSequence; + this.#onActiveOperationsSettled = options.onActiveOperationsSettled; this.#rushVersion = options.rushVersion; this.#writer = new OrderedClientWriter(options.client, options.onWriteFailure); } @@ -125,10 +146,34 @@ export class PhasedRequestEventSink implements _IOperationGraphEventSink { public onIterationScheduled(records: Iterable): void { this.#completedOperations = 0; this.#totalOperations = 0; + this.#pendingOperationIds.clear(); + this.#settled = false; for (const record of records) { - if (this.#activeOperationIds.has(record.operation.name) && !record.silent) { + const operationId: string = record.operation.name; + if (!this.#activeOperationIds.has(operationId)) { + continue; + } + if (!record.silent) { this.#totalOperations++; } + if (!TERMINAL_OPERATION_STATUSES.has(record.status)) { + this.#pendingOperationIds.add(operationId); + } + } + } + + public onOperationCompleted(result: IOperationExecutionResult): void { + if (!this.#pendingOperationIds.delete(result.operation.name) || this.#settled) { + return; + } + if (result.status === OperationStatus.Aborted) { + // The iteration is being aborted or failed to start; leave this client's result to the batch. + this.#settled = true; + return; + } + if (this.#pendingOperationIds.size === 0) { + this.#settled = true; + this.#onActiveOperationsSettled?.(); } } diff --git a/libraries/rush-daemon/src/PhasedRequestRouter.ts b/libraries/rush-daemon/src/PhasedRequestRouter.ts index 0903eaf136..9ba5d73c63 100644 --- a/libraries/rush-daemon/src/PhasedRequestRouter.ts +++ b/libraries/rush-daemon/src/PhasedRequestRouter.ts @@ -74,6 +74,11 @@ interface IBatchEntry extends IPreparedPhasedRequest { abortRequested: boolean; completed: boolean; executionStarted: boolean; + /** + * Set when this entry's result starts being produced, so it is produced exactly once. An entry can finish while + * its batch's iteration is still running for other participants; see `#finishSettledEntry`. + */ + finishPromise: Promise | undefined; outputError: unknown; participated: boolean; reject: (error: unknown) => void; @@ -259,6 +264,7 @@ class PhasedRequestBatchCoordinator { abortRequested: false, completed: false, executionStarted: false, + finishPromise: undefined, outputError: undefined, participated: false, reject, @@ -407,6 +413,7 @@ class PhasedRequestBatchCoordinator { client: entry.client, getNextSequence: () => entry.client.getNextEventSequence(), onWriteFailure: (error: Error) => this.#deactivateEntry(entry, false, error), + onActiveOperationsSettled: () => this.#finishSettledEntry(entry), rushVersion: this.#workspaceSession.metadata.rushVersion }); entry.unsubscribe = this.#multiplexer.subscribe(entry.requestSink); @@ -433,7 +440,7 @@ class PhasedRequestBatchCoordinator { }) ); const executionPromise: Promise = this.#graph.executeScheduledIterationAsync(); - if (!participants.some((entry: IBatchEntry) => this.#isEntryLive(entry))) { + if (!participants.some((entry: IBatchEntry) => this.#needsIteration(entry))) { // Let executeScheduledIterationAsync promote the scheduled iteration before aborting it. await Promise.resolve(); this.#requestIterationAbort(); @@ -509,7 +516,7 @@ class PhasedRequestBatchCoordinator { entry.executionStarted && this.#currentBatch && (this.#graph.hasScheduledIteration || this.#graph.status === OperationStatus.Executing) && - !this.#currentBatch.some((candidate: IBatchEntry) => this.#isEntryLive(candidate)) + !this.#currentBatch.some((candidate: IBatchEntry) => this.#needsIteration(candidate)) ) { this.#requestIterationAbort(); } @@ -519,6 +526,44 @@ class PhasedRequestBatchCoordinator { return !entry.abortRequested && !entry.client.abortSignal.aborted && entry.outputError === undefined; } + /** Whether a live participant still waits for the running iteration to produce its result. */ + #needsIteration(entry: IBatchEntry): boolean { + return entry.finishPromise === undefined && this.#isEntryLive(entry); + } + + /** + * Publishes a coalesced participant's result as soon as every operation in its own selection has completed, + * instead of holding it until the shared iteration finishes the other participants' larger selections. + * + * @remarks + * The sink invokes this after the operations' final events and log chunks were enqueued, and `#finishEntryAsync` + * drains them before writing the result. The iteration, graph lease and execution lease stay owned by the batch. + * The last participant that needs the iteration keeps the ordinary contract: its result follows iteration end + * and execution lease release, so single-client requests and warm-state retention are unchanged. + */ + #finishSettledEntry(entry: IBatchEntry): void { + if ( + !this.#needsIteration(entry) || + !entry.participated || + !this.#currentBatch?.some( + (candidate: IBatchEntry) => candidate !== entry && this.#needsIteration(candidate) + ) + ) { + return; + } + entry.unsubscribe?.(); + entry.unsubscribe = undefined; + entry.finishPromise = this.#produceResultAsync(entry, true, undefined, [], undefined, true).catch( + (error: unknown) => { + // Unlike a batch-wide failure, an early result's failure concerns only this client. + if (!entry.completed) { + this.#completeEntry(entry); + entry.reject(error); + } + } + ); + } + #requestIterationAbort(): void { const abortPromise: Promise = this.#graph.abortCurrentIterationAsync(); this.#abortTail = Promise.all([this.#abortTail, abortPromise]) @@ -528,12 +573,31 @@ class PhasedRequestBatchCoordinator { }); } - async #finishEntryAsync( + #finishEntryAsync( entry: IBatchEntry, batchScheduled: boolean, executionError: unknown, batchCleanupErrors: ReadonlyArray = [], beforeResultAsync?: () => Promise + ): Promise { + entry.finishPromise ??= this.#produceResultAsync( + entry, + batchScheduled, + executionError, + batchCleanupErrors, + beforeResultAsync, + false + ); + return entry.finishPromise; + } + + async #produceResultAsync( + entry: IBatchEntry, + batchScheduled: boolean, + executionError: unknown, + batchCleanupErrors: ReadonlyArray, + beforeResultAsync: (() => Promise) | undefined, + iterationInProgress: boolean ): Promise { if (entry.completed) { return; @@ -558,7 +622,8 @@ class PhasedRequestBatchCoordinator { entry.selection.activeOperations, this.#graph, entry.requestSink, - aborted && entry.participated + aborted && entry.participated, + iterationInProgress ) : []; const result: IDaemonPhasedRequestResult = createPhasedCommandResult({ @@ -581,6 +646,13 @@ class PhasedRequestBatchCoordinator { } async #rejectEntryAsync(entry: IBatchEntry, error: unknown): Promise { + if (entry.finishPromise) { + try { + await entry.finishPromise; + } catch { + // An interrupted result is replaced by the batch failure below. + } + } if (entry.completed) { return; } @@ -611,7 +683,8 @@ function createBatchReleaseBarrier( batch: ReadonlyArray, releaseAsync: () => Promise ): () => Promise { - let remaining: number = batch.filter((entry) => !entry.completed).length; + // Entries that already started their result (e.g. published early) never arrive at the barrier. + let remaining: number = batch.filter((entry) => !entry.completed && !entry.finishPromise).length; if (remaining === 0) return releaseAsync; let arrive: () => void = () => undefined; const allDrained: Promise = new Promise((resolve) => { @@ -851,7 +924,8 @@ function collectOperationOutcomes( activeOperations: ReadonlyArray, graph: IOperationGraph, requestSink: PhasedRequestEventSink, - fillMissingAsAborted: boolean = false + fillMissingAsAborted: boolean = false, + iterationInProgress: boolean = false ): ReadonlyArray { const outcomes: IPhasedOperationOutcome[] = []; for (const operation of [...activeOperations].sort(compareOperations)) { @@ -862,7 +936,10 @@ function collectOperationOutcomes( let errorMessage: string | undefined; if ( observed !== undefined && - (retained === undefined || OBSERVED_STATUS_OVERRIDES_RETAINED.has(observed.status)) + // While the iteration still runs, retained results may predate this iteration. + (iterationInProgress || + retained === undefined || + OBSERVED_STATUS_OVERRIDES_RETAINED.has(observed.status)) ) { status = observed.status; errorMessage = observed.executionResult.error?.message; diff --git a/libraries/rush-daemon/src/test/PhasedRequestBatching.test.ts b/libraries/rush-daemon/src/test/PhasedRequestBatching.test.ts index 9ed2adef1c..81ef14a48f 100644 --- a/libraries/rush-daemon/src/test/PhasedRequestBatching.test.ts +++ b/libraries/rush-daemon/src/test/PhasedRequestBatching.test.ts @@ -8,7 +8,7 @@ import type { IDaemonPhasedRequest, IDaemonPhasedRequestResult } from '@rushstack/rush-daemon-protocol'; -import { RUSHD_OPERATION_HEADER } from '@rushstack/rush-daemon-protocol'; +import { RUSHD_OPERATION_HEADER, RUSHD_OPERATION_STREAM_CLOSED } from '@rushstack/rush-daemon-protocol'; import { OperationStatus } from '@microsoft/rush-lib'; import { PhasedRequestRouter } from '../PhasedRequestRouter'; @@ -18,7 +18,7 @@ import { TestPhasedRequestClient, createRoutingFixture } from './PhasedRequestRouterTestUtilities'; -import type { ITestRoutingFixture } from './PhasedRequestRouterTestUtilities'; +import type { ITestClientWrite, ITestRoutingFixture } from './PhasedRequestRouterTestUtilities'; const OPERATION_A: string = 'project-a (_phase:test)'; const OPERATION_B: string = 'project-b (_phase:test)'; @@ -57,6 +57,7 @@ function createRequest( function createFixture(options?: { readonly actionAAsync?: (terminal: ITerminal) => Promise; + readonly actionBAsync?: (terminal: ITerminal) => Promise; readonly actionCAsync?: (terminal: ITerminal) => Promise; readonly statusA?: OperationStatus; }): ITestRoutingFixture { @@ -70,7 +71,7 @@ function createFixture(options?: { options?.actionAAsync ) ], - [OPERATION_B, new TestOperationRunner(OPERATION_B)], + [OPERATION_B, new TestOperationRunner(OPERATION_B, OperationStatus.Success, options?.actionBAsync)], [OPERATION_C, new TestOperationRunner(OPERATION_C, OperationStatus.Success, options?.actionCAsync)] ]), [[OPERATION_B, OPERATION_A]] @@ -81,6 +82,45 @@ function getResultOperationIds(result: IDaemonPhasedRequestResult): ReadonlyArra return result.operationResults.map(({ operationId }) => operationId); } +interface IExecutionLeaseTracker { + readonly events: string[]; +} + +function trackExecutionLease(fixture: ITestRoutingFixture): IExecutionLeaseTracker { + const events: string[] = []; + fixture.session.acquireExecutionLeaseAsync = async (): Promise => { + events.push('acquired'); + return { + [Symbol.asyncDispose]: async (): Promise => { + events.push('released'); + } + }; + }; + return { events }; +} + +function trackResult( + resultPromise: Promise, + label: string, + events: string[] +): Promise { + return resultPromise.then((result: IDaemonPhasedRequestResult) => { + events.push(`result:${label}`); + return result; + }); +} + +function isStreamClosedEvent(write: ITestClientWrite, operationId: string): boolean { + const payload: unknown = write.event?.payload; + return ( + typeof payload === 'object' && + payload !== null && + (payload as { name?: unknown }).name === RUSHD_OPERATION_STREAM_CLOSED && + write.event !== undefined && + eventOperationId(write.event) === operationId + ); +} + function eventOperationId(event: IDaemonEventEnvelope): string | undefined { if (event.scope?.operationId) { return event.scope.operationId; @@ -145,6 +185,131 @@ describe('shared phased request batching', () => { ]); }); + it('publishes a coalesced client result as soon as its own closure settles', async () => { + const releaseB: IDeferred = createDeferred(); + const fixture: ITestRoutingFixture = createFixture({ + actionAAsync: async (terminal: ITerminal): Promise => terminal.writeLine('from-a'), + actionBAsync: async (): Promise => releaseB.promise + }); + const lease: IExecutionLeaseTracker = trackExecutionLease(fixture); + const scheduleSpy: jest.SpyInstance = jest.spyOn(fixture.graph, 'scheduleIterationAsync'); + const router: PhasedRequestRouter = new PhasedRequestRouter(fixture.session); + const clientA: TestPhasedRequestClient = new TestPhasedRequestClient('one'); + const clientB: TestPhasedRequestClient = new TestPhasedRequestClient('two'); + // A slow reader must still receive all of its operation output before its early result. + clientA.onWriteAsync = async (): Promise => { + await new Promise((resolve) => setImmediate(resolve)); + }; + + const resultAPromise: Promise = trackResult( + router.executeAsync(createRequest('a', OPERATION_A), clientA), + 'a', + lease.events + ); + const resultBPromise: Promise = trackResult( + router.executeAsync(createRequest('b', OPERATION_B), clientB), + 'b', + lease.events + ); + + const resultA: IDaemonPhasedRequestResult = await resultAPromise; + expect(lease.events).toEqual(['acquired', 'result:a']); + expect(fixture.graph.status).toBe(OperationStatus.Executing); + expect(resultA).toMatchObject({ exitCode: 0, outcome: 'success', scheduled: true }); + expect(resultA.operationResults).toEqual([ + expect.objectContaining({ operationId: OPERATION_A, status: OperationStatus.Success }) + ]); + expect(getWrittenOperationIds(clientA)).toEqual(new Set([OPERATION_A])); + expect(clientA.writes.some((write: ITestClientWrite) => write.text?.includes('from-a'))).toBe(true); + expect(clientA.writes.findIndex((write) => isStreamClosedEvent(write, OPERATION_A))).toBeGreaterThan(-1); + expect(clientA.writes[clientA.writes.length - 1]?.result).toBe(resultA); + expect(getHeaderData(clientA)).toEqual([ + { completedOperations: 1, operationId: OPERATION_A, totalOperations: 1 } + ]); + const clientAWriteCount: number = clientA.writes.length; + + releaseB.resolve(); + const resultB: IDaemonPhasedRequestResult = await resultBPromise; + // The last participant keeps the ordinary contract: its result follows execution lease release. + expect(lease.events).toEqual(['acquired', 'result:a', 'released', 'result:b']); + expect(scheduleSpy).toHaveBeenCalledTimes(1); + expect(fixture.runners.get(OPERATION_A)?.runCount).toBe(1); + expect(fixture.runners.get(OPERATION_B)?.runCount).toBe(1); + expect(resultB).toMatchObject({ exitCode: 0, outcome: 'success' }); + expect(getResultOperationIds(resultB)).toEqual([OPERATION_A, OPERATION_B]); + expect(clientA.writes).toHaveLength(clientAWriteCount); + }); + + it('publishes an early failure result while the batch continues for other clients', async () => { + const releaseC: IDeferred = createDeferred(); + const fixture: ITestRoutingFixture = createFixture({ + actionCAsync: async (): Promise => releaseC.promise, + statusA: OperationStatus.Failure + }); + fixture.graph.parallelism = 2; + const lease: IExecutionLeaseTracker = trackExecutionLease(fixture); + const router: PhasedRequestRouter = new PhasedRequestRouter(fixture.session); + + const failedPromise: Promise = trackResult( + router.executeAsync(createRequest('failed', OPERATION_A), new TestPhasedRequestClient('one')), + 'failed', + lease.events + ); + const continuingPromise: Promise = trackResult( + router.executeAsync(createRequest('continuing', OPERATION_C), new TestPhasedRequestClient('two')), + 'continuing', + lease.events + ); + + const failed: IDaemonPhasedRequestResult = await failedPromise; + expect(lease.events).toEqual(['acquired', 'result:failed']); + expect(failed).toMatchObject({ aborted: false, exitCode: 1, outcome: 'failure' }); + expect(failed.operationResults).toEqual([ + expect.objectContaining({ operationId: OPERATION_A, status: OperationStatus.Failure }) + ]); + + releaseC.resolve(); + const continuing: IDaemonPhasedRequestResult = await continuingPromise; + expect(lease.events).toEqual(['acquired', 'result:failed', 'released', 'result:continuing']); + expect(continuing).toMatchObject({ exitCode: 0, outcome: 'success' }); + expect(getResultOperationIds(continuing)).toEqual([OPERATION_C]); + }); + + it('aborts the iteration when the only client still needing it cancels after an early result', async () => { + const operationCStarted: IDeferred = createDeferred(); + const releaseC: IDeferred = createDeferred(); + const fixture: ITestRoutingFixture = createFixture({ + actionCAsync: async (): Promise => { + operationCStarted.resolve(); + await releaseC.promise; + } + }); + fixture.graph.parallelism = 2; + const abortSpy: jest.SpyInstance = jest.spyOn(fixture.graph, 'abortCurrentIterationAsync'); + const router: PhasedRequestRouter = new PhasedRequestRouter(fixture.session); + const cancelledClient: TestPhasedRequestClient = new TestPhasedRequestClient('two'); + + const finishedPromise: Promise = router.executeAsync( + createRequest('finished', OPERATION_A), + new TestPhasedRequestClient('one') + ); + const cancelledPromise: Promise = router.executeAsync( + createRequest('cancelled', OPERATION_C), + cancelledClient + ); + const finished: IDaemonPhasedRequestResult = await finishedPromise; + await operationCStarted.promise; + expect(finished).toMatchObject({ exitCode: 0, outcome: 'success' }); + const abortCallCountBeforeCancellation: number = abortSpy.mock.calls.length; + + cancelledClient.abortController.abort(); + expect(abortSpy.mock.calls.length).toBeGreaterThan(abortCallCountBeforeCancellation); + releaseC.resolve(); + const cancelled: IDaemonPhasedRequestResult = await cancelledPromise; + + expect(cancelled).toMatchObject({ aborted: true, outcome: 'aborted' }); + }); + it('derives shared and disjoint failure results from each client subset', async () => { const fixture: ITestRoutingFixture = createFixture({ statusA: OperationStatus.Failure }); const router: PhasedRequestRouter = new PhasedRequestRouter(fixture.session); @@ -193,20 +358,27 @@ describe('shared phased request batching', () => { it('unsubscribes one mid-run cancellation without aborting work required by another client', async () => { const operationStarted: IDeferred = createDeferred(); + const operationCStarted: IDeferred = createDeferred(); const releaseOperation: IDeferred = createDeferred(); const fixture: ITestRoutingFixture = createFixture({ actionAAsync: async (): Promise => { operationStarted.resolve(); await releaseOperation.promise; + }, + // Keep the continuing client's work outstanding until after the cancellation. + actionCAsync: async (): Promise => { + operationCStarted.resolve(); + await releaseOperation.promise; } }); + fixture.graph.parallelism = 2; const cancelledClient: TestPhasedRequestClient = new TestPhasedRequestClient('one'); const continuingClient: TestPhasedRequestClient = new TestPhasedRequestClient('two'); const abortSpy: jest.SpyInstance = jest.spyOn(fixture.graph, 'abortCurrentIterationAsync'); const router: PhasedRequestRouter = new PhasedRequestRouter(fixture.session); const cancelled = router.executeAsync(createRequest('cancelled', OPERATION_A), cancelledClient); const continuing = router.executeAsync(createRequest('continuing', OPERATION_C), continuingClient); - await operationStarted.promise; + await Promise.all([operationStarted.promise, operationCStarted.promise]); const abortCallCountBeforeCancellation: number = abortSpy.mock.calls.length; cancelledClient.abortController.abort(); diff --git a/libraries/rush-daemon/src/test/PhasedRequestEventSink.test.ts b/libraries/rush-daemon/src/test/PhasedRequestEventSink.test.ts index ac5bfdfe06..bb2bcb191d 100644 --- a/libraries/rush-daemon/src/test/PhasedRequestEventSink.test.ts +++ b/libraries/rush-daemon/src/test/PhasedRequestEventSink.test.ts @@ -2,12 +2,14 @@ // See LICENSE in the project root for license information. import type { IDaemonEventEnvelope } from '@rushstack/rush-daemon-protocol'; +import { OperationStatus, type IOperationExecutionResult } from '@microsoft/rush-lib'; import { PhasedRequestEventSink } from '../PhasedRequestEventSink'; import { TestPhasedRequestClient } from './PhasedRequestRouterTestUtilities'; const ACTIVE_OPERATION: string = 'project-a (_phase:test)'; const OTHER_OPERATION: string = 'project-b (_phase:test)'; +const SECOND_ACTIVE_OPERATION: string = 'project-c (_phase:test)'; function createSink(client: TestPhasedRequestClient): PhasedRequestEventSink { return new PhasedRequestEventSink({ @@ -44,3 +46,51 @@ it('forwards unscoped and active activity while filtering other operation activi ]); expect(activities.every(({ required }) => required)).toBe(true); }); + +function createRecord(operationId: string, status: OperationStatus): IOperationExecutionResult { + return { operation: { name: operationId }, silent: false, status } as unknown as IOperationExecutionResult; +} + +function createSettlingSink(onSettled: () => void): PhasedRequestEventSink { + const client: TestPhasedRequestClient = new TestPhasedRequestClient(); + return new PhasedRequestEventSink({ + activeOperationIds: new Set([ACTIVE_OPERATION, SECOND_ACTIVE_OPERATION]), + client, + getNextSequence: () => client.getNextEventSequence(), + onActiveOperationsSettled: onSettled, + onWriteFailure: () => undefined, + rushVersion: '5.178.1' + }); +} + +it('reports settlement once, after every scheduled active operation completed', () => { + const onSettled: jest.Mock = jest.fn(); + const sink: PhasedRequestEventSink = createSettlingSink(onSettled); + sink.onIterationScheduled([ + createRecord(ACTIVE_OPERATION, OperationStatus.Waiting), + createRecord(SECOND_ACTIVE_OPERATION, OperationStatus.Ready), + createRecord(OTHER_OPERATION, OperationStatus.Waiting) + ]); + + sink.onOperationCompleted(createRecord(ACTIVE_OPERATION, OperationStatus.Success)); + sink.onOperationCompleted(createRecord(OTHER_OPERATION, OperationStatus.Success)); + expect(onSettled).not.toHaveBeenCalled(); + sink.onOperationCompleted(createRecord(SECOND_ACTIVE_OPERATION, OperationStatus.Failure)); + sink.onOperationCompleted(createRecord(SECOND_ACTIVE_OPERATION, OperationStatus.Failure)); + + expect(onSettled).toHaveBeenCalledTimes(1); +}); + +it('does not report settlement when an active operation was aborted', () => { + const onSettled: jest.Mock = jest.fn(); + const sink: PhasedRequestEventSink = createSettlingSink(onSettled); + sink.onIterationScheduled([ + createRecord(ACTIVE_OPERATION, OperationStatus.Waiting), + createRecord(SECOND_ACTIVE_OPERATION, OperationStatus.Waiting) + ]); + + sink.onOperationCompleted(createRecord(ACTIVE_OPERATION, OperationStatus.Aborted)); + sink.onOperationCompleted(createRecord(SECOND_ACTIVE_OPERATION, OperationStatus.Success)); + + expect(onSettled).not.toHaveBeenCalled(); +}); \ No newline at end of file diff --git a/libraries/rush-daemon/src/test/PhasedRequestRouterTestUtilities.ts b/libraries/rush-daemon/src/test/PhasedRequestRouterTestUtilities.ts index c61dc8e33d..41f17440d7 100644 --- a/libraries/rush-daemon/src/test/PhasedRequestRouterTestUtilities.ts +++ b/libraries/rush-daemon/src/test/PhasedRequestRouterTestUtilities.ts @@ -194,6 +194,7 @@ export class TestRoutingWorkspaceSession implements IWorkspaceSession { public readonly rushSession: RushSession | undefined = undefined; public readonly operationGraph: IOperationGraph; public onReconcileAsync: (() => Promise) | undefined; + public acquireExecutionLeaseAsync: (() => Promise) | undefined; public constructor(operationGraph: IOperationGraph) { this.operationGraph = operationGraph;