From 3ac64643e07bc49a8b89b940765672a526be6ed6 Mon Sep 17 00:00:00 2001 From: Aiden Cline Date: Tue, 14 Jul 2026 21:18:55 +0000 Subject: [PATCH] fix(core): keep interrupted sessions stopped --- packages/core/src/session.ts | 9 +++- packages/core/src/session/execution.ts | 4 +- packages/core/src/session/run-coordinator.ts | 47 ++++++++++++++----- packages/core/test/session-prompt.test.ts | 16 +++++-- .../core/test/session-run-coordinator.test.ts | 39 +++++++++++++-- packages/core/test/session-runner.test.ts | 18 +++++++ 6 files changed, 110 insertions(+), 23 deletions(-) diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index 2dabfb2d6f..1d35802820 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -379,7 +379,7 @@ const layer = Layer.effect( ) if (!SessionInput.equivalent(admitted, expected)) return yield* new PromptConflictError({ sessionID: input.sessionID, messageID }) - if (input.resume !== false) yield* execution.wake(admitted.sessionID) + if (input.resume !== false) yield* execution.wake(admitted.sessionID, admitted.admittedSeq) return admitted }), ), @@ -428,7 +428,12 @@ const layer = Layer.effect( yield* execution.resume(sessionID) }), interrupt: Effect.fn("V2Session.interrupt")((sessionID) => - Effect.uninterruptible(execution.interrupt(sessionID)), + Effect.uninterruptible( + Effect.gen(function* () { + if (!(yield* store.get(sessionID))) return yield* execution.interrupt(sessionID) + yield* execution.interrupt(sessionID, yield* EventV2.latestSequence(db, sessionID)) + }), + ), ), revert: { stage: Effect.fn("V2Session.revert.stage")(function* (input) { diff --git a/packages/core/src/session/execution.ts b/packages/core/src/session/execution.ts index 5938c37726..e307ab36a4 100644 --- a/packages/core/src/session/execution.ts +++ b/packages/core/src/session/execution.ts @@ -12,9 +12,9 @@ export interface Interface { /** Starts execution while idle or joins the active execution. */ readonly resume: (sessionID: SessionSchema.ID) => Effect.Effect /** Registers newly recorded work. Repeated wakeups may coalesce. */ - readonly wake: (sessionID: SessionSchema.ID) => Effect.Effect + readonly wake: (sessionID: SessionSchema.ID, seq?: number) => Effect.Effect /** Interrupt active work owned by this process. Idle interruption is a no-op. */ - readonly interrupt: (sessionID: SessionSchema.ID) => Effect.Effect + readonly interrupt: (sessionID: SessionSchema.ID, seq?: number) => Effect.Effect } /** Routes execution from a Session ID to the runner owned by that Session's Location. */ diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index 2f89aff9e3..65d0aba7e2 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -9,28 +9,39 @@ export interface Coordinator { /** Starts execution while idle or joins the active execution. */ readonly run: (key: Key) => Effect.Effect /** Registers one coalesced follow-up after newly recorded work. */ - readonly wake: (key: Key) => Effect.Effect + readonly wake: (key: Key, seq?: number) => Effect.Effect /** Stops active execution and waits for its cleanup. */ - readonly interrupt: (key: Key) => Effect.Effect + readonly interrupt: (key: Key, seq?: number) => Effect.Effect } +type Wake = { readonly seq?: number } + type Entry = { readonly done: Deferred.Deferred + currentWakeSeq?: number owner?: Fiber.Fiber - pendingWake: boolean + pendingWake?: Wake + interruptSeq?: number stopping: boolean } +const coalesceWake = (left: Wake | undefined, seq: number | undefined): Wake => { + if (left === undefined) return { seq } + if (left.seq === undefined || seq === undefined) return {} + return { seq: Math.max(left.seq, seq) } +} + export const make = (options: { readonly drain: (key: Key, force: boolean) => Effect.Effect }): Effect.Effect, never, Scope.Scope> => Effect.gen(function* () { const active = new Map>() + const interruptSeq = new Map() const fork = yield* FiberSet.makeRuntime() - const makeEntry = (): Entry => ({ + const makeEntry = (currentWakeSeq?: number): Entry => ({ done: Deferred.makeUnsafe(), - pendingWake: false, + currentWakeSeq, stopping: false, }) @@ -49,13 +60,15 @@ export const make = (options: { } const settle = (key: Key, entry: Entry, exit: Exit.Exit) => { - if (Exit.isSuccess(exit) && !entry.stopping && entry.pendingWake) { - entry.pendingWake = false + if (Exit.isSuccess(exit) && !entry.stopping && entry.pendingWake !== undefined) { + const pending = entry.pendingWake + entry.pendingWake = undefined + entry.currentWakeSeq = pending.seq start(key, entry, false, true) return } - const successor = entry.pendingWake ? makeEntry() : undefined + const successor = entry.pendingWake === undefined ? undefined : makeEntry(entry.pendingWake.seq) if (successor === undefined) active.delete(key) else { active.set(key, successor) @@ -78,25 +91,33 @@ export const make = (options: { return restore(Deferred.await(next.done)) }) - const wake = (key: Key) => + const wake = (key: Key, seq?: number) => Effect.sync(() => { + const latest = interruptSeq.get(key) + if (latest !== undefined && (seq === undefined || seq <= latest)) return const entry = active.get(key) if (entry !== undefined) { - entry.pendingWake = true + if (entry.stopping && entry.interruptSeq !== undefined && (seq === undefined || seq <= entry.interruptSeq)) + return + entry.pendingWake = coalesceWake(entry.pendingWake, seq) return } - const next = makeEntry() + const next = makeEntry(seq) active.set(key, next) start(key, next, false) }) - const interrupt = (key: Key): Effect.Effect => + const interrupt = (key: Key, seq?: number): Effect.Effect => Effect.suspend(() => { + if (seq !== undefined) interruptSeq.set(key, Math.max(interruptSeq.get(key) ?? seq, seq)) const entry = active.get(key) if (entry?.owner === undefined) return Effect.void + if (seq !== undefined && entry.currentWakeSeq !== undefined && entry.currentWakeSeq > seq) return Effect.void entry.stopping = true - entry.pendingWake = false + entry.interruptSeq = seq + if (seq === undefined || entry.pendingWake?.seq === undefined || entry.pendingWake.seq <= seq) + entry.pendingWake = undefined return Fiber.interrupt(entry.owner) }) diff --git a/packages/core/test/session-prompt.test.ts b/packages/core/test/session-prompt.test.ts index c6bc9430b3..1ebcea2fd9 100644 --- a/packages/core/test/session-prompt.test.ts +++ b/packages/core/test/session-prompt.test.ts @@ -22,7 +22,9 @@ import { testEffect } from "./lib/effect" const executionCalls: SessionV2.ID[] = [] const interruptCalls: SessionV2.ID[] = [] +const interruptSeqs: Array = [] const wakeCalls: SessionV2.ID[] = [] +const wakeSeqs: Array = [] const activeSessions = new Set() const execution = Layer.succeed( SessionExecution.Service, @@ -32,13 +34,15 @@ const execution = Layer.succeed( Effect.sync(() => { executionCalls.push(sessionID) }), - interrupt: (sessionID) => + interrupt: (sessionID, seq) => Effect.sync(() => { interruptCalls.push(sessionID) + interruptSeqs.push(seq) }), - wake: (sessionID) => + wake: (sessionID, seq) => Effect.sync(() => { wakeCalls.push(sessionID) + wakeSeqs.push(seq) }), }), ) @@ -123,9 +127,11 @@ describe("SessionV2.prompt", () => { yield* setup const session = yield* SessionV2.Service interruptCalls.length = 0 + interruptSeqs.length = 0 yield* session.interrupt(sessionID) expect(interruptCalls).toEqual([sessionID]) + expect(interruptSeqs).toEqual([-1]) expect(yield* session.messages({ sessionID })).toEqual([]) }), ) @@ -134,9 +140,11 @@ describe("SessionV2.prompt", () => { Effect.gen(function* () { const session = yield* SessionV2.Service interruptCalls.length = 0 + interruptSeqs.length = 0 yield* session.interrupt(SessionV2.ID.make("ses_missing")) expect(interruptCalls).toEqual([SessionV2.ID.make("ses_missing")]) + expect(interruptSeqs).toEqual([undefined]) }), ) @@ -542,11 +550,13 @@ describe("SessionV2.prompt", () => { const session = yield* SessionV2.Service executionCalls.length = 0 wakeCalls.length = 0 + wakeSeqs.length = 0 - yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Run by default" }) }) + const message = yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Run by default" }) }) expect(executionCalls).toEqual([]) expect(wakeCalls).toEqual([sessionID]) + expect(wakeSeqs).toEqual([message.admittedSeq]) }), ) diff --git a/packages/core/test/session-run-coordinator.test.ts b/packages/core/test/session-run-coordinator.test.ts index dfbeda664c..6f49b0f06e 100644 --- a/packages/core/test/session-run-coordinator.test.ts +++ b/packages/core/test/session-run-coordinator.test.ts @@ -244,6 +244,39 @@ describe("SessionRunCoordinator", () => { ), ) + it.effect("suppresses a stale wake registered during interruption cleanup", () => + Effect.scoped( + Effect.gen(function* () { + const firstStarted = yield* Deferred.make() + const cleanupStarted = yield* Deferred.make() + const cleanupGate = yield* Deferred.make() + let runs = 0 + const coordinator = yield* SessionRunCoordinator.make({ + drain: () => + Effect.sync(() => ++runs).pipe( + Effect.andThen(Deferred.succeed(firstStarted, undefined)), + Effect.andThen(Effect.never), + Effect.onInterrupt(() => + Deferred.succeed(cleanupStarted, undefined).pipe(Effect.andThen(Deferred.await(cleanupGate))), + ), + ), + }) + + 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", 1) + yield* Deferred.succeed(cleanupGate, undefined) + yield* Fiber.join(interrupt) + yield* Effect.yieldNow + + expect(runs).toBe(1) + expect(Array.from(yield* coordinator.active)).toEqual([]) + }), + ), + ) + it.effect("runs a wake registered during interruption cleanup", () => Effect.scoped( Effect.gen(function* () { @@ -268,11 +301,11 @@ describe("SessionRunCoordinator", () => { ), }) - yield* coordinator.wake("session") + yield* coordinator.wake("session", 1) yield* Deferred.await(firstStarted) - const interrupt = yield* coordinator.interrupt("session").pipe(Effect.forkChild) + const interrupt = yield* coordinator.interrupt("session", 2).pipe(Effect.forkChild) yield* Deferred.await(cleanupStarted) - yield* coordinator.wake("session") + yield* coordinator.wake("session", 3) yield* Deferred.succeed(cleanupGate, undefined) yield* Fiber.join(interrupt) yield* Deferred.await(secondStarted) diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index 0515d55cf5..f05f19271a 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -1898,6 +1898,24 @@ describe("SessionRunnerLLM", () => { }), ) + it.effect("does not wake an admitted prompt older than an interrupt", () => + Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + const id = SessionMessage.ID.create() + const prompt = Prompt.make({ text: "Remain stopped" }) + requests.length = 0 + streamGate = yield* Deferred.make() + + yield* session.prompt({ id, sessionID, prompt, resume: false }) + yield* session.interrupt(sessionID) + yield* session.prompt({ id, sessionID, prompt }) + + expect(Array.from(yield* session.active)).not.toContain(sessionID) + expect(requests).toHaveLength(0) + }), + ) + it.effect("preserves durable queued input for a later wake after interruption", () => Effect.gen(function* () { yield* setup