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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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"
}
15 changes: 11 additions & 4 deletions libraries/rush-daemon/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
49 changes: 47 additions & 2 deletions libraries/rush-daemon/src/PhasedRequestEventSink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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<typeof TextEncoder> = new TextEncoder();
// Mirrors rush-lib's TERMINAL_STATUSES, which is not part of its public API.
const TERMINAL_OPERATION_STATUSES: ReadonlySet<OperationStatus> = new Set([
OperationStatus.Success,
OperationStatus.SuccessWithWarning,
OperationStatus.Skipped,
OperationStatus.Blocked,
OperationStatus.FromCache,
OperationStatus.Failure,
OperationStatus.NoOp,
OperationStatus.Aborted
]);

interface IObservedOperationResult {
readonly executionResult: IOperationExecutionResult;
Expand Down Expand Up @@ -91,7 +102,10 @@ export class PhasedRequestEventSink implements _IOperationGraphEventSink {
readonly #observedResults: Map<Operation, IObservedOperationResult> = new Map();
readonly #rushVersion: string;
readonly #writer: OrderedClientWriter;
readonly #onActiveOperationsSettled: (() => void) | undefined;
readonly #pendingOperationIds: Set<string> = new Set();
#completedOperations: number = 0;
#settled: boolean = false;
#totalOperations: number = 0;

public constructor(options: {
Expand All @@ -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);
}
Expand All @@ -125,10 +146,34 @@ export class PhasedRequestEventSink implements _IOperationGraphEventSink {
public onIterationScheduled(records: Iterable<IOperationExecutionResult>): 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?.();
}
}

Expand Down
91 changes: 84 additions & 7 deletions libraries/rush-daemon/src/PhasedRequestRouter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> | undefined;
outputError: unknown;
participated: boolean;
reject: (error: unknown) => void;
Expand Down Expand Up @@ -259,6 +264,7 @@ class PhasedRequestBatchCoordinator {
abortRequested: false,
completed: false,
executionStarted: false,
finishPromise: undefined,
outputError: undefined,
participated: false,
reject,
Expand Down Expand Up @@ -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);
Expand All @@ -433,7 +440,7 @@ class PhasedRequestBatchCoordinator {
})
);
const executionPromise: Promise<boolean> = 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();
Expand Down Expand Up @@ -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();
}
Expand All @@ -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<void> = this.#graph.abortCurrentIterationAsync();
this.#abortTail = Promise.all([this.#abortTail, abortPromise])
Expand All @@ -528,12 +573,31 @@ class PhasedRequestBatchCoordinator {
});
}

async #finishEntryAsync(
#finishEntryAsync(
entry: IBatchEntry,
batchScheduled: boolean,
executionError: unknown,
batchCleanupErrors: ReadonlyArray<unknown> = [],
beforeResultAsync?: () => Promise<void>
): Promise<void> {
entry.finishPromise ??= this.#produceResultAsync(
entry,
batchScheduled,
executionError,
batchCleanupErrors,
beforeResultAsync,
false
);
return entry.finishPromise;
}

async #produceResultAsync(
entry: IBatchEntry,
batchScheduled: boolean,
executionError: unknown,
batchCleanupErrors: ReadonlyArray<unknown>,
beforeResultAsync: (() => Promise<void>) | undefined,
iterationInProgress: boolean
): Promise<void> {
if (entry.completed) {
return;
Expand All @@ -558,7 +622,8 @@ class PhasedRequestBatchCoordinator {
entry.selection.activeOperations,
this.#graph,
entry.requestSink,
aborted && entry.participated
aborted && entry.participated,
iterationInProgress
)
: [];
const result: IDaemonPhasedRequestResult = createPhasedCommandResult({
Expand All @@ -581,6 +646,13 @@ class PhasedRequestBatchCoordinator {
}

async #rejectEntryAsync(entry: IBatchEntry, error: unknown): Promise<void> {
if (entry.finishPromise) {
try {
await entry.finishPromise;
} catch {
// An interrupted result is replaced by the batch failure below.
}
}
if (entry.completed) {
return;
}
Expand Down Expand Up @@ -611,7 +683,8 @@ function createBatchReleaseBarrier(
batch: ReadonlyArray<IBatchEntry>,
releaseAsync: () => Promise<void>
): () => Promise<void> {
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<void> = new Promise<void>((resolve) => {
Expand Down Expand Up @@ -851,7 +924,8 @@ function collectOperationOutcomes(
activeOperations: ReadonlyArray<Operation>,
graph: IOperationGraph,
requestSink: PhasedRequestEventSink,
fillMissingAsAborted: boolean = false
fillMissingAsAborted: boolean = false,
iterationInProgress: boolean = false
): ReadonlyArray<IPhasedOperationOutcome> {
const outcomes: IPhasedOperationOutcome[] = [];
for (const operation of [...activeOperations].sort(compareOperations)) {
Expand All @@ -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;
Expand Down
Loading
Loading