diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index 5068b913ad..5a93cb8612 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -1,268 +1,384 @@ export * as SessionRunCoordinator from "./run-coordinator" -import { Cause, Context, Deferred, Effect, Exit, Fiber, FiberSet, Layer, Scope } from "effect" +import { + Cause, + Context, + Data, + Deferred, + Effect, + Equal, + Exit, + Fiber, + FiberSet, + Layer, + Scope, + SynchronizedRef, +} from "effect" import { SessionRunner } from "./runner" import { SessionSchema } from "./schema" export type Mode = "run" | "wake" -/** Why one drain generation should run. Explicit runs dominate advisory wakes when demands coalesce. */ -type Demand = { readonly _tag: "run" } | { readonly _tag: "wake"; readonly seq?: number } - -/** - * Runs at most one drain chain per key while allowing different keys to drain concurrently. - * - * For each key: - * - * idle --run/wake--> draining --run/wake--> draining + one coalesced rerun --> idle - * - * `run` is an explicit drain request. It starts a chain or joins the current chain and - * upgrades a pending follow-up so the caller receives explicit-run semantics. - * - * `wake` reports that durable work may now be available. It starts a chain while idle or - * requests one coalesced follow-up while draining. Repeated wakes collapse together. - * - * `interrupt` stops the current ownership chain. Advisory wakes from before the interrupt - * boundary are suppressed; advisory wakes after the boundary run after cleanup. - */ export interface Coordinator { - /** Starts or joins one explicit drain generation. */ readonly run: (key: Key) => Effect.Effect - /** Coalesces one wake-up after durable work is recorded. */ readonly wake: (key: Key, seq?: number) => Effect.Effect - /** Waits until the current ownership chain settles. */ readonly awaitIdle: (key: Key) => Effect.Effect - /** Interrupts the active ownership chain without automatically draining pending wakes. */ readonly interrupt: (key: Key, seq?: number) => Effect.Effect } -/** 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> - current: Demand - pending?: Demand - explicitWaiter?: Deferred.Deferred - interruptSeq?: number - owner?: Fiber.Fiber - stopping: boolean +/** @internal */ +export class Demand extends Data.Class<{ + readonly explicit: boolean + readonly wakeSeq?: number + readonly unsequencedWake: boolean +}> { + static readonly empty = new Demand({ explicit: false, wakeSeq: undefined, unsequencedWake: false }) + static readonly run = nonEmpty(new Demand({ explicit: true, wakeSeq: undefined, unsequencedWake: false })) + + static wake(seq?: number) { + return nonEmpty(new Demand({ explicit: false, wakeSeq: seq, unsequencedWake: seq === undefined })) + } + + combine(other: Demand) { + return new Demand({ + explicit: this.explicit || other.explicit, + wakeSeq: + this.wakeSeq === undefined + ? other.wakeSeq + : other.wakeSeq === undefined + ? this.wakeSeq + : Math.max(this.wakeSeq, other.wakeSeq), + unsequencedWake: this.unsequencedWake || other.unsequencedWake, + }) + } + + afterBoundary(boundary?: number) { + return new Demand({ + explicit: false, + wakeSeq: + boundary !== undefined && this.wakeSeq !== undefined && this.wakeSeq > boundary ? this.wakeSeq : undefined, + unsequencedWake: false, + }) + } + + isNonEmpty(): this is NonEmptyDemand { + return this.explicit || this.wakeSeq !== undefined || this.unsequencedWake + } + + get mode(): Mode { + return this.explicit ? "run" : "wake" + } } -/** Combines follow-up demand: runs dominate, while wakes retain the newest durable admission sequence. */ -const coalesce = (left: Demand | undefined, right: Demand): Demand => { - if (left?._tag === "run" || right._tag === "run") return { _tag: "run" } - return { _tag: "wake", seq: maxSeq(left?.seq, right.seq) } +type NonEmptyDemand = Demand & + ({ readonly explicit: true } | { readonly wakeSeq: number } | { readonly unsequencedWake: true }) + +function nonEmpty(demand: Demand): NonEmptyDemand { + if (!demand.isNonEmpty()) throw new Error("Session run demand must not be empty") + return demand } -const maxSeq = (left: number | undefined, right: number | undefined) => { - if (left === undefined) return right - if (right === undefined) return left - return Math.max(left, right) +type Lifecycle = + | { readonly _tag: "Running"; readonly token: object; readonly owner: Deferred.Deferred> } + | { + readonly _tag: "Stopping" + readonly token: object + readonly owner: Deferred.Deferred> + readonly boundary?: number + } + +type Lane = { + readonly current: NonEmptyDemand + readonly pending: Demand + readonly lifecycle: Lifecycle + readonly terminal: Deferred.Deferred> + readonly waiter?: Deferred.Deferred> +} + +type State = { + readonly closed: boolean + readonly lanes: ReadonlyMap> + readonly interruptSeq: ReadonlyMap +} + +type Start = { + readonly key: Key + readonly demand: NonEmptyDemand + readonly successor: boolean + readonly token: object + readonly owner: Deferred.Deferred> + readonly ready: Deferred.Deferred + readonly terminal: Deferred.Deferred> +} + +type RunRequest = + | { readonly _tag: "Closed" } + | { readonly _tag: "Await"; readonly terminal: Deferred.Deferred> } + | { readonly _tag: "Retry"; readonly terminal: Deferred.Deferred> } + | { readonly _tag: "Start"; readonly start: Start; readonly terminal: Deferred.Deferred> } + +type Completion = { + readonly start?: Start + readonly terminal?: Deferred.Deferred> + readonly waiter?: Deferred.Deferred> + readonly report?: Cause.Cause } -/** Constructs a scoped coordinator. Every in-memory transition is synchronous. */ export const make = (options: { readonly drain: (key: Key, mode: Mode) => Effect.Effect readonly onFailure?: (key: Key, cause: Cause.Cause) => Effect.Effect }): Effect.Effect, never, Scope.Scope> => Effect.gen(function* () { - const active = new Map>() - const interruptSeq = new Map() - const report = yield* FiberSet.makeRuntime() + const state = yield* SynchronizedRef.make>({ + closed: false, + lanes: new Map(), + interruptSeq: new Map(), + }) const fork = yield* FiberSet.makeRuntime() const shutdown = Deferred.makeUnsafe() - let closed = false - yield* Effect.addFinalizer(() => - Effect.sync(() => { - closed = true - Deferred.doneUnsafe(shutdown, Effect.void) - active.clear() - interruptSeq.clear() - }), - ) - const makeEntry = (current: Demand, explicitWaiter?: Deferred.Deferred): Entry => ({ - done: Deferred.makeUnsafe(), - settled: Deferred.makeUnsafe>(), - current, - explicitWaiter, - stopping: false, - }) - - const start = (key: Key, entry: Entry, demand: Demand, successor = false) => { - const ready = Deferred.makeUnsafe() - const drain = Effect.suspend(() => options.drain(key, demand._tag)) - // Initial work retains immediate-start behavior but cannot run before ownership is published. - // Observer-started successors yield once so synchronous drains cannot recurse on the JS stack. - const owner = fork( - (successor - ? Effect.yieldNow.pipe(Effect.andThen(drain)) - : Deferred.await(ready).pipe(Effect.andThen(drain)) - ).pipe( - Effect.onExit((exit) => Effect.sync(() => settle(key, entry, demand, exit))), - Effect.exit, - Effect.asVoid, - ), - ) - entry.owner = owner - if (!successor) Deferred.doneUnsafe(ready, Effect.void) + const updateLane = (current: State, key: Key, lane?: Lane): State => { + const lanes = new Map(current.lanes) + if (lane === undefined) lanes.delete(key) + else lanes.set(key, lane) + return { ...current, lanes } } - const settle = (key: Key, entry: Entry, demand: Demand, exit: Exit.Exit) => { - if (closed) { - Deferred.doneUnsafe(entry.done, exit) - Deferred.doneUnsafe(entry.settled, Effect.succeed(exit)) - return + const start = (input: { + readonly state: State + readonly key: Key + readonly demand: NonEmptyDemand + readonly terminal?: Deferred.Deferred> + readonly waiter?: Deferred.Deferred> + readonly successor?: boolean + }) => { + const instruction: Start = { + key: input.key, + demand: input.demand, + successor: input.successor ?? false, + token: {}, + owner: Deferred.makeUnsafe>(), + ready: Deferred.makeUnsafe(), + terminal: input.terminal ?? Deferred.makeUnsafe>(), } - if (demand._tag === "run" && entry.explicitWaiter !== undefined) { - Deferred.doneUnsafe(entry.explicitWaiter, exit) - entry.explicitWaiter = undefined + return { + state: updateLane(input.state, input.key, { + current: input.demand, + pending: Demand.empty, + lifecycle: { _tag: "Running", token: instruction.token, owner: instruction.owner }, + terminal: instruction.terminal, + waiter: input.waiter, + }), + start: instruction, + result: instruction.terminal, } - if (entry.stopping && demand._tag === "wake" && entry.explicitWaiter !== undefined) { - Deferred.doneUnsafe(entry.explicitWaiter, exit) - entry.explicitWaiter = undefined - } - if (active.get(key) !== entry) { - Deferred.doneUnsafe(entry.done, exit) - Deferred.doneUnsafe(entry.settled, Effect.succeed(exit)) - return - } - if (exit._tag === "Success" && !entry.stopping) { - if (entry.pending !== undefined) { - const pending = entry.pending - entry.pending = undefined - entry.current = pending - start(key, entry, pending, true) - return + } + + const launch = (instruction: Start) => + Effect.gen(function* () { + const fiber = fork( + Deferred.await(instruction.ready).pipe( + Effect.andThen(instruction.successor ? Effect.yieldNow : Effect.void), + Effect.andThen(Effect.suspend(() => options.drain(instruction.key, instruction.demand.mode))), + Effect.onExit((exit) => complete(instruction.key, instruction.token, exit)), + Effect.exit, + Effect.asVoid, + ), + ) + yield* Deferred.succeed(instruction.owner, fiber) + yield* Deferred.succeed(instruction.ready, undefined) + }) + + const complete = (key: Key, token: object, exit: Exit.Exit): Effect.Effect => { + return SynchronizedRef.modify(state, (current): readonly [Completion, State] => { + const lane = current.lanes.get(key) + if (lane === undefined || lane.lifecycle.token !== token) return [{}, current] + + const deliberateInterrupt = + lane.lifecycle._tag === "Stopping" && exit._tag === "Failure" && Cause.hasInterruptsOnly(exit.cause) + const report = + exit._tag === "Failure" && !deliberateInterrupt && !lane.current.explicit ? exit.cause : undefined + const completesWaiter = lane.current.explicit || (lane.lifecycle._tag === "Stopping" && !lane.current.explicit) + const waiter = completesWaiter ? undefined : lane.waiter + + if (exit._tag === "Success" && lane.lifecycle._tag === "Running" && lane.pending.isNonEmpty()) { + const next = start({ + state: current, + key, + demand: lane.pending, + terminal: lane.terminal, + waiter, + successor: true, + }) + return [{ start: next.start, waiter: completesWaiter ? lane.waiter : undefined, report }, next.state] } - active.delete(key) - Deferred.doneUnsafe(entry.done, exit) - Deferred.doneUnsafe(entry.settled, Effect.succeed(exit)) - return - } - const successor = entry.pending !== undefined ? makeEntry(entry.pending, entry.explicitWaiter) : 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)) - if ( - exit._tag === "Failure" && - !(entry.stopping && Cause.hasInterruptsOnly(exit.cause)) && - demand._tag === "wake" && - options.onFailure !== undefined - ) { - report(Effect.suspend(() => options.onFailure!(key, exit.cause))) - } + const next = lane.pending.isNonEmpty() + ? start({ state: current, key, demand: lane.pending, waiter, successor: true }) + : { state: updateLane(current, key) } + return [ + { + start: "start" in next ? next.start : undefined, + terminal: lane.terminal, + waiter: completesWaiter ? lane.waiter : undefined, + report, + }, + next.state, + ] + }).pipe(Effect.flatMap((instruction) => executeCompletion(key, exit, instruction))) } + const executeCompletion = (key: Key, exit: Exit.Exit, instruction: Completion) => + Effect.gen(function* () { + if (instruction.start !== undefined) yield* launch(instruction.start) + if (instruction.waiter !== undefined) yield* Deferred.succeed(instruction.waiter, exit) + if (instruction.terminal !== undefined) yield* Deferred.succeed(instruction.terminal, exit) + if (instruction.report !== undefined && options.onFailure !== undefined) { + const onFailure = options.onFailure + const cause = instruction.report + fork(Effect.suspend(() => onFailure(key, cause)).pipe(Effect.exit, Effect.asVoid)) + } + }) + + const awaitTerminal = (terminal: Deferred.Deferred>) => + Effect.raceFirst( + Deferred.await(terminal).pipe( + Effect.flatMap( + Exit.match({ + onSuccess: Effect.succeed, + onFailure: Effect.failCause, + }), + ), + ), + Deferred.await(shutdown).pipe(Effect.andThen(Effect.interrupt)), + ) + + const run = (key: Key): Effect.Effect => + Effect.suspend(() => + Effect.uninterruptibleMask((restore) => { + return SynchronizedRef.modify(state, (current): readonly [RunRequest, State] => { + if (current.closed) return [{ _tag: "Closed" }, current] + const lane = current.lanes.get(key) + if (lane?.lifecycle._tag === "Stopping") return [{ _tag: "Retry", terminal: lane.terminal }, current] + if (lane?.current.explicit) return [{ _tag: "Await", terminal: lane.terminal }, current] + if (lane !== undefined) { + const terminal = lane.waiter ?? Deferred.makeUnsafe>() + const pending = lane.pending.combine(Demand.run) + if (Equal.equals(pending, lane.pending) && lane.waiter !== undefined) + return [{ _tag: "Await", terminal }, current] + return [{ _tag: "Await", terminal }, updateLane(current, key, { ...lane, pending, waiter: terminal })] + } + const next = start({ state: current, key, demand: Demand.run }) + return [{ _tag: "Start", start: next.start, terminal: next.result }, next.state] + }).pipe( + Effect.flatMap((request) => { + if (request._tag === "Closed") return Effect.interrupt + if (request._tag === "Start") + return launch(request.start).pipe(Effect.andThen(awaitTerminal(request.terminal))) + if (request._tag === "Await") return awaitTerminal(request.terminal) + return Effect.raceFirst( + Deferred.await(request.terminal).pipe(Effect.as(true)), + Deferred.await(shutdown).pipe(Effect.as(false)), + ).pipe(Effect.flatMap((retry) => (retry ? run(key) : Effect.interrupt))) + }), + restore, + ) + }), + ) + const wake = (key: Key, seq?: number) => - Effect.sync(() => { - if (closed) return - if (!isAfterInterrupt(key, seq)) return - const entry = active.get(key) - if (entry !== undefined) { - if (!acceptsWake(entry, seq)) return - entry.pending = coalesce(entry.pending, { _tag: "wake", seq }) - return - } + Effect.uninterruptible( + Effect.suspend(() => { + return SynchronizedRef.modify(state, (current): readonly [Start | undefined, State] => { + if (current.closed) return [undefined, current] + const boundary = current.interruptSeq.get(key) + if (boundary !== undefined && (seq === undefined || seq <= boundary)) return [undefined, current] + const lane = current.lanes.get(key) + if (lane === undefined) { + const next = start({ state: current, key, demand: Demand.wake(seq) }) + return [next.start, next.state] + } + if ( + lane.lifecycle._tag === "Stopping" && + (lane.lifecycle.boundary === undefined || seq === undefined || seq <= lane.lifecycle.boundary) + ) + return [undefined, current] + const pending = lane.pending.combine(Demand.wake(seq)) + if (Equal.equals(pending, lane.pending)) return [undefined, current] + return [undefined, updateLane(current, key, { ...lane, pending })] + }).pipe(Effect.flatMap((instruction) => (instruction === undefined ? Effect.void : launch(instruction)))) + }), + ) - const next = makeEntry({ _tag: "wake", seq }) - active.set(key, next) - start(key, next, next.current) - }) + const interrupt = (key: Key, seq?: number) => + Effect.uninterruptible( + SynchronizedRef.modify(state, (current) => { + if (current.closed) return [undefined, current] as const + const latest = current.interruptSeq.get(key) + const lane = current.lanes.get(key) + if (seq !== undefined && latest !== undefined && seq <= latest) + return [lane?.lifecycle._tag === "Stopping" ? lane.lifecycle.owner : undefined, current] as const + + const bounded = (() => { + if (seq === undefined) return current + const interruptSeq = new Map(current.interruptSeq) + interruptSeq.set(key, seq) + return { ...current, interruptSeq } + })() + if (lane === undefined) return [undefined, bounded] as const + if ( + !lane.current.explicit && + seq !== undefined && + lane.current.wakeSeq !== undefined && + lane.current.wakeSeq > seq + ) + return [undefined, bounded] as const + + const pending = lane.current.afterBoundary(seq).combine(lane.pending.afterBoundary(seq)) + const boundary = + lane.lifecycle._tag === "Stopping" && lane.lifecycle.boundary !== undefined && seq !== undefined + ? Math.max(lane.lifecycle.boundary, seq) + : lane.lifecycle._tag === "Stopping" && seq === undefined + ? lane.lifecycle.boundary + : seq + return [ + lane.lifecycle.owner, + updateLane(bounded, key, { + ...lane, + pending, + lifecycle: { _tag: "Stopping", token: lane.lifecycle.token, owner: lane.lifecycle.owner, boundary }, + }), + ] as const + }).pipe( + Effect.flatMap((owner) => + owner === undefined ? Effect.void : Deferred.await(owner).pipe(Effect.flatMap(Fiber.interrupt)), + ), + ), + ) const awaitIdle = (key: Key): Effect.Effect => Effect.gen(function* () { - let firstFailure: Cause.Cause | undefined - while (!closed) { - const entry = active.get(key) - if (entry === undefined) break + let failure: Cause.Cause | undefined + while (true) { + const terminal = (yield* SynchronizedRef.get(state)).lanes.get(key)?.terminal + if (terminal === undefined) break const exit = yield* Effect.raceFirst( - Deferred.await(entry.settled), + Deferred.await(terminal), Deferred.await(shutdown).pipe(Effect.as(Exit.void)), ) - if (closed) break - if (exit._tag === "Failure" && firstFailure === undefined) firstFailure = exit.cause + if (exit._tag === "Failure" && failure === undefined) failure = exit.cause } - if (firstFailure !== undefined) return yield* Effect.failCause(firstFailure) + if (failure !== undefined) return yield* Effect.failCause(failure) }) - const interrupt = (key: Key, seq?: number): Effect.Effect => - Effect.suspend(() => { - const entry = active.get(key) - const latest = interruptSeq.get(key) - if (seq !== undefined && latest !== undefined && seq <= latest) - return entry?.stopping && entry.owner !== undefined ? Fiber.interrupt(entry.owner) : Effect.void - if (seq !== undefined) interruptSeq.set(key, seq) - if (entry?.owner === undefined) return Effect.void - if ( - seq !== undefined && - entry.current._tag === "wake" && - entry.current.seq !== undefined && - entry.current.seq > seq - ) - return Effect.void - if (entry.stopping) { - entry.interruptSeq = maxSeq(entry.interruptSeq, seq) - suppressPendingAtOrBefore(entry, seq) - return Fiber.interrupt(entry.owner) - } - entry.stopping = true - entry.interruptSeq = seq - suppressPendingAtOrBefore(entry, seq) - return Fiber.interrupt(entry.owner) - }) + yield* Effect.addFinalizer(() => + SynchronizedRef.modify(state, (_current) => [ + undefined, + { closed: true, lanes: new Map(), interruptSeq: new Map() } satisfies State, + ]).pipe(Effect.andThen(Deferred.succeed(shutdown, undefined))), + ) - return { run, wake, awaitIdle, interrupt } - - function run(key: Key): Effect.Effect { - return Effect.uninterruptibleMask((restore) => { - if (closed) return Effect.interrupt - const entry = active.get(key) - if (entry !== undefined) { - if (entry.stopping) { - return restore(Deferred.await(entry.settled).pipe(Effect.andThen(run(key)))) - } - if (entry.current._tag === "wake") { - entry.pending = coalesce(entry.pending, { _tag: "run" }) - entry.explicitWaiter ??= Deferred.makeUnsafe() - return restore(awaitRun(entry.explicitWaiter)) - } - return restore(awaitRun(entry.done)) - } - - const next = makeEntry({ _tag: "run" }) - active.set(key, next) - start(key, next, next.current) - return restore(awaitRun(next.done)) - }) - } - - function awaitRun(done: Deferred.Deferred): Effect.Effect { - return Effect.raceFirst(Deferred.await(done), Deferred.await(shutdown).pipe(Effect.andThen(Effect.interrupt))) - } - - function acceptsWake(entry: Entry, seq: number | undefined) { - return !entry.stopping || (entry.interruptSeq !== undefined && seq !== undefined && seq > entry.interruptSeq) - } - - function isAfterInterrupt(key: Key, seq: number | undefined) { - const latest = interruptSeq.get(key) - return latest === undefined || (seq !== undefined && seq > latest) - } - - function suppressPendingAtOrBefore(entry: Entry, seq: number | undefined) { - if ( - entry.pending?._tag === "wake" && - seq !== undefined && - entry.pending.seq !== undefined && - entry.pending.seq > seq - ) - return - entry.pending = undefined - } + return { run, wake, interrupt, awaitIdle } }) export interface Interface extends Coordinator {} diff --git a/packages/core/test/session-run-coordinator.test.ts b/packages/core/test/session-run-coordinator.test.ts index 39fb5779d3..502f6bcc85 100644 --- a/packages/core/test/session-run-coordinator.test.ts +++ b/packages/core/test/session-run-coordinator.test.ts @@ -1,10 +1,34 @@ -import { describe, expect } from "bun:test" -import { Cause, Deferred, Effect, Exit, Fiber, Layer, Scope } from "effect" +import { describe, expect, test } from "bun:test" +import { Cause, Deferred, Effect, Equal, Exit, Fiber, Layer, Scope } from "effect" import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator" import { testEffect } from "./lib/effect" const it = testEffect(Layer.empty) +describe("SessionRunCoordinator.Demand", () => { + const Demand = SessionRunCoordinator.Demand + + test("combines associatively with an identity", () => { + const left = Demand.run.combine(Demand.wake(1)) + const right = Demand.wake().combine(Demand.wake(3)) + + expect(Equal.equals(Demand.empty.combine(left), left)).toBeTrue() + expect(Equal.equals(left.combine(Demand.empty), left)).toBeTrue() + expect(Equal.equals(left.combine(right), right.combine(left))).toBeTrue() + expect( + Equal.equals(left.combine(right).combine(Demand.wake(2)), left.combine(right.combine(Demand.wake(2)))), + ).toBeTrue() + }) + + test("keeps only sequenced wakes newer than an interrupt boundary", () => { + const demand = Demand.run.combine(Demand.wake()).combine(Demand.wake(3)) + + expect(Equal.equals(demand.afterBoundary(2), Demand.wake(3))).toBeTrue() + expect(Equal.equals(demand.afterBoundary(3), Demand.empty)).toBeTrue() + expect(Equal.equals(demand.afterBoundary(), Demand.empty)).toBeTrue() + }) +}) + describe("SessionRunCoordinator", () => { it.effect("joins concurrent resumes for one key", () => Effect.scoped( @@ -29,6 +53,51 @@ describe("SessionRunCoordinator", () => { ), ) + it.effect("allocates fresh ownership when one run effect is reused", () => + Effect.scoped( + Effect.gen(function* () { + let runs = 0 + const coordinator = yield* SessionRunCoordinator.make({ drain: () => Effect.sync(() => runs++) }) + const run = coordinator.run("session") + + yield* run + yield* run + + expect(runs).toBe(2) + }), + ), + ) + + it.effect("captures awaitIdle chains safely while settlement races", () => + Effect.scoped( + Effect.gen(function* () { + const iterations = 500 + const gates = Array.from({ length: iterations }, () => Deferred.makeUnsafe()) + let runs = 0 + const coordinator = yield* SessionRunCoordinator.make({ + drain: () => + Effect.suspend(() => { + const gate = gates[runs++] + return gate === undefined ? Effect.die("Missing test gate") : Deferred.await(gate) + }), + }) + + for (let index = 0; index < iterations; index++) { + const run = yield* coordinator.run("session").pipe(Effect.forkChild) + yield* Effect.yieldNow + const idle = yield* coordinator.awaitIdle("session").pipe(Effect.forkChild({ startImmediately: true })) + const gate = gates[index] + if (gate === undefined) yield* Effect.die("Missing test gate") + yield* Deferred.succeed(gate, undefined) + yield* Effect.all([Fiber.join(run), Fiber.join(idle)]) + } + + expect(runs).toBe(iterations) + return undefined + }), + ), + ) + it.effect("starts a drain when woken while idle", () => Effect.scoped( Effect.gen(function* () { @@ -122,6 +191,106 @@ describe("SessionRunCoordinator", () => { ), ) + it.effect("preserves a newer wake coalesced behind a pending explicit run", () => + Effect.scoped( + Effect.gen(function* () { + const firstStarted = yield* Deferred.make() + const secondStarted = yield* Deferred.make() + const modes: SessionRunCoordinator.Mode[] = [] + const coordinator = yield* SessionRunCoordinator.make({ + drain: (_key, mode) => + Effect.sync(() => modes.push(mode)).pipe( + Effect.flatMap((run) => + run === 1 + ? Deferred.succeed(firstStarted, undefined).pipe(Effect.andThen(Effect.never)) + : Deferred.succeed(secondStarted, undefined), + ), + ), + }) + + yield* coordinator.wake("session", 1) + yield* Deferred.await(firstStarted) + const run = yield* coordinator.run("session").pipe(Effect.exit, Effect.forkChild) + yield* Effect.yieldNow + yield* coordinator.wake("session", 3) + yield* coordinator.interrupt("session", 2) + yield* Deferred.await(secondStarted) + yield* coordinator.awaitIdle("session").pipe(Effect.exit) + + const runExit = yield* Fiber.join(run) + expect(Exit.isFailure(runExit) && Cause.hasInterruptsOnly(runExit.cause)).toBeTrue() + expect(modes).toEqual(["wake", "wake"]) + }), + ), + ) + + it.effect("preserves a newer wake from an interrupted active combined demand", () => + Effect.scoped( + Effect.gen(function* () { + const firstGate = yield* Deferred.make() + const secondStarted = yield* Deferred.make() + const thirdStarted = yield* Deferred.make() + const modes: SessionRunCoordinator.Mode[] = [] + const coordinator = yield* SessionRunCoordinator.make({ + drain: (_key, mode) => + Effect.sync(() => modes.push(mode)).pipe( + Effect.flatMap((run) => { + if (run === 1) return Deferred.await(firstGate) + if (run === 2) return Deferred.succeed(secondStarted, undefined).pipe(Effect.andThen(Effect.never)) + return Deferred.succeed(thirdStarted, undefined) + }), + ), + }) + + yield* coordinator.wake("session", 1) + const run = yield* coordinator.run("session").pipe(Effect.exit, Effect.forkChild) + yield* Effect.yieldNow + yield* coordinator.wake("session", 3) + yield* Deferred.succeed(firstGate, undefined) + yield* Deferred.await(secondStarted) + yield* coordinator.interrupt("session", 2) + yield* Deferred.await(thirdStarted) + yield* coordinator.awaitIdle("session").pipe(Effect.exit) + + const runExit = yield* Fiber.join(run) + expect(Exit.isFailure(runExit) && Cause.hasInterruptsOnly(runExit.cause)).toBeTrue() + expect(modes).toEqual(["wake", "run", "wake"]) + }), + ), + ) + + it.effect("suppresses an older wake from an interrupted active combined demand", () => + Effect.scoped( + Effect.gen(function* () { + const firstGate = yield* Deferred.make() + const secondStarted = yield* Deferred.make() + const modes: SessionRunCoordinator.Mode[] = [] + const coordinator = yield* SessionRunCoordinator.make({ + drain: (_key, mode) => + Effect.sync(() => modes.push(mode)).pipe( + Effect.flatMap((run) => { + if (run === 1) return Deferred.await(firstGate) + return Deferred.succeed(secondStarted, undefined).pipe(Effect.andThen(Effect.never)) + }), + ), + }) + + yield* coordinator.wake("session", 1) + const run = yield* coordinator.run("session").pipe(Effect.exit, Effect.forkChild) + yield* Effect.yieldNow + yield* coordinator.wake("session", 2) + yield* Deferred.succeed(firstGate, undefined) + yield* Deferred.await(secondStarted) + yield* coordinator.interrupt("session", 2) + yield* coordinator.awaitIdle("session").pipe(Effect.exit) + + const runExit = yield* Fiber.join(run) + expect(Exit.isFailure(runExit) && Cause.hasInterruptsOnly(runExit.cause)).toBeTrue() + expect(modes).toEqual(["wake", "run"]) + }), + ), + ) + it.effect("interrupts only the requested key", () => Effect.scoped( Effect.gen(function* () { @@ -316,6 +485,53 @@ describe("SessionRunCoordinator", () => { ), ) + it.effect("does not let an interrupted attempt completion settle its successor", () => + Effect.scoped( + Effect.gen(function* () { + const firstStarted = yield* Deferred.make() + const cleanupStarted = yield* Deferred.make() + const cleanupGate = yield* Deferred.make() + const secondStarted = yield* Deferred.make() + const secondGate = yield* Deferred.make() + const idleSettled = yield* Deferred.make() + let runs = 0 + const coordinator = yield* SessionRunCoordinator.make({ + drain: () => + Effect.sync(() => ++runs).pipe( + Effect.flatMap((run) => + run === 1 + ? Deferred.succeed(firstStarted, undefined).pipe( + Effect.andThen(Effect.never), + Effect.onInterrupt(() => + Deferred.succeed(cleanupStarted, undefined).pipe(Effect.andThen(Deferred.await(cleanupGate))), + ), + ) + : Deferred.succeed(secondStarted, undefined).pipe(Effect.andThen(Deferred.await(secondGate))), + ), + ), + }) + + yield* coordinator.wake("session", 1) + yield* Deferred.await(firstStarted) + const interrupt = yield* coordinator.interrupt("session", 2).pipe(Effect.forkChild) + yield* Deferred.await(cleanupStarted) + yield* coordinator.wake("session", 3) + yield* Deferred.succeed(cleanupGate, undefined) + yield* Fiber.join(interrupt) + yield* Deferred.await(secondStarted) + const idle = yield* coordinator + .awaitIdle("session") + .pipe(Effect.ensuring(Deferred.succeed(idleSettled, undefined)), Effect.forkChild) + + yield* Effect.yieldNow + expect(yield* Deferred.isDone(idleSettled)).toBeFalse() + yield* Deferred.succeed(secondGate, undefined) + yield* Fiber.join(idle) + expect(runs).toBe(2) + }), + ), + ) + it.effect("interrupts an explicit run queued before the interruption request", () => Effect.scoped( Effect.gen(function* () { @@ -847,6 +1063,37 @@ describe("SessionRunCoordinator", () => { }), ) + it.effect("settles a post-stop run waiter when its owning scope closes", () => + Effect.gen(function* () { + const scope = yield* Scope.make() + const started = yield* Deferred.make() + const cleanupStarted = yield* Deferred.make() + const cleanupGate = yield* Deferred.make() + const coordinator = yield* SessionRunCoordinator.make({ + drain: () => + Deferred.succeed(started, undefined).pipe( + Effect.andThen(Effect.never), + Effect.onInterrupt(() => + Deferred.succeed(cleanupStarted, undefined).pipe(Effect.andThen(Deferred.await(cleanupGate))), + ), + ), + }).pipe(Scope.provide(scope)) + + yield* coordinator.wake("session") + yield* Deferred.await(started) + const interrupt = yield* coordinator.interrupt("session").pipe(Effect.forkChild) + yield* Deferred.await(cleanupStarted) + const run = yield* coordinator.run("session").pipe(Effect.forkChild) + const close = yield* Scope.close(scope, Exit.void).pipe(Effect.forkChild) + + const runExit = yield* Fiber.await(run) + expect(Exit.isFailure(runExit) && Cause.hasInterruptsOnly(runExit.cause)).toBeTrue() + yield* Deferred.succeed(cleanupGate, undefined) + yield* Fiber.join(interrupt) + yield* Fiber.join(close) + }), + ) + it.effect("does not start work after its owning scope closes", () => Effect.gen(function* () { const scope = yield* Scope.make()