diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index 7a3c971bc4..d1ca442f13 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -39,14 +39,14 @@ export interface Coordinator { /** One Session's process-local execution lane: one active demand and at most one coalesced follow-up. */ type Entry = { readonly done: Deferred.Deferred - readonly settled: Deferred.Deferred> + readonly settled: Deferred.Deferred | undefined> current: Demand pending?: Demand explicitWaiter?: Deferred.Deferred interruptSeq?: number owner?: Fiber.Fiber stopping: boolean - advisoryRetriesRemaining: number + advisoryRetryAvailable: boolean } /** Combines follow-up demand: runs dominate, while wakes retain the newest durable admission sequence. */ @@ -82,17 +82,16 @@ export const make = (options: { }), ) - const makeEntry = ( - current: Demand, - explicitWaiter?: Deferred.Deferred, - advisoryRetriesRemaining = current._tag === "wake" ? 1 : 0, - ): Entry => ({ + const makeEntry = (current: Demand, options?: { + readonly explicitWaiter?: Deferred.Deferred + readonly advisoryRetryAvailable?: boolean + }): Entry => ({ done: Deferred.makeUnsafe(), - settled: Deferred.makeUnsafe>(), + settled: Deferred.makeUnsafe | undefined>(), current, - explicitWaiter, + explicitWaiter: options?.explicitWaiter, stopping: false, - advisoryRetriesRemaining, + advisoryRetryAvailable: options?.advisoryRetryAvailable ?? current._tag === "wake", }) const start = (key: Key, entry: Entry, demand: Demand, successor = false) => { @@ -138,7 +137,7 @@ export const make = (options: { const pending = entry.pending entry.pending = undefined entry.current = pending - entry.advisoryRetriesRemaining = pending._tag === "wake" ? 1 : 0 + entry.advisoryRetryAvailable = pending._tag === "wake" start(key, entry, pending, true) return } @@ -150,15 +149,16 @@ export const make = (options: { const successor = entry.pending !== undefined - ? makeEntry(entry.pending, entry.explicitWaiter) - : exit._tag === "Failure" && demand._tag === "wake" && !entry.stopping && entry.advisoryRetriesRemaining > 0 - ? makeEntry(demand, entry.explicitWaiter, entry.advisoryRetriesRemaining - 1) + ? makeEntry(entry.pending, { explicitWaiter: entry.explicitWaiter }) + : exit._tag === "Failure" && demand._tag === "wake" && !entry.stopping && entry.advisoryRetryAvailable + ? makeEntry(demand, { explicitWaiter: entry.explicitWaiter, advisoryRetryAvailable: false }) : undefined + const retrying = successor !== undefined && entry.pending === undefined if (successor === undefined) active.delete(key) else active.set(key, successor) if (successor !== undefined) start(key, successor, successor.current, true) Deferred.doneUnsafe(entry.done, exit) - Deferred.doneUnsafe(entry.settled, Effect.succeed(exit)) + Deferred.doneUnsafe(entry.settled, Effect.succeed(retrying ? undefined : exit)) if ( exit._tag === "Failure" && !(entry.stopping && Cause.hasInterruptsOnly(exit.cause)) && @@ -197,7 +197,7 @@ export const make = (options: { Deferred.await(shutdown).pipe(Effect.as(Exit.void)), ) if (closed) break - if (exit._tag === "Failure" && firstFailure === undefined) firstFailure = exit.cause + if (exit?._tag === "Failure" && firstFailure === undefined) firstFailure = exit.cause } if (firstFailure !== undefined) return yield* Effect.failCause(firstFailure) }) diff --git a/packages/core/test/session-run-coordinator.test.ts b/packages/core/test/session-run-coordinator.test.ts index 31fc40642e..b30901e431 100644 --- a/packages/core/test/session-run-coordinator.test.ts +++ b/packages/core/test/session-run-coordinator.test.ts @@ -750,9 +750,10 @@ describe("SessionRunCoordinator", () => { yield* coordinator.wake("session", 1) yield* coordinator.wake("session", 2) yield* Deferred.succeed(firstGate, undefined) - yield* coordinator.awaitIdle("session").pipe(Effect.exit) + const idle = yield* coordinator.awaitIdle("session").pipe(Effect.exit) - expect(runs).toBe(3) + expect(runs).toBe(3) + expect(Exit.isSuccess(idle)).toBeTrue() }), ), )