diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index 5068b913ad..d1ca442f13 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -39,13 +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 + advisoryRetryAvailable: boolean } /** Combines follow-up demand: runs dominate, while wakes retain the newest durable admission sequence. */ @@ -81,12 +82,16 @@ export const make = (options: { }), ) - const makeEntry = (current: Demand, explicitWaiter?: Deferred.Deferred): 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, + advisoryRetryAvailable: options?.advisoryRetryAvailable ?? current._tag === "wake", }) const start = (key: Key, entry: Entry, demand: Demand, successor = false) => { @@ -132,6 +137,7 @@ export const make = (options: { const pending = entry.pending entry.pending = undefined entry.current = pending + entry.advisoryRetryAvailable = pending._tag === "wake" start(key, entry, pending, true) return } @@ -141,16 +147,23 @@ export const make = (options: { return } - const successor = entry.pending !== undefined ? makeEntry(entry.pending, entry.explicitWaiter) : undefined + const successor = + entry.pending !== undefined + ? 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)) && demand._tag === "wake" && + successor === undefined && options.onFailure !== undefined ) { report(Effect.suspend(() => options.onFailure!(key, exit.cause))) @@ -184,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 39fb5779d3..c4eefc1058 100644 --- a/packages/core/test/session-run-coordinator.test.ts +++ b/packages/core/test/session-run-coordinator.test.ts @@ -33,7 +33,9 @@ describe("SessionRunCoordinator", () => { Effect.scoped( Effect.gen(function* () { const drained = yield* Deferred.make() - const coordinator = yield* SessionRunCoordinator.make({ drain: () => Deferred.succeed(drained, undefined) }) + const coordinator = yield* SessionRunCoordinator.make({ + drain: () => Deferred.succeed(drained, undefined), + }) yield* coordinator.wake("session") yield* Deferred.await(drained) @@ -44,7 +46,9 @@ describe("SessionRunCoordinator", () => { it.effect("does nothing when interrupted while idle", () => Effect.scoped( Effect.gen(function* () { - const coordinator = yield* SessionRunCoordinator.make({ drain: () => Effect.void }) + const coordinator = yield* SessionRunCoordinator.make({ + drain: () => Effect.void, + }) yield* coordinator.interrupt("session") }), @@ -55,7 +59,9 @@ describe("SessionRunCoordinator", () => { Effect.scoped( Effect.gen(function* () { let runs = 0 - const coordinator = yield* SessionRunCoordinator.make({ drain: () => Effect.sync(() => runs++) }) + const coordinator = yield* SessionRunCoordinator.make({ + drain: () => Effect.sync(() => runs++), + }) yield* coordinator.interrupt("session", 2) yield* coordinator.wake("session", 1) @@ -722,6 +728,36 @@ describe("SessionRunCoordinator", () => { ), ) + it.effect("retries a failed advisory successor once without another wake", () => + Effect.scoped( + Effect.gen(function* () { + const firstGate = yield* Deferred.make() + let runs = 0 + const failure = new Error("transient wake failure") + const coordinator = yield* SessionRunCoordinator.make({ + drain: () => + Effect.sync(() => ++runs).pipe( + Effect.flatMap((run) => + run === 1 + ? Deferred.await(firstGate).pipe(Effect.andThen(Effect.fail(failure))) + : run === 2 + ? Effect.fail(failure) + : Effect.void, + ), + ), + }) + + yield* coordinator.wake("session", 1) + yield* coordinator.wake("session", 2) + yield* Deferred.succeed(firstGate, undefined) + const idle = yield* coordinator.awaitIdle("session").pipe(Effect.exit) + + expect(runs).toBe(3) + expect(Exit.isSuccess(idle)).toBeTrue() + }), + ), + ) + it.effect("upgrades an active wake when an explicit run joins it", () => Effect.scoped( Effect.gen(function* () { @@ -916,8 +952,9 @@ describe("SessionRunCoordinator", () => { const failure = new Error("wake failed") const reported: Cause.Cause[] = [] const reportedOnce = yield* Deferred.make() + let runs = 0 const coordinator = yield* SessionRunCoordinator.make({ - drain: () => Effect.fail(failure), + drain: () => Effect.sync(() => runs++).pipe(Effect.andThen(Effect.fail(failure))), onFailure: (_key, cause) => Effect.sync(() => reported.push(cause)).pipe(Effect.andThen(Deferred.succeed(reportedOnce, undefined))), }) @@ -927,6 +964,7 @@ describe("SessionRunCoordinator", () => { yield* Effect.yieldNow expect(reported).toHaveLength(1) + expect(runs).toBe(2) expect(Cause.squash(reported[0]!)).toBe(failure) }), ),