diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index 5068b913ad8..7a3c971bc4c 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -46,6 +46,7 @@ type Entry = { interruptSeq?: number owner?: Fiber.Fiber stopping: boolean + advisoryRetriesRemaining: number } /** Combines follow-up demand: runs dominate, while wakes retain the newest durable admission sequence. */ @@ -81,12 +82,17 @@ export const make = (options: { }), ) - const makeEntry = (current: Demand, explicitWaiter?: Deferred.Deferred): Entry => ({ + const makeEntry = ( + current: Demand, + explicitWaiter?: Deferred.Deferred, + advisoryRetriesRemaining = current._tag === "wake" ? 1 : 0, + ): Entry => ({ done: Deferred.makeUnsafe(), settled: Deferred.makeUnsafe>(), current, explicitWaiter, stopping: false, + advisoryRetriesRemaining, }) const start = (key: Key, entry: Entry, demand: Demand, successor = false) => { @@ -132,6 +138,7 @@ export const make = (options: { const pending = entry.pending entry.pending = undefined entry.current = pending + entry.advisoryRetriesRemaining = pending._tag === "wake" ? 1 : 0 start(key, entry, pending, true) return } @@ -141,7 +148,12 @@ export const make = (options: { return } - const successor = entry.pending !== undefined ? makeEntry(entry.pending, entry.explicitWaiter) : undefined + 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) + : undefined if (successor === undefined) active.delete(key) else active.set(key, successor) if (successor !== undefined) start(key, successor, successor.current, true) @@ -151,6 +163,7 @@ export const make = (options: { 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))) diff --git a/packages/core/test/session-run-coordinator.test.ts b/packages/core/test/session-run-coordinator.test.ts index 39fb5779d3a..31fc40642ea 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,35 @@ 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) + yield* coordinator.awaitIdle("session").pipe(Effect.exit) + + expect(runs).toBe(3) + }), + ), + ) + it.effect("upgrades an active wake when an explicit run joins it", () => Effect.scoped( Effect.gen(function* () { @@ -916,8 +951,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 +963,7 @@ describe("SessionRunCoordinator", () => { yield* Effect.yieldNow expect(reported).toHaveLength(1) + expect(runs).toBe(2) expect(Cause.squash(reported[0]!)).toBe(failure) }), ),