diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index bddee4e39ae..d9b936c1f28 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -98,9 +98,11 @@ export const make = (options: { // The leading yield lets `owner` be assigned before the drain can settle, and keeps // failing self-waking executions from growing the stack across successor starts. // Drains start one tick after wake; callers observe progress through events or run. + // Mask the yield with startup so shutdown cannot skip the write-ahead started hook. execution.owner = fork( Effect.yieldNow.pipe( - Effect.andThen(Effect.uninterruptible(options.started?.(key) ?? Effect.void)), + Effect.andThen(options.started?.(key) ?? Effect.void), + Effect.uninterruptible, Effect.andThen(loop(key, execution, force)), Effect.onExit((exit) => Effect.sync(() => { diff --git a/packages/core/test/session-execution.test.ts b/packages/core/test/session-execution.test.ts index 23fa86d195c..ad5e7605a7a 100644 --- a/packages/core/test/session-execution.test.ts +++ b/packages/core/test/session-execution.test.ts @@ -163,111 +163,129 @@ describe("SessionExecution lifecycle", () => { }), ) - it.live("recovers a move admitted during user-interruption cleanup when shutdown prevents its successor", () => - Effect.gen(function* () { - const database = yield* Database.Service - const bus = yield* Bus.Service - const store = yield* SessionStore.Service - const admission = yield* SessionInbox.Service - const jobs = yield* Job.Service - const destination = AbsolutePath.make((yield* tmpdirScoped()).path) - const sessionID = Session.ID.make("ses_shutdown_successor") - yield* seedSessions(database, [sessionID], { resume_attempts: 2 }) - const lifecycle: string[] = [] - yield* bus.project(SessionEvent.Execution.Started, () => Effect.sync(() => void lifecycle.push("started"))) - yield* bus.project(SessionEvent.Execution.Interrupted, (event) => - Effect.sync(() => void lifecycle.push(event.data.reason)), - ) + for (const boundary of ["cleanup", "successor"]) { + it.live(`recovers a cleanup-era move when shutdown interrupts ${boundary}`, () => + Effect.gen(function* () { + const database = yield* Database.Service + const bus = yield* Bus.Service + const store = yield* SessionStore.Service + const admission = yield* SessionInbox.Service + const jobs = yield* Job.Service + const destination = AbsolutePath.make((yield* tmpdirScoped()).path) + const sessionID = Session.ID.make("ses_shutdown_successor") + yield* seedSessions(database, [sessionID], { resume_attempts: 2 }) + const lifecycle: string[] = [] + yield* bus.project(SessionEvent.Execution.Started, () => Effect.sync(() => void lifecycle.push("started"))) + yield* bus.project(SessionEvent.Execution.Interrupted, (event) => + Effect.sync(() => void lifecycle.push(event.data.reason)), + ) - const draining = yield* Deferred.make() - const cleanup = yield* Deferred.make() - const release = yield* Deferred.make() - const scope = yield* Scope.make() - yield* Effect.addFinalizer(() => - Deferred.succeed(release, undefined).pipe(Effect.andThen(Scope.close(scope, Exit.void))), - ) - const drains: string[] = [] - const context = yield* buildExecution(scope, () => - Effect.sync(() => void drains.push("original")).pipe( - Effect.andThen(Deferred.succeed(draining, undefined)), - Effect.andThen(Effect.never), - Effect.onInterrupt(() => Deferred.succeed(cleanup, undefined).pipe(Effect.andThen(Deferred.await(release)))), - ), - ) - const execution = Context.get(context, SessionExecution.Service) - const sessions = Context.get( - yield* Layer.buildWithScope( - AppNodeBuilder.build(Session.node, [ - [Database.node, Layer.succeed(Database.Service, database)], - [Bus.node, Layer.succeed(Bus.Service, bus)], - [SessionStore.node, Layer.succeed(SessionStore.Service, store)], - [SessionInbox.node, Layer.succeed(SessionInbox.Service, admission)], - [Job.node, Layer.succeed(Job.Service, jobs)], - // The shared Bus already has the production Session projectors. - [SessionProjector.node, Layer.empty], - [Project.node, globalProjectNode], - [SessionExecution.node, Layer.succeed(SessionExecution.Service, execution)], - [ - LocationServiceMap.node, - Layer.effect( - LocationServiceMap.Service, - LayerMap.make( - (ref: Location.Ref) => - // Move validation only needs Location from the destination graph. - // oxlint-disable-next-line typescript-eslint/no-unsafe-type-assertion - LayerNode.compile(Location.boundNode(ref), [ - [Project.node, globalProjectNode], - ]) as unknown as Layer.Layer, + const draining = yield* Deferred.make() + const cleanup = yield* Deferred.make() + const release = yield* Deferred.make() + const scope = yield* Scope.make() + yield* Effect.addFinalizer(() => + Deferred.succeed(release, undefined).pipe(Effect.andThen(Scope.close(scope, Exit.void))), + ) + const drains: string[] = [] + const context = yield* buildExecution(scope, () => + Effect.sync(() => void drains.push("original")).pipe( + Effect.andThen(Deferred.succeed(draining, undefined)), + Effect.andThen(Effect.never), + Effect.onInterrupt(() => + Deferred.succeed(cleanup, undefined).pipe(Effect.andThen(Deferred.await(release))), + ), + ), + ) + const execution = Context.get(context, SessionExecution.Service) + const sessions = Context.get( + yield* Layer.buildWithScope( + AppNodeBuilder.build(Session.node, [ + [Database.node, Layer.succeed(Database.Service, database)], + [Bus.node, Layer.succeed(Bus.Service, bus)], + [SessionStore.node, Layer.succeed(SessionStore.Service, store)], + [SessionInbox.node, Layer.succeed(SessionInbox.Service, admission)], + [Job.node, Layer.succeed(Job.Service, jobs)], + // The shared Bus already has the production Session projectors. + [SessionProjector.node, Layer.empty], + [Project.node, globalProjectNode], + [SessionExecution.node, Layer.succeed(SessionExecution.Service, execution)], + [ + LocationServiceMap.node, + Layer.effect( + LocationServiceMap.Service, + LayerMap.make( + (ref: Location.Ref) => + // Move validation only needs Location from the destination graph. + // oxlint-disable-next-line typescript-eslint/no-unsafe-type-assertion + LayerNode.compile(Location.boundNode(ref), [ + [Project.node, globalProjectNode], + ]) as unknown as Layer.Layer, + ), ), - ), - ], - ]).pipe(Layer.fresh), - scope, - ), - Session.Service, - ) + ], + ]).pipe(Layer.fresh), + scope, + ), + Session.Service, + ) - yield* execution.wake(sessionID) - yield* Deferred.await(draining) - yield* sessions.interrupt(sessionID) - yield* Deferred.await(cleanup) - yield* sessions.move({ sessionID, directory: destination }) - expect(yield* sessions.inbox(sessionID)).toMatchObject([{ type: "move" }]) + yield* execution.wake(sessionID) + yield* Deferred.await(draining) + // A joined caller observes the retiring execution's exit after its successor is forked. + const joined = yield* execution.resume(sessionID).pipe( + Effect.onExit(() => (boundary === "successor" ? Scope.close(scope, Exit.void) : Effect.void)), + Effect.exit, + Effect.forkChild({ startImmediately: true }), + ) + yield* sessions.interrupt(sessionID) + yield* Deferred.await(cleanup) + yield* sessions.move({ sessionID, directory: destination }) + expect(yield* sessions.inbox(sessionID)).toMatchObject([{ type: "move" }]) - const closing = yield* Scope.close(scope, Exit.void).pipe(Effect.forkChild({ startImmediately: true })) - yield* Effect.yieldNow - yield* Deferred.succeed(release, undefined) - yield* Fiber.join(closing) + const closing = + boundary === "cleanup" + ? yield* Scope.close(scope, Exit.void).pipe(Effect.forkChild({ startImmediately: true })) + : undefined + if (closing) yield* Effect.yieldNow + yield* Deferred.succeed(release, undefined) + if (closing) yield* Fiber.join(closing) + yield* Fiber.join(joined) - expect({ active: yield* execution.isActive(sessionID), claimed: (yield* claims(database))[sessionID] }).toEqual({ - active: false, - claimed: true, - }) - yield* execution.awaitIdle(sessionID) - expect(drains).toEqual(["original"]) - expect(lifecycle).toEqual(["started", "user"]) - // The interrupted intent releases its old recovery budget; the pending move is new work. - expect(yield* attempts(database, sessionID)).toBe(0) + expect({ active: yield* execution.isActive(sessionID), claimed: (yield* claims(database))[sessionID] }).toEqual( + { + active: false, + claimed: true, + }, + ) + yield* execution.awaitIdle(sessionID) + expect(drains).toEqual(["original"]) + expect(lifecycle).toEqual( + boundary === "cleanup" ? ["started", "user"] : ["started", "user", "started", "shutdown"], + ) + // The interrupted intent releases its old recovery budget; the pending move is new work. + expect(yield* attempts(database, sessionID)).toBe(0) - const restartedScope = yield* Scope.make() - yield* Effect.addFinalizer(() => Scope.close(restartedScope, Exit.void)) - const resumed = yield* Deferred.make() - const pending: SessionInbox.Item["type"][] = [] - const restarted = yield* buildExecution(restartedScope, () => - Effect.gen(function* () { - drains.push("restarted") - pending.push(...(yield* SessionInbox.list(database.db, sessionID)).map((item) => item.type)) - yield* Deferred.succeed(resumed, undefined) - }), - ) - yield* Context.get(restarted, SessionRestart.Service).resumeSuspendedSessions - yield* Deferred.await(resumed) - yield* Context.get(restarted, SessionExecution.Service).awaitIdle(sessionID) - expect(drains).toEqual(["original", "restarted"]) - expect(pending).toEqual(["move"]) - expect((yield* claims(database))[sessionID]).toBe(false) - }), - ) + const restartedScope = yield* Scope.make() + yield* Effect.addFinalizer(() => Scope.close(restartedScope, Exit.void)) + const resumed = yield* Deferred.make() + const pending: SessionInbox.Item["type"][] = [] + const restarted = yield* buildExecution(restartedScope, () => + Effect.gen(function* () { + drains.push("restarted") + pending.push(...(yield* SessionInbox.list(database.db, sessionID)).map((item) => item.type)) + yield* Deferred.succeed(resumed, undefined) + }), + ) + yield* Context.get(restarted, SessionRestart.Service).resumeSuspendedSessions + yield* Deferred.await(resumed) + yield* Context.get(restarted, SessionExecution.Service).awaitIdle(sessionID) + expect(drains).toEqual(["original", "restarted"]) + expect(pending).toEqual(["move"]) + expect((yield* claims(database))[sessionID]).toBe(false) + }), + ) + } it.effect("does not resume a user-cancelled background child whose notification was not admitted", () => Effect.gen(function* () { diff --git a/packages/core/test/session-run-coordinator.test.ts b/packages/core/test/session-run-coordinator.test.ts index 7f9d3c0a183..dc867079f2d 100644 --- a/packages/core/test/session-run-coordinator.test.ts +++ b/packages/core/test/session-run-coordinator.test.ts @@ -279,11 +279,7 @@ describe("SessionRunCoordinator", () => { }, } const coordinator = yield* SessionRunCoordinator.make({ - started: () => { - // Count scheduling even if shutdown interrupts the successor before its first drain. - lifecycle.push("scheduled") - return Effect.void - }, + started: () => Effect.sync(() => void lifecycle.push("started")), drain: () => Deferred.succeed(running, undefined).pipe( Effect.andThen(Effect.never),