From cacb08c1be68c7436f3f812a9fb24bf3d931ad51 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sat, 6 Jun 2026 22:30:40 -0400 Subject: [PATCH] refactor(core): isolate provider turn settlement --- packages/core/src/session/runner/run-turn.ts | 149 ++++++++---------- .../session/runner/settle-provider-turn.ts | 71 +++++++++ 2 files changed, 133 insertions(+), 87 deletions(-) create mode 100644 packages/core/src/session/runner/settle-provider-turn.ts diff --git a/packages/core/src/session/runner/run-turn.ts b/packages/core/src/session/runner/run-turn.ts index cbcda8ffd4..404fa4dfed 100644 --- a/packages/core/src/session/runner/run-turn.ts +++ b/packages/core/src/session/runner/run-turn.ts @@ -18,7 +18,7 @@ import { isContextOverflowFailure, type ProviderErrorEvent, } from "@opencode-ai/llm" -import { Cause, DateTime, Effect, FiberSet, Option, Schema, Semaphore, Stream } from "effect" +import { DateTime, Effect, Schema, Semaphore, Stream } from "effect" import { AgentV2 } from "../../agent" import { Config } from "../../config" import { Database } from "../../database/database" @@ -26,11 +26,9 @@ import { EventV2 } from "../../event" import { Location } from "../../location" import { ModelV2 } from "../../model" import { ProviderV2 } from "../../provider" -import { QuestionV2 } from "../../question" import { SkillGuidance } from "../../skill/guidance" import { SystemContext } from "../../system-context/index" import { SystemContextRegistry } from "../../system-context/registry" -import { ToolOutputStore } from "../../tool-output-store" import { ToolRegistry } from "../../tool/registry" import { SessionCompaction } from "../compaction" import { SessionContextEpoch } from "../context-epoch" @@ -42,6 +40,7 @@ import { SessionStore } from "../store" import type { RunError } from "./index" import { SessionRunnerModel } from "./model" import { createLLMEventPublisher } from "./publish-llm-event" +import { SettleProviderTurn } from "./settle-provider-turn" import { toLLMMessages } from "./to-llm-message" export interface Input { @@ -73,10 +72,6 @@ export const make = Effect.gen(function* () { if (!session) return yield* Effect.die(`Session not found: ${sessionID}`) return session }) - const awaitToolFibers = (fibers: FiberSet.FiberSet) => - Effect.raceFirst(FiberSet.join(fibers), FiberSet.awaitEmpty(fibers)) - const isQuestionRejected = (cause: Cause.Cause) => - cause.reasons.some((reason) => Cause.isDieReason(reason) && reason.defect instanceof QuestionV2.RejectedError) const stale = Symbol("stale turn preparation") const retryAgentMismatch = (effect: Effect.Effect) => effect.pipe( @@ -181,21 +176,19 @@ export const make = Effect.gen(function* () { withPublication(publisher.publish(event, outputPaths)) const failUnsettled = (message: string, providerExecuted = false) => withPublication(publisher.failUnsettledTools(message, providerExecuted)) - const toolFibers = yield* FiberSet.make() let needsContinuation = false let overflowFailure: ProviderErrorEvent | undefined - const startTool = Effect.fnUntraced(function* (event: Extract) { + const toolEffect = Effect.fnUntraced(function* (event: Extract) { needsContinuation = true const assistantMessageID = yield* publisher.assistantMessageID(event.id) - yield* Effect.uninterruptibleMask((restore) => - restore( - prepared.toolMaterialization.settle({ - sessionID: prepared.session.id, - agent: prepared.agent.id, - assistantMessageID, - call: event, - }), - ).pipe( + return yield* prepared.toolMaterialization + .settle({ + sessionID: prepared.session.id, + agent: prepared.agent.id, + assistantMessageID, + call: event, + }) + .pipe( Effect.flatMap((settlement) => publish( LLMEvent.toolResult({ @@ -207,85 +200,67 @@ export const make = Effect.gen(function* () { settlement.outputPaths ?? [], ), ), - ), - ).pipe(FiberSet.run(toolFibers)) + ) }) - const providerStream = llm.stream(prepared.request).pipe( - Stream.runForEach((event) => + const result = yield* SettleProviderTurn.run({ + stream: (runTool) => + llm.stream(prepared.request).pipe( + Stream.runForEach((event) => + Effect.gen(function* () { + if (overflowFailure || publisher.hasProviderError()) return + if ( + LLMEvent.is.providerError(event) && + isContextOverflowFailure(event) && + !publisher.hasAssistantStarted() + ) { + overflowFailure = event + return + } + yield* publish(event) + if (event.type !== "tool-call" || event.providerExecuted) return + yield* runTool(toolEffect(event)) + }), + ), + Effect.ensuring(withPublication(publisher.flush())), + ), + recoverOverflow: (failure, restore) => Effect.gen(function* () { - if (overflowFailure || publisher.hasProviderError()) return - if (LLMEvent.is.providerError(event) && isContextOverflowFailure(event) && !publisher.hasAssistantStarted()) { - overflowFailure = event - return - } - yield* publish(event) - if (event.type !== "tool-call" || event.providerExecuted) return - yield* startTool(event) - }), - ), - Effect.ensuring(withPublication(publisher.flush())), - ) - - // Keep cleanup protected after the response ends so no started tool is - // forgotten, while the response stream and tool work remain interruptible. - return yield* Effect.uninterruptibleMask((restore) => - Effect.gen(function* () { - const stream = yield* restore(providerStream).pipe(Effect.exit) - const failure = - stream._tag === "Failure" ? Option.getOrUndefined(Cause.findErrorOption(stream.cause)) : undefined - if ( - canRecoverOverflow && - !publisher.hasAssistantStarted() && - isContextOverflowFailure(overflowFailure ?? failure) && - (yield* restore( + if (!canRecoverOverflow || publisher.hasAssistantStarted()) return false + if (!isContextOverflowFailure(overflowFailure ?? failure)) return false + return yield* restore( compaction.compactAfterOverflow({ sessionID: prepared.session.id, entries: prepared.entries, model: prepared.model, request: prepared.request, }), - )) - ) - return AttemptResult.cases.CompactedOverflow.make({}) - if (overflowFailure) yield* publish(overflowFailure) - const llmFailure = failure instanceof LLMError ? failure : undefined - if (llmFailure && !publisher.hasProviderError()) { - yield* failUnsettled("Provider did not return a tool result", true) - yield* withPublication( - events.publish(SessionEvent.Step.Failed, { - sessionID: prepared.session.id, - timestamp: yield* DateTime.now, - assistantMessageID: yield* publisher.startAssistant(), - error: { type: "unknown", message: llmFailure.reason.message }, - }), ) - } - const streamInterrupted = stream._tag === "Failure" && Cause.hasInterrupts(stream.cause) - if (streamInterrupted) yield* FiberSet.clear(toolFibers) - const settled = yield* restore(awaitToolFibers(toolFibers)).pipe(Effect.exit) - if (settled._tag === "Failure" && isQuestionRejected(settled.cause)) { - yield* FiberSet.clear(toolFibers) - yield* failUnsettled("Tool execution interrupted") - return yield* Effect.interrupt - } - const toolInterrupted = settled._tag === "Failure" && Cause.hasInterrupts(settled.cause) - if (toolInterrupted) yield* FiberSet.clear(toolFibers) - if (streamInterrupted || toolInterrupted || publisher.hasProviderError()) - yield* failUnsettled("Tool execution interrupted") - if (settled._tag === "Failure" && !toolInterrupted) { - const failure = Cause.squash(settled.cause) - const message = failure instanceof Error ? failure.message : String(failure) - yield* failUnsettled(`Tool execution failed: ${message}`) - } - if (stream._tag === "Success" && !publisher.hasProviderError()) - yield* failUnsettled("Provider did not return a tool result", true) - if (stream._tag === "Failure") return yield* Effect.failCause(stream.cause) - if (settled._tag === "Failure") return yield* Effect.failCause(settled.cause) - return AttemptResult.cases.Complete.make({ + }), + projectProviderFailure: (failure) => + Effect.gen(function* () { + if (overflowFailure) yield* publish(overflowFailure) + if (failure instanceof LLMError && !publisher.hasProviderError()) { + yield* failUnsettled("Provider did not return a tool result", true) + yield* withPublication( + events.publish(SessionEvent.Step.Failed, { + sessionID: prepared.session.id, + timestamp: yield* DateTime.now, + assistantMessageID: yield* publisher.startAssistant(), + error: { type: "unknown", message: failure.reason.message }, + }), + ) + } + }), + hasProviderError: publisher.hasProviderError, + failUnsettled, + }) + return SettleProviderTurn.Result.match(result, { + RecoveredOverflow: () => AttemptResult.cases.CompactedOverflow.make({}), + Complete: () => + AttemptResult.cases.Complete.make({ needsContinuation: !publisher.hasProviderError() && needsContinuation, - }) - }), - ) + }), + }) }, Effect.scoped) const run = Effect.fn("SessionRunner.runTurn")(function* (input: Input): Effect.fn.Return { diff --git a/packages/core/src/session/runner/settle-provider-turn.ts b/packages/core/src/session/runner/settle-provider-turn.ts new file mode 100644 index 0000000000..2e6302bfb2 --- /dev/null +++ b/packages/core/src/session/runner/settle-provider-turn.ts @@ -0,0 +1,71 @@ +export * as SettleProviderTurn from "./settle-provider-turn" + +import { Cause, Effect, FiberSet, Option, Schema } from "effect" +import { QuestionV2 } from "../../question" +import { ToolOutputStore } from "../../tool-output-store" + +export const Result = Schema.TaggedUnion({ + Complete: {}, + RecoveredOverflow: {}, +}) + +const isQuestionRejected = (cause: Cause.Cause) => + cause.reasons.some((reason) => Cause.isDieReason(reason) && reason.defect instanceof QuestionV2.RejectedError) + +const awaitTools = (fibers: FiberSet.FiberSet) => + Effect.raceFirst(FiberSet.join(fibers), FiberSet.awaitEmpty(fibers)) + +/** + * Runs one provider response together with every local tool it starts. + * + * The provider and tools remain interruptible. Once the provider stops, cleanup + * cannot be interrupted before every started tool is observed, cancelled, or + * settled and its original failure is propagated. + */ +export const run = Effect.fn("SessionRunner.settleProviderTurn")(function* (input: { + readonly stream: ( + runTool: (effect: Effect.Effect) => Effect.Effect, + ) => Effect.Effect + readonly recoverOverflow: ( + failure: unknown, + restore: (effect: Effect.Effect) => Effect.Effect, + ) => Effect.Effect + readonly projectProviderFailure: (failure: unknown) => Effect.Effect + readonly hasProviderError: () => boolean + readonly failUnsettled: (message: string, providerExecuted?: boolean) => Effect.Effect +}) { + const tools = yield* FiberSet.make() + return yield* Effect.uninterruptibleMask((restore) => + Effect.gen(function* () { + const stream = yield* restore(input.stream((effect) => effect.pipe(FiberSet.run(tools)))).pipe(Effect.exit) + const failure = stream._tag === "Failure" ? Option.getOrUndefined(Cause.findErrorOption(stream.cause)) : undefined + if (yield* input.recoverOverflow(failure, restore)) return Result.cases.RecoveredOverflow.make({}) + + yield* input.projectProviderFailure(failure) + const streamInterrupted = stream._tag === "Failure" && Cause.hasInterrupts(stream.cause) + if (streamInterrupted) yield* FiberSet.clear(tools) + const settled = yield* restore(awaitTools(tools)).pipe(Effect.exit) + if (settled._tag === "Failure" && isQuestionRejected(settled.cause)) { + yield* FiberSet.clear(tools) + yield* input.failUnsettled("Tool execution interrupted") + return yield* Effect.interrupt + } + + const toolInterrupted = settled._tag === "Failure" && Cause.hasInterrupts(settled.cause) + if (toolInterrupted) yield* FiberSet.clear(tools) + if (streamInterrupted || toolInterrupted || input.hasProviderError()) + yield* input.failUnsettled("Tool execution interrupted") + if (settled._tag === "Failure" && !toolInterrupted) { + const failure = Cause.squash(settled.cause) + yield* input.failUnsettled( + `Tool execution failed: ${failure instanceof Error ? failure.message : String(failure)}`, + ) + } + if (stream._tag === "Success" && !input.hasProviderError()) + yield* input.failUnsettled("Provider did not return a tool result", true) + if (stream._tag === "Failure") return yield* Effect.failCause(stream.cause) + if (settled._tag === "Failure") return yield* Effect.failCause(settled.cause) + return Result.cases.Complete.make({}) + }), + ) +}, Effect.scoped)