From 2240c77547c2fa8ce92b002562dd0aacbd2369bf Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 31 Aug 2026 12:53:54 -0400 Subject: [PATCH] test(core): guard atomic successor settlement Keep the open-coordinator decision and settlement in one synchronous callback. Separating them with an Effect boundary permits shutdown to skip both successor start and suspension. Sweep scheduler-yield boundaries in the retiring execution and document why settlement cannot be lifted into unconditional ensuring. --- packages/core/src/session/run-coordinator.ts | 2 + .../core/test/session-run-coordinator.test.ts | 72 ++++++++++++++++++- 2 files changed, 73 insertions(+), 1 deletion(-) diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index 71d04dd1eed..bddee4e39ae 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -113,6 +113,8 @@ export const make = (options: { return Effect.suspend(() => options.suspended?.(key) ?? Effect.void).pipe( Effect.ensuring(Effect.sync(() => settle(key, execution, exit))), ) + // Keep this decision and settlement synchronous. An Effect boundary could yield + // to shutdown after skipping suspension but before starting the successor. settle(key, execution, exit) return Effect.void }), diff --git a/packages/core/test/session-run-coordinator.test.ts b/packages/core/test/session-run-coordinator.test.ts index 5ece5996af8..7f9d3c0a183 100644 --- a/packages/core/test/session-run-coordinator.test.ts +++ b/packages/core/test/session-run-coordinator.test.ts @@ -1,5 +1,5 @@ import { describe, expect } from "bun:test" -import { Cause, Deferred, Effect, Exit, Fiber, Layer, Scope } from "effect" +import { Cause, Deferred, Effect, Exit, Fiber, Layer, Scheduler, Scope } from "effect" import { SessionInbox } from "@opencode-ai/core/session/inbox" import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator" import { testEffect } from "./lib/effect" @@ -250,6 +250,76 @@ describe("SessionRunCoordinator", () => { ) } + for (const cutoff of Array.from({ length: 40 }, (_, index) => index + 1)) { + it.effect(`shutdown at settlement operation ${cutoff} preserves its pending successor`, () => + Effect.gen(function* () { + const running = yield* Deferred.make() + const cleanup = yield* Deferred.make() + const release = yield* Deferred.make() + const closeNow = yield* Deferred.make() + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => + Deferred.succeed(release, undefined).pipe(Effect.andThen(Scope.close(scope, Exit.void))), + ) + const base = new Scheduler.MixedScheduler() + const dispatcher = base.makeDispatcher() + const lifecycle: string[] = [] + let owner: number | undefined + let seen = 0 + const scheduler: Scheduler.Scheduler = { + executionMode: base.executionMode, + makeDispatcher: () => base.makeDispatcher(), + shouldYield: (fiber) => { + if (fiber.id === owner && ++seen === cutoff) { + // Close on the next scheduler tick, never reentrantly inside the current operation. + dispatcher.scheduleTask(() => Deferred.doneUnsafe(closeNow, Exit.void), 0) + return true + } + return base.shouldYield(fiber) + }, + } + const coordinator = yield* SessionRunCoordinator.make({ + started: () => { + // Count scheduling even if shutdown interrupts the successor before its first drain. + lifecycle.push("scheduled") + return Effect.void + }, + drain: () => + Deferred.succeed(running, undefined).pipe( + Effect.andThen(Effect.never), + Effect.onInterrupt(() => + Deferred.succeed(cleanup, undefined).pipe(Effect.andThen(Deferred.await(release))), + ), + ), + settled: () => + Effect.withFiber((fiber) => { + owner ??= fiber.id + return Effect.void + }), + suspended: () => Effect.sync(() => void lifecycle.push("suspended")), + }).pipe(Scope.provide(scope), Effect.provideService(Scheduler.Scheduler, scheduler)) + + yield* coordinator.wake("session") + yield* Deferred.await(running) + yield* coordinator.interrupt("session") + yield* Deferred.await(cleanup) + yield* coordinator.wake("session") + const closing = yield* Deferred.await(closeNow).pipe( + Effect.andThen(Scope.close(scope, Exit.void)), + Effect.forkChild({ startImmediately: true }), + ) + yield* Deferred.succeed(release, undefined) + yield* Effect.yieldNow + // If settlement finished before the cutoff, close the already-scheduled successor. + yield* Deferred.succeed(closeNow, undefined) + yield* Fiber.join(closing) + + expect(lifecycle.length).toBe(2) + expect(yield* coordinator.active).toEqual(new Set()) + }), + ) + } + it.effect("coalesces wakes received during active execution", () => Effect.gen(function* () { const firstStarted = yield* Deferred.make()