| import { |
| LLM, |
| LLMClient, |
| LLMError, |
| LLMEvent, |
| Message, |
| SystemPart, |
| isContextOverflowFailure, |
| type ProviderErrorEvent, |
| } from "@opencode-ai/llm" |
| import { Cause, DateTime, Effect, FiberSet, Layer, Option, Semaphore, Stream } from "effect" |
| import { AgentV2 } from "../../agent" |
| import { Config } from "../../config" |
| import { Database } from "../../database/database" |
| import { EventV2 } from "../../event" |
| import { Location } from "../../location" |
| import { ModelV2 } from "../../model" |
| import { PermissionV2 } from "../../permission" |
| import { ProviderV2 } from "../../provider" |
| import { QuestionV2 } from "../../question" |
| import { SystemContext } from "../../system-context/index" |
| import { SystemContextRegistry } from "../../system-context/registry" |
| import { SkillGuidance } from "../../skill/guidance" |
| import { ReferenceGuidance } from "../../reference/guidance" |
| import { ToolRegistry } from "../../tool/registry" |
| import { ToolOutputStore } from "../../tool-output-store" |
| import { SessionContextEpoch } from "../context-epoch" |
| import { SessionCompaction } from "../compaction" |
| import { SessionEvent } from "../event" |
| import { SessionHistory } from "../history" |
| import { SessionInput } from "../input" |
| import { SessionSchema } from "../schema" |
| import { SessionStore } from "../store" |
| import { type RunError, Service } from "./index" |
| import { SessionRunnerModel } from "./model" |
| import { createLLMEventPublisher } from "./publish-llm-event" |
| import { toLLMMessages } from "./to-llm-message" |
| import { MAX_STEPS_PROMPT } from "./max-steps" |
| import { Snapshot } from "../../snapshot" |
| import { makeLocationNode } from "../../effect/app-node" |
| import { llmClient } from "../../effect/app-node-platform" |
|
|
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
|
|
| const layer = Layer.effect( |
| Service, |
| Effect.gen(function* () { |
| const events = yield* EventV2.Service |
| const llm = yield* LLMClient.Service |
| const agents = yield* AgentV2.Service |
| const tools = yield* ToolRegistry.Service |
| const models = yield* SessionRunnerModel.Service |
| const store = yield* SessionStore.Service |
| const location = yield* Location.Service |
| const systemContext = yield* SystemContextRegistry.Service |
| const skillGuidance = yield* SkillGuidance.Service |
| const referenceGuidance = yield* ReferenceGuidance.Service |
| const config = yield* Config.Service |
| const snapshots = yield* Snapshot.Service |
| const db = (yield* Database.Service).db |
| const compaction = SessionCompaction.make({ events, llm, config: yield* config.entries() }) |
| const getSession = Effect.fn("SessionRunner.getSession")(function* (sessionID: SessionSchema.ID) { |
| const session = yield* store.get(sessionID) |
| if (!session) return yield* Effect.die(`Session not found: ${sessionID}`) |
| return session |
| }) |
|
|
| const getContext = Effect.fn("SessionRunner.getContext")(function* (sessionID: SessionSchema.ID) { |
| return yield* store.context(sessionID) |
| }) |
| const failInterruptedTools = Effect.fn("SessionRunner.failInterruptedTools")(function* ( |
| sessionID: SessionSchema.ID, |
| ) { |
| for (const message of yield* getContext(sessionID)) { |
| if (message.type !== "assistant") continue |
| for (const tool of message.content) { |
| if (tool.type !== "tool" || (tool.state.status !== "pending" && tool.state.status !== "running")) continue |
| yield* events.publish(SessionEvent.Tool.Failed, { |
| sessionID, |
| timestamp: yield* DateTime.now, |
| assistantMessageID: message.id, |
| callID: tool.id, |
| error: { type: "unknown", message: "Tool execution interrupted" }, |
| provider: { |
| executed: tool.provider?.executed === true, |
| ...(tool.provider?.metadata === undefined ? {} : { metadata: tool.provider.metadata }), |
| }, |
| }) |
| } |
| } |
| }) |
|
|
| const awaitToolFibers = (fibers: FiberSet.FiberSet<void, ToolOutputStore.Error>) => |
| Effect.raceFirst(FiberSet.join(fibers), FiberSet.awaitEmpty(fibers)) |
|
|
| |
| const isUserDeclined = (cause: Cause.Cause<unknown>) => |
| cause.reasons.some( |
| (reason) => |
| Cause.isDieReason(reason) && |
| (reason.defect instanceof PermissionV2.DeclinedError || reason.defect instanceof QuestionV2.RejectedError), |
| ) |
|
|
| type TurnTransition = |
| |
| | { readonly _tag: "ContinueAfterCompaction"; readonly step: number } |
| |
| | { readonly _tag: "ContinueAfterOverflowCompaction"; readonly step: number } |
|
|
| class TurnTransitionError extends Error { |
| constructor(readonly transition: TurnTransition) { |
| super() |
| } |
| } |
|
|
| const continueAfterCompaction = (step: number) => new TurnTransitionError({ _tag: "ContinueAfterCompaction", step }) |
| const continueAfterOverflowCompaction = (step: number) => |
| new TurnTransitionError({ _tag: "ContinueAfterOverflowCompaction", step }) |
|
|
| const loadSystemContext = (agent: AgentV2.Selection) => |
| Effect.all([systemContext.load(), skillGuidance.load(agent), referenceGuidance.load()], { |
| concurrency: "unbounded", |
| }).pipe(Effect.map(SystemContext.combine)) |
|
|
| const runTurnAttempt = Effect.fn("SessionRunner.runTurn")(function* ( |
| sessionID: SessionSchema.ID, |
| promotion: SessionInput.Delivery | undefined, |
| step: number, |
| recoverOverflow?: typeof compaction.compactAfterOverflow, |
| ) { |
| const session = yield* getSession(sessionID) |
| if (session.location.directory !== location.directory || session.location.workspaceID !== location.workspaceID) |
| return yield* Effect.interrupt |
| const agent = yield* agents.select(session.agent) |
| const initialized = yield* SessionContextEpoch.initialize(db, loadSystemContext(agent), session.id) |
| const toolFibers = yield* FiberSet.make<void, ToolOutputStore.Error>() |
| let needsContinuation = false |
| let currentStep = step |
| if (promotion) { |
| const cutoff = yield* EventV2.latestSequence(db, session.id) |
| let promoted = 0 |
| if (promotion === "steer") promoted = yield* SessionInput.promoteSteers(db, events, session.id, cutoff) |
| if (promotion === "queue") { |
| promoted += Number(yield* SessionInput.promoteNextQueued(db, events, session.id)) |
| promoted += yield* SessionInput.promoteSteers(db, events, session.id, cutoff) |
| } |
| if (promoted > 0) currentStep = 1 |
| } |
| const system = |
| initialized ?? (yield* SessionContextEpoch.prepare(db, events, loadSystemContext(agent), session.id)) |
| const model = yield* models.resolve(session) |
| const entries = yield* SessionHistory.entriesForRunner(db, session.id, system.baselineSeq) |
| const context = entries.map((entry) => entry.message) |
| const isLastStep = agent.info?.steps !== undefined && currentStep >= agent.info.steps |
| const toolMaterialization = isLastStep ? undefined : yield* tools.materialize(agent.info?.permissions) |
| const promptCacheKey = /^ses_[0-9a-f]{64}$/.test(session.id) ? session.id.slice(4) : session.id |
| const request = LLM.request({ |
| model, |
| http: { |
| headers: { |
| "x-session-affinity": session.id, |
| "X-Session-Id": session.id, |
| ...(session.parentID ? { "x-parent-session-id": session.parentID } : {}), |
| }, |
| }, |
| providerOptions: { openai: { promptCacheKey } }, |
| system: [agent.info?.system, system.baseline] |
| .filter((part): part is string => part !== undefined && part.length > 0) |
| .map(SystemPart.make), |
| messages: [...toLLMMessages(context, model), ...(isLastStep ? [Message.assistant(MAX_STEPS_PROMPT)] : [])], |
| tools: toolMaterialization?.definitions ?? [], |
| toolChoice: isLastStep ? "none" : undefined, |
| }) |
| if (yield* compaction.compactIfNeeded({ sessionID: session.id, entries, model, request })) |
| return yield* Effect.die(continueAfterCompaction(currentStep)) |
| const startSnapshot = yield* snapshots.capture() |
| const publisher = createLLMEventPublisher(events, { |
| sessionID: session.id, |
| agent: agent.id, |
| model: { |
| id: ModelV2.ID.make(model.id), |
| providerID: ProviderV2.ID.make(model.provider), |
| ...(session.model?.variant === undefined ? {} : { variant: session.model.variant }), |
| }, |
| snapshot: startSnapshot, |
| }) |
| const withPublication = Semaphore.makeUnsafe(1).withPermit |
| const publish = (event: LLMEvent, outputPaths: ReadonlyArray<string> = []) => |
| withPublication(publisher.publish(event, outputPaths)) |
| let overflowFailure: ProviderErrorEvent | undefined |
| const providerStream = llm.stream(request).pipe( |
| Stream.runForEach((event) => |
| Effect.gen(function* () { |
| if (overflowFailure || publisher.hasProviderError()) return |
| if (LLMEvent.is.providerError(event)) { |
| if (isContextOverflowFailure(event) && !publisher.hasAssistantStarted()) { |
| overflowFailure = event |
| return |
| } |
| } |
| yield* publish(event) |
| if (event.type !== "tool-call" || event.providerExecuted) return |
| if (!toolMaterialization) { |
| yield* withPublication(publisher.failUnsettledTools("Tools are disabled after the maximum agent steps")) |
| return |
| } |
| needsContinuation = true |
| const assistantMessageID = yield* publisher.assistantMessageID(event.id) |
| yield* Effect.uninterruptibleMask((restore) => |
| restore( |
| toolMaterialization.settle({ |
| sessionID: session.id, |
| agent: agent.id, |
| assistantMessageID, |
| call: event, |
| }), |
| ).pipe( |
| Effect.flatMap((settlement) => |
| publish( |
| LLMEvent.toolResult({ |
| id: event.id, |
| name: event.name, |
| result: settlement.result, |
| output: settlement.output, |
| }), |
| settlement.outputPaths ?? [], |
| ), |
| ), |
| ), |
| ).pipe(FiberSet.run(toolFibers)) |
| }), |
| ), |
| Effect.ensuring(withPublication(publisher.flush())), |
| ) |
|
|
| 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 ( |
| recoverOverflow && |
| !publisher.hasAssistantStarted() && |
| isContextOverflowFailure(overflowFailure ?? failure) && |
| (yield* restore(recoverOverflow({ sessionID: session.id, entries, model, request }))) |
| ) |
| return yield* Effect.die(continueAfterOverflowCompaction(currentStep)) |
| if (overflowFailure) yield* publish(overflowFailure) |
| const llmFailure = failure instanceof LLMError ? failure : undefined |
| if (llmFailure && !publisher.hasProviderError()) { |
| yield* withPublication(publisher.failUnsettledTools("Provider did not return a tool result", true)) |
| yield* withPublication(publisher.failAssistant(llmFailure.reason.message)) |
| } |
| if (stream._tag === "Failure" && Cause.hasInterrupts(stream.cause)) yield* FiberSet.clear(toolFibers) |
| const settled = yield* restore(awaitToolFibers(toolFibers)).pipe(Effect.exit) |
| if (settled._tag === "Failure" && isUserDeclined(settled.cause)) { |
| yield* FiberSet.clear(toolFibers) |
| yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted")) |
| return yield* Effect.interrupt |
| } |
| if ( |
| (stream._tag === "Failure" && Cause.hasInterrupts(stream.cause)) || |
| (settled._tag === "Failure" && Cause.hasInterrupts(settled.cause)) |
| ) { |
| yield* FiberSet.clear(toolFibers) |
| yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted")) |
| if (publisher.hasActiveAssistant()) |
| yield* withPublication(publisher.failAssistant("Provider turn interrupted")) |
| } |
| if (settled._tag === "Failure" && !Cause.hasInterrupts(settled.cause)) { |
| const failure = Cause.squash(settled.cause) |
| const message = failure instanceof Error ? failure.message : String(failure) |
| yield* withPublication(publisher.failUnsettledTools(`Tool execution failed: ${message}`)) |
| } |
| const stepSettlement = publisher.stepSettlement() |
| if (stepSettlement && !publisher.hasProviderError()) { |
| const endSnapshot = yield* snapshots.capture() |
| const files = |
| startSnapshot && endSnapshot |
| ? yield* snapshots |
| .files({ from: startSnapshot, to: endSnapshot }) |
| .pipe(Effect.catch(() => Effect.succeed(undefined))) |
| : undefined |
| yield* withPublication( |
| events.publish(SessionEvent.Step.Ended, { |
| sessionID: session.id, |
| timestamp: yield* DateTime.now, |
| assistantMessageID: yield* publisher.startAssistant(), |
| finish: stepSettlement.finish, |
| cost: 0, |
| tokens: stepSettlement.tokens, |
| snapshot: endSnapshot, |
| files, |
| }), |
| ) |
| } |
| if (publisher.hasProviderError()) |
| yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted")) |
| if (stream._tag === "Success" && !publisher.hasProviderError()) |
| yield* withPublication(publisher.failUnsettledTools("Provider did not return a tool result", true)) |
| if (stream._tag === "Failure") return yield* Effect.failCause(stream.cause) |
| if (settled._tag === "Failure" && Cause.hasInterrupts(settled.cause)) |
| return yield* Effect.failCause(settled.cause) |
| return { needsContinuation: !publisher.hasProviderError() && needsContinuation, step: currentStep } |
| }), |
| ) |
| }, Effect.scoped) |
| type RunTurn = ( |
| sessionID: SessionSchema.ID, |
| promotion: SessionInput.Delivery | undefined, |
| step: number, |
| ) => Effect.Effect<{ readonly needsContinuation: boolean; readonly step: number }, RunError> |
|
|
| const runAfterOverflowCompaction: RunTurn = Effect.fnUntraced(function* (sessionID, promotion, step) { |
| return yield* runTurnAttempt(sessionID, promotion, step).pipe( |
| Effect.catchDefect( |
| Effect.fnUntraced(function* (defect) { |
| if (!(defect instanceof TurnTransitionError)) return yield* Effect.die(defect) |
| if (defect.transition._tag === "ContinueAfterOverflowCompaction") |
| return yield* Effect.die("Post-compaction provider attempt cannot recover another overflow") |
| yield* Effect.yieldNow |
| return yield* runAfterOverflowCompaction(sessionID, undefined, defect.transition.step) |
| }), |
| ), |
| ) |
| }) |
|
|
| const runTurn: RunTurn = Effect.fnUntraced(function* (sessionID, promotion, step) { |
| return yield* runTurnAttempt(sessionID, promotion, step, compaction.compactAfterOverflow).pipe( |
| Effect.catchDefect( |
| Effect.fnUntraced(function* (defect) { |
| if (!(defect instanceof TurnTransitionError)) return yield* Effect.die(defect) |
| yield* Effect.yieldNow |
| if (defect.transition._tag === "ContinueAfterOverflowCompaction") |
| return yield* runAfterOverflowCompaction(sessionID, undefined, defect.transition.step) |
| return yield* runTurn(sessionID, undefined, defect.transition.step) |
| }), |
| ), |
| ) |
| }) |
|
|
| const run = Effect.fn("SessionRunner.run")(function* (input: { |
| readonly sessionID: SessionSchema.ID |
| readonly force: boolean |
| }) { |
| const hasSteer = yield* SessionInput.hasPending(db, input.sessionID, "steer") |
| const hasQueue = hasSteer ? false : yield* SessionInput.hasPending(db, input.sessionID, "queue") |
| if (!input.force && !hasSteer && !hasQueue) return |
| yield* failInterruptedTools(input.sessionID) |
| let promotion: SessionInput.Delivery | undefined = hasSteer ? "steer" : hasQueue ? "queue" : undefined |
| let shouldRun = input.force || hasSteer || hasQueue |
| while (shouldRun) { |
| let needsContinuation = true |
| let step = 1 |
| while (needsContinuation) { |
| const result = yield* runTurn(input.sessionID, promotion, step) |
| needsContinuation = result.needsContinuation |
| step = result.step + 1 |
| promotion = "steer" |
| if (!needsContinuation) needsContinuation = yield* SessionInput.hasPending(db, input.sessionID, "steer") |
| } |
| shouldRun = yield* SessionInput.hasPending(db, input.sessionID, "queue") |
| promotion = shouldRun ? "queue" : undefined |
| } |
| }) |
|
|
| return Service.of({ |
| run, |
| }) |
| }), |
| ) |
|
|
| export const node = makeLocationNode({ |
| service: Service, |
| layer, |
| deps: [ |
| EventV2.node, |
| llmClient, |
| AgentV2.node, |
| ToolRegistry.node, |
| SessionRunnerModel.node, |
| SessionStore.node, |
| Location.node, |
| SystemContextRegistry.node, |
| SkillGuidance.node, |
| ReferenceGuidance.node, |
| Config.node, |
| Snapshot.node, |
| Database.node, |
| ], |
| }) |
|
|