diff --git a/packages/core/src/session/event.ts b/packages/core/src/session/event.ts index dab7562ce6a..aea8c009678 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 5bc3fc982cc..f933d43d6ca 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 206cfa57e07..0d8e9f2a66c 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 63155345ad6..bbe1ce75710 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 78dc5030b56..e22da3be54d 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 f69b331358c..d1ca442f136 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))) } }