From 911664db67e0ff420bfe52fc405cc466481006e4 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sat, 6 Jun 2026 08:09:36 -0400 Subject: [PATCH 1/5] fix(core): retry failed session wakes --- packages/core/src/session/run-coordinator.ts | 17 ++++++- .../core/test/session-run-coordinator.test.ts | 45 +++++++++++++++++-- 2 files changed, 56 insertions(+), 6 deletions(-) diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index 5068b913ad..7a3c971bc4 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 39fb5779d3..31fc40642e 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) }), ), From ecf6c532d5192a6b6d7867a9a1bb907647557b25 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sat, 6 Jun 2026 08:13:08 -0400 Subject: [PATCH 2/5] fix(core): hide recovered wake failures --- packages/core/src/session/run-coordinator.ts | 32 +++++++++---------- .../core/test/session-run-coordinator.test.ts | 5 +-- 2 files changed, 19 insertions(+), 18 deletions(-) diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index 7a3c971bc4..d1ca442f13 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -39,14 +39,14 @@ export interface Coordinator { /** 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> + readonly settled: Deferred.Deferred | undefined> current: Demand pending?: Demand explicitWaiter?: Deferred.Deferred interruptSeq?: number owner?: Fiber.Fiber stopping: boolean - advisoryRetriesRemaining: number + advisoryRetryAvailable: boolean } /** Combines follow-up demand: runs dominate, while wakes retain the newest durable admission sequence. */ @@ -82,17 +82,16 @@ export const make = (options: { }), ) - const makeEntry = ( - current: Demand, - explicitWaiter?: Deferred.Deferred, - advisoryRetriesRemaining = current._tag === "wake" ? 1 : 0, - ): Entry => ({ + const makeEntry = (current: Demand, options?: { + readonly explicitWaiter?: Deferred.Deferred + readonly advisoryRetryAvailable?: boolean + }): Entry => ({ done: Deferred.makeUnsafe(), - settled: Deferred.makeUnsafe>(), + settled: Deferred.makeUnsafe | undefined>(), current, - explicitWaiter, + explicitWaiter: options?.explicitWaiter, stopping: false, - advisoryRetriesRemaining, + advisoryRetryAvailable: options?.advisoryRetryAvailable ?? current._tag === "wake", }) const start = (key: Key, entry: Entry, demand: Demand, successor = false) => { @@ -138,7 +137,7 @@ export const make = (options: { const pending = entry.pending entry.pending = undefined entry.current = pending - entry.advisoryRetriesRemaining = pending._tag === "wake" ? 1 : 0 + entry.advisoryRetryAvailable = pending._tag === "wake" start(key, entry, pending, true) return } @@ -150,15 +149,16 @@ export const make = (options: { 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) + ? makeEntry(entry.pending, { explicitWaiter: entry.explicitWaiter }) + : exit._tag === "Failure" && demand._tag === "wake" && !entry.stopping && entry.advisoryRetryAvailable + ? makeEntry(demand, { explicitWaiter: entry.explicitWaiter, advisoryRetryAvailable: false }) : undefined + const retrying = successor !== undefined && entry.pending === 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)) + Deferred.doneUnsafe(entry.settled, Effect.succeed(retrying ? undefined : exit)) if ( exit._tag === "Failure" && !(entry.stopping && Cause.hasInterruptsOnly(exit.cause)) && @@ -197,7 +197,7 @@ export const make = (options: { Deferred.await(shutdown).pipe(Effect.as(Exit.void)), ) if (closed) break - if (exit._tag === "Failure" && firstFailure === undefined) firstFailure = exit.cause + if (exit?._tag === "Failure" && firstFailure === undefined) firstFailure = exit.cause } if (firstFailure !== undefined) return yield* Effect.failCause(firstFailure) }) diff --git a/packages/core/test/session-run-coordinator.test.ts b/packages/core/test/session-run-coordinator.test.ts index 31fc40642e..b30901e431 100644 --- a/packages/core/test/session-run-coordinator.test.ts +++ b/packages/core/test/session-run-coordinator.test.ts @@ -750,9 +750,10 @@ describe("SessionRunCoordinator", () => { yield* coordinator.wake("session", 1) yield* coordinator.wake("session", 2) yield* Deferred.succeed(firstGate, undefined) - yield* coordinator.awaitIdle("session").pipe(Effect.exit) + const idle = yield* coordinator.awaitIdle("session").pipe(Effect.exit) - expect(runs).toBe(3) + expect(runs).toBe(3) + expect(Exit.isSuccess(idle)).toBeTrue() }), ), ) From fd07e583848d35409bfa802d26079e5666d39832 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sat, 6 Jun 2026 09:48:31 -0400 Subject: [PATCH 3/5] feat(core): publish terminal session run failures --- packages/core/src/session/event.ts | 18 +++++++++ packages/core/src/session/execution/local.ts | 42 +++++++++++++++++--- packages/core/src/session/input.ts | 14 +++++++ packages/core/src/session/message-updater.ts | 1 + packages/core/src/session/projector.ts | 1 + packages/core/src/session/run-coordinator.ts | 8 +++- 6 files changed, 76 insertions(+), 8 deletions(-) diff --git a/packages/core/src/session/event.ts b/packages/core/src/session/event.ts index aea8c00967..dab7562ce6 100644 --- a/packages/core/src/session/event.ts +++ b/packages/core/src/session/event.ts @@ -119,6 +119,23 @@ export namespace PromptLifecycle { export type Promoted = typeof Promoted.Type } +export namespace Run { + export const Failed = EventV2.define({ + type: "session.next.run.failed", + ...options, + schema: { + ...Base, + reason: Schema.Literals(["execution-failed", "step-limit-exceeded", "unknown"]), + input: Schema.Struct({ + messageID: SessionMessageID.ID, + admittedSeq: NonNegativeInt, + promotedSeq: NonNegativeInt.pipe(Schema.optional), + }).pipe(Schema.optional), + }, + }) + export type Failed = typeof Failed.Type +} + export const InterruptRequested = EventV2.define({ type: "session.next.interrupt.requested", ...options, @@ -476,6 +493,7 @@ const DurableDefinitions = [ Prompted, PromptLifecycle.Admitted, PromptLifecycle.Promoted, + Run.Failed, InterruptRequested, ContextUpdated, Synthetic, diff --git a/packages/core/src/session/execution/local.ts b/packages/core/src/session/execution/local.ts index f933d43d6c..5bc3fc982c 100644 --- a/packages/core/src/session/execution/local.ts +++ b/packages/core/src/session/execution/local.ts @@ -1,10 +1,15 @@ -import { Effect, Layer } from "effect" +import { Cause, DateTime, Effect, Layer, Option } from "effect" +import { LLMError } from "@opencode-ai/llm" +import { Database } from "../../database/database" +import { EventV2 } from "../../event" import { LocationServiceMap } from "../../location-layer" import { SessionRunCoordinator } from "../run-coordinator" import { SessionRunner } from "../runner" import { SessionSchema } from "../schema" import { SessionStore } from "../store" import { SessionExecution } from "../execution" +import { SessionEvent } from "../event" +import { SessionInput } from "../input" /** Current-process routing for implicit-local Locations. Future remote placement belongs here. */ export const layer = Layer.effect( @@ -12,6 +17,8 @@ export const layer = Layer.effect( Effect.gen(function* () { const store = yield* SessionStore.Service const locations = yield* LocationServiceMap + const database = yield* Database.Service + const events = yield* EventV2.Service const coordinator = yield* SessionRunCoordinator.make({ drain: Effect.fnUntraced(function* (sessionID: SessionSchema.ID, mode) { const session = yield* store.get(sessionID) @@ -20,11 +27,34 @@ export const layer = Layer.effect( Effect.provide(locations.get(session.location)), ) }), - onFailure: (sessionID, cause) => - Effect.logError("Failed to drain Session").pipe( - Effect.annotateLogs("sessionID", sessionID), - Effect.annotateLogs("cause", cause), - ), + onFailure: (sessionID, cause, context) => + Effect.gen(function* () { + yield* Effect.logError("Failed to drain Session").pipe( + Effect.annotateLogs("sessionID", sessionID), + Effect.annotateLogs("cause", cause), + ) + if (Cause.hasInterruptsOnly(cause)) return + const error = Option.getOrUndefined(Cause.findErrorOption(cause)) + // Provider failures already publish Step.Failed before escaping the runner. + if (error instanceof LLMError) return + const input = context.seq === undefined + ? undefined + : yield* SessionInput.findByAdmittedSeq(database.db, sessionID, context.seq) + yield* events.publish(SessionEvent.Run.Failed, { + sessionID, + timestamp: yield* DateTime.now, + reason: error instanceof SessionRunner.StepLimitExceededError ? "step-limit-exceeded" : "execution-failed", + ...(input === undefined + ? {} + : { + input: { + messageID: input.id, + admittedSeq: input.admittedSeq, + ...(input.promotedSeq === undefined ? {} : { promotedSeq: input.promotedSeq }), + }, + }), + }) + }), }) return SessionExecution.Service.of({ diff --git a/packages/core/src/session/input.ts b/packages/core/src/session/input.ts index 0d8e9f2a66..206cfa57e0 100644 --- a/packages/core/src/session/input.ts +++ b/packages/core/src/session/input.ts @@ -47,6 +47,20 @@ export const find = Effect.fn("SessionInput.find")(function* (db: DatabaseServic return row === undefined ? undefined : fromRow(row) }) +export const findByAdmittedSeq = Effect.fn("SessionInput.findByAdmittedSeq")(function* ( + db: DatabaseService, + sessionID: SessionSchema.ID, + admittedSeq: number, +) { + const row = yield* db + .select() + .from(SessionInputTable) + .where(and(eq(SessionInputTable.session_id, sessionID), eq(SessionInputTable.admitted_seq, admittedSeq))) + .get() + .pipe(Effect.orDie) + return row === undefined ? undefined : fromRow(row) +}) + export class LifecycleConflict extends Schema.TaggedErrorClass()("SessionInput.LifecycleConflict", { id: SessionMessage.ID, }) {} diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index bbe1ce7571..63155345ad 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -139,6 +139,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) { }, "session.next.prompt.admitted": () => Effect.void, "session.next.prompt.promoted": () => Effect.void, + "session.next.run.failed": () => Effect.void, "session.next.interrupt.requested": () => Effect.void, "session.next.context.updated": (event) => adapter.appendMessage( diff --git a/packages/core/src/session/projector.ts b/packages/core/src/session/projector.ts index e22da3be54..78dc5030b5 100644 --- a/packages/core/src/session/projector.ts +++ b/packages/core/src/session/projector.ts @@ -410,6 +410,7 @@ export const layer = Layer.effectDiscard( ) }), ) + yield* events.project(SessionEvent.Run.Failed, () => Effect.void) yield* events.project(SessionEvent.InterruptRequested, () => Effect.void) yield* events.project(SessionEvent.ContextUpdated, (event) => { if (!event.replay || event.seq === undefined) return run(db, event) diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index d1ca442f13..f69b331358 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -64,7 +64,11 @@ const maxSeq = (left: number | undefined, right: number | undefined) => { /** 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 + readonly onFailure?: ( + key: Key, + cause: Cause.Cause, + context: { readonly mode: "wake"; readonly seq?: number }, + ) => Effect.Effect }): Effect.Effect, never, Scope.Scope> => Effect.gen(function* () { const active = new Map>() @@ -166,7 +170,7 @@ export const make = (options: { successor === undefined && options.onFailure !== undefined ) { - report(Effect.suspend(() => options.onFailure!(key, exit.cause))) + report(Effect.suspend(() => options.onFailure!(key, exit.cause, { mode: "wake", seq: demand.seq }))) } } From 93ca32f41a8b11a49b1be3dcf421a88d0709062c Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sat, 6 Jun 2026 22:01:43 -0400 Subject: [PATCH 4/5] Revert "feat(core): publish terminal session run failures" This reverts commit fd07e583848d35409bfa802d26079e5666d39832. --- packages/core/src/session/event.ts | 18 --------- packages/core/src/session/execution/local.ts | 42 +++----------------- packages/core/src/session/input.ts | 14 ------- packages/core/src/session/message-updater.ts | 1 - packages/core/src/session/projector.ts | 1 - packages/core/src/session/run-coordinator.ts | 8 +--- 6 files changed, 8 insertions(+), 76 deletions(-) diff --git a/packages/core/src/session/event.ts b/packages/core/src/session/event.ts index dab7562ce6..aea8c00967 100644 --- a/packages/core/src/session/event.ts +++ b/packages/core/src/session/event.ts @@ -119,23 +119,6 @@ export namespace PromptLifecycle { export type Promoted = typeof Promoted.Type } -export namespace Run { - export const Failed = EventV2.define({ - type: "session.next.run.failed", - ...options, - schema: { - ...Base, - reason: Schema.Literals(["execution-failed", "step-limit-exceeded", "unknown"]), - input: Schema.Struct({ - messageID: SessionMessageID.ID, - admittedSeq: NonNegativeInt, - promotedSeq: NonNegativeInt.pipe(Schema.optional), - }).pipe(Schema.optional), - }, - }) - export type Failed = typeof Failed.Type -} - export const InterruptRequested = EventV2.define({ type: "session.next.interrupt.requested", ...options, @@ -493,7 +476,6 @@ const DurableDefinitions = [ Prompted, PromptLifecycle.Admitted, PromptLifecycle.Promoted, - Run.Failed, InterruptRequested, ContextUpdated, Synthetic, diff --git a/packages/core/src/session/execution/local.ts b/packages/core/src/session/execution/local.ts index 5bc3fc982c..f933d43d6c 100644 --- a/packages/core/src/session/execution/local.ts +++ b/packages/core/src/session/execution/local.ts @@ -1,15 +1,10 @@ -import { Cause, DateTime, Effect, Layer, Option } from "effect" -import { LLMError } from "@opencode-ai/llm" -import { Database } from "../../database/database" -import { EventV2 } from "../../event" +import { Effect, Layer } from "effect" import { LocationServiceMap } from "../../location-layer" import { SessionRunCoordinator } from "../run-coordinator" import { SessionRunner } from "../runner" import { SessionSchema } from "../schema" import { SessionStore } from "../store" import { SessionExecution } from "../execution" -import { SessionEvent } from "../event" -import { SessionInput } from "../input" /** Current-process routing for implicit-local Locations. Future remote placement belongs here. */ export const layer = Layer.effect( @@ -17,8 +12,6 @@ export const layer = Layer.effect( Effect.gen(function* () { const store = yield* SessionStore.Service const locations = yield* LocationServiceMap - const database = yield* Database.Service - const events = yield* EventV2.Service const coordinator = yield* SessionRunCoordinator.make({ drain: Effect.fnUntraced(function* (sessionID: SessionSchema.ID, mode) { const session = yield* store.get(sessionID) @@ -27,34 +20,11 @@ export const layer = Layer.effect( Effect.provide(locations.get(session.location)), ) }), - onFailure: (sessionID, cause, context) => - Effect.gen(function* () { - yield* Effect.logError("Failed to drain Session").pipe( - Effect.annotateLogs("sessionID", sessionID), - Effect.annotateLogs("cause", cause), - ) - if (Cause.hasInterruptsOnly(cause)) return - const error = Option.getOrUndefined(Cause.findErrorOption(cause)) - // Provider failures already publish Step.Failed before escaping the runner. - if (error instanceof LLMError) return - const input = context.seq === undefined - ? undefined - : yield* SessionInput.findByAdmittedSeq(database.db, sessionID, context.seq) - yield* events.publish(SessionEvent.Run.Failed, { - sessionID, - timestamp: yield* DateTime.now, - reason: error instanceof SessionRunner.StepLimitExceededError ? "step-limit-exceeded" : "execution-failed", - ...(input === undefined - ? {} - : { - input: { - messageID: input.id, - admittedSeq: input.admittedSeq, - ...(input.promotedSeq === undefined ? {} : { promotedSeq: input.promotedSeq }), - }, - }), - }) - }), + onFailure: (sessionID, cause) => + Effect.logError("Failed to drain Session").pipe( + Effect.annotateLogs("sessionID", sessionID), + Effect.annotateLogs("cause", cause), + ), }) return SessionExecution.Service.of({ diff --git a/packages/core/src/session/input.ts b/packages/core/src/session/input.ts index 206cfa57e0..0d8e9f2a66 100644 --- a/packages/core/src/session/input.ts +++ b/packages/core/src/session/input.ts @@ -47,20 +47,6 @@ export const find = Effect.fn("SessionInput.find")(function* (db: DatabaseServic return row === undefined ? undefined : fromRow(row) }) -export const findByAdmittedSeq = Effect.fn("SessionInput.findByAdmittedSeq")(function* ( - db: DatabaseService, - sessionID: SessionSchema.ID, - admittedSeq: number, -) { - const row = yield* db - .select() - .from(SessionInputTable) - .where(and(eq(SessionInputTable.session_id, sessionID), eq(SessionInputTable.admitted_seq, admittedSeq))) - .get() - .pipe(Effect.orDie) - return row === undefined ? undefined : fromRow(row) -}) - export class LifecycleConflict extends Schema.TaggedErrorClass()("SessionInput.LifecycleConflict", { id: SessionMessage.ID, }) {} diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index 63155345ad..bbe1ce7571 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -139,7 +139,6 @@ export function update(adapter: Adapter, event: SessionEvent.Event) { }, "session.next.prompt.admitted": () => Effect.void, "session.next.prompt.promoted": () => Effect.void, - "session.next.run.failed": () => Effect.void, "session.next.interrupt.requested": () => Effect.void, "session.next.context.updated": (event) => adapter.appendMessage( diff --git a/packages/core/src/session/projector.ts b/packages/core/src/session/projector.ts index 78dc5030b5..e22da3be54 100644 --- a/packages/core/src/session/projector.ts +++ b/packages/core/src/session/projector.ts @@ -410,7 +410,6 @@ export const layer = Layer.effectDiscard( ) }), ) - yield* events.project(SessionEvent.Run.Failed, () => Effect.void) yield* events.project(SessionEvent.InterruptRequested, () => Effect.void) yield* events.project(SessionEvent.ContextUpdated, (event) => { if (!event.replay || event.seq === undefined) return run(db, event) diff --git a/packages/core/src/session/run-coordinator.ts b/packages/core/src/session/run-coordinator.ts index f69b331358..d1ca442f13 100644 --- a/packages/core/src/session/run-coordinator.ts +++ b/packages/core/src/session/run-coordinator.ts @@ -64,11 +64,7 @@ const maxSeq = (left: number | undefined, right: number | undefined) => { /** 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, - context: { readonly mode: "wake"; readonly seq?: number }, - ) => Effect.Effect + readonly onFailure?: (key: Key, cause: Cause.Cause) => Effect.Effect }): Effect.Effect, never, Scope.Scope> => Effect.gen(function* () { const active = new Map>() @@ -170,7 +166,7 @@ export const make = (options: { successor === undefined && options.onFailure !== undefined ) { - report(Effect.suspend(() => options.onFailure!(key, exit.cause, { mode: "wake", seq: demand.seq }))) + report(Effect.suspend(() => options.onFailure!(key, exit.cause))) } } From ab3850d1b4a640e3706ab9f3328744c47a004e96 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sat, 6 Jun 2026 22:02:34 -0400 Subject: [PATCH 5/5] test(core): format wake retry coverage --- packages/core/test/session-run-coordinator.test.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/core/test/session-run-coordinator.test.ts b/packages/core/test/session-run-coordinator.test.ts index b30901e431..c4eefc1058 100644 --- a/packages/core/test/session-run-coordinator.test.ts +++ b/packages/core/test/session-run-coordinator.test.ts @@ -750,10 +750,10 @@ describe("SessionRunCoordinator", () => { yield* coordinator.wake("session", 1) yield* coordinator.wake("session", 2) yield* Deferred.succeed(firstGate, undefined) - const idle = yield* coordinator.awaitIdle("session").pipe(Effect.exit) + const idle = yield* coordinator.awaitIdle("session").pipe(Effect.exit) - expect(runs).toBe(3) - expect(Exit.isSuccess(idle)).toBeTrue() + expect(runs).toBe(3) + expect(Exit.isSuccess(idle)).toBeTrue() }), ), )