import { Cause, Deferred, Effect, Exit, Layer, Context, Scope } from "effect" import * as Stream from "effect/Stream" import { Agent } from "@/agent/agent" import { Bus } from "@/bus" import { Config } from "@/config/config" import { Permission } from "@/permission" import { Plugin } from "@/plugin" import { Snapshot } from "@/snapshot" import * as Session from "./session" import { LLM } from "./llm" import { MessageV2 } from "./message-v2" import { Image } from "@/image/image" import { isOverflow } from "./overflow" import { PartID } from "./schema" import type { SessionID } from "./schema" import { SessionRetry } from "./retry" import { SessionStatus } from "./status" import { SessionSummary } from "./summary" import type { Provider } from "@/provider/provider" import { Question } from "@/question" import { errorMessage } from "@/util/error" import * as Log from "@opencode-ai/core/util/log" import { isRecord } from "@/util/record" import { SyncEvent } from "@/sync" import { SessionEvent } from "@/v2/session-event" import { Modelv2 } from "@/v2/model" import * as DateTime from "effect/DateTime" import { Flag } from "@opencode-ai/core/flag/flag" const DOOM_LOOP_THRESHOLD = 3 const log = Log.create({ service: "session.processor" }) export type Result = "compact" | "stop" | "continue" export type Event = LLM.Event export interface Handle { readonly message: MessageV2.Assistant readonly updateToolCall: ( toolCallID: string, update: (part: MessageV2.ToolPart) => MessageV2.ToolPart, ) => Effect.Effect readonly completeToolCall: ( toolCallID: string, output: { title: string metadata: Record output: string attachments?: MessageV2.FilePart[] }, ) => Effect.Effect readonly process: (streamInput: LLM.StreamInput) => Effect.Effect } type Input = { assistantMessage: MessageV2.Assistant sessionID: SessionID model: Provider.Model } export interface Interface { readonly create: (input: Input) => Effect.Effect } type ToolCall = { partID: MessageV2.ToolPart["id"] messageID: MessageV2.ToolPart["messageID"] sessionID: MessageV2.ToolPart["sessionID"] done: Deferred.Deferred } interface ProcessorContext extends Input { toolcalls: Record shouldBreak: boolean snapshot: string | undefined blocked: boolean needsCompaction: boolean currentText: MessageV2.TextPart | undefined reasoningMap: Record } type StreamEvent = Event export class Service extends Context.Service()("@opencode/SessionProcessor") {} export const layer: Layer.Layer< Service, never, | Session.Service | Config.Service | Bus.Service | Snapshot.Service | Agent.Service | LLM.Service | Permission.Service | Plugin.Service | Image.Service | SessionSummary.Service | SessionStatus.Service | SyncEvent.Service > = Layer.effect( Service, Effect.gen(function* () { const session = yield* Session.Service const config = yield* Config.Service const bus = yield* Bus.Service const snapshot = yield* Snapshot.Service const agents = yield* Agent.Service const llm = yield* LLM.Service const permission = yield* Permission.Service const plugin = yield* Plugin.Service const summary = yield* SessionSummary.Service const scope = yield* Scope.Scope const status = yield* SessionStatus.Service const image = yield* Image.Service const sync = yield* SyncEvent.Service const create = Effect.fn("SessionProcessor.create")(function* (input: Input) { // Pre-capture snapshot before the LLM stream starts. The AI SDK // may execute tools internally before emitting start-step events, // so capturing inside the event handler can be too late. const initialSnapshot = yield* snapshot.track() const ctx: ProcessorContext = { assistantMessage: input.assistantMessage, sessionID: input.sessionID, model: input.model, toolcalls: {}, shouldBreak: false, snapshot: initialSnapshot, blocked: false, needsCompaction: false, currentText: undefined, reasoningMap: {}, } let aborted = false const slog = log.clone().tag("session.id", input.sessionID).tag("messageID", input.assistantMessage.id) const parse = (e: unknown) => MessageV2.fromError(e, { providerID: input.model.providerID, aborted, }) const settleToolCall = Effect.fn("SessionProcessor.settleToolCall")(function* (toolCallID: string) { const done = ctx.toolcalls[toolCallID]?.done delete ctx.toolcalls[toolCallID] if (done) yield* Deferred.succeed(done, undefined).pipe(Effect.ignore) }) const readToolCall = Effect.fn("SessionProcessor.readToolCall")(function* (toolCallID: string) { const call = ctx.toolcalls[toolCallID] if (!call) return const part = yield* session.getPart({ partID: call.partID, messageID: call.messageID, sessionID: call.sessionID, }) if (!part || part.type !== "tool") { delete ctx.toolcalls[toolCallID] return } return { call, part } }) const updateToolCall = Effect.fn("SessionProcessor.updateToolCall")(function* ( toolCallID: string, update: (part: MessageV2.ToolPart) => MessageV2.ToolPart, ) { const match = yield* readToolCall(toolCallID) if (!match) return const part = yield* session.updatePart(update(match.part)) ctx.toolcalls[toolCallID] = { ...match.call, partID: part.id, messageID: part.messageID, sessionID: part.sessionID, } return part }) const completeToolCall = Effect.fn("SessionProcessor.completeToolCall")(function* ( toolCallID: string, output: { title: string metadata: Record output: string attachments?: MessageV2.FilePart[] }, ) { const match = yield* readToolCall(toolCallID) if (!match || match.part.state.status !== "running") return yield* session.updatePart({ ...match.part, state: { status: "completed", input: match.part.state.input, output: output.output, metadata: output.metadata, title: output.title, time: { start: match.part.state.time.start, end: Date.now() }, attachments: output.attachments, }, }) yield* settleToolCall(toolCallID) }) const failToolCall = Effect.fn("SessionProcessor.failToolCall")(function* (toolCallID: string, error: unknown) { const match = yield* readToolCall(toolCallID) if (!match || match.part.state.status !== "running") return false yield* session.updatePart({ ...match.part, state: { status: "error", input: match.part.state.input, error: errorMessage(error), time: { start: match.part.state.time.start, end: Date.now() }, }, }) if (error instanceof Permission.RejectedError || error instanceof Question.RejectedError) { ctx.blocked = ctx.shouldBreak } yield* settleToolCall(toolCallID) return true }) const handleEvent = Effect.fnUntraced(function* (value: StreamEvent) { switch (value.type) { case "start": yield* status.set(ctx.sessionID, { type: "busy" }) return case "reasoning-start": if (value.id in ctx.reasoningMap) return // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Reasoning.Started.Sync, { sessionID: ctx.sessionID, reasoningID: value.id, timestamp: DateTime.makeUnsafe(Date.now()), }) } ctx.reasoningMap[value.id] = { id: PartID.ascending(), messageID: ctx.assistantMessage.id, sessionID: ctx.assistantMessage.sessionID, type: "reasoning", text: "", time: { start: Date.now() }, metadata: value.providerMetadata, } yield* session.updatePart(ctx.reasoningMap[value.id]) return case "reasoning-delta": if (!(value.id in ctx.reasoningMap)) return ctx.reasoningMap[value.id].text += value.text if (value.providerMetadata) ctx.reasoningMap[value.id].metadata = value.providerMetadata yield* session.updatePartDelta({ sessionID: ctx.reasoningMap[value.id].sessionID, messageID: ctx.reasoningMap[value.id].messageID, partID: ctx.reasoningMap[value.id].id, field: "text", delta: value.text, }) return case "reasoning-end": if (!(value.id in ctx.reasoningMap)) return // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Reasoning.Ended.Sync, { sessionID: ctx.sessionID, reasoningID: value.id, text: ctx.reasoningMap[value.id].text, timestamp: DateTime.makeUnsafe(Date.now()), }) } // oxlint-disable-next-line no-self-assign -- reactivity trigger ctx.reasoningMap[value.id].text = ctx.reasoningMap[value.id].text ctx.reasoningMap[value.id].time = { ...ctx.reasoningMap[value.id].time, end: Date.now() } if (value.providerMetadata) ctx.reasoningMap[value.id].metadata = value.providerMetadata yield* session.updatePart(ctx.reasoningMap[value.id]) delete ctx.reasoningMap[value.id] return case "tool-input-start": if (ctx.assistantMessage.summary) { throw new Error(`Tool call not allowed while generating summary: ${value.toolName}`) } // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Tool.Input.Started.Sync, { sessionID: ctx.sessionID, callID: value.id, name: value.toolName, timestamp: DateTime.makeUnsafe(Date.now()), }) } const part = yield* session.updatePart({ id: ctx.toolcalls[value.id]?.partID ?? PartID.ascending(), messageID: ctx.assistantMessage.id, sessionID: ctx.assistantMessage.sessionID, type: "tool", tool: value.toolName, callID: value.id, state: { status: "pending", input: {}, raw: "" }, metadata: value.providerExecuted ? { providerExecuted: true } : undefined, } satisfies MessageV2.ToolPart) ctx.toolcalls[value.id] = { done: yield* Deferred.make(), partID: part.id, messageID: part.messageID, sessionID: part.sessionID, } return case "tool-input-delta": return case "tool-input-end": { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Tool.Input.Ended.Sync, { sessionID: ctx.sessionID, callID: value.id, text: "", timestamp: DateTime.makeUnsafe(Date.now()), }) } return } case "tool-call": { if (ctx.assistantMessage.summary) { throw new Error(`Tool call not allowed while generating summary: ${value.toolName}`) } const toolCall = yield* readToolCall(value.toolCallId) // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Tool.Called.Sync, { sessionID: ctx.sessionID, callID: value.toolCallId, tool: value.toolName, input: value.input, provider: { executed: toolCall?.part.metadata?.providerExecuted === true, ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), }, timestamp: DateTime.makeUnsafe(Date.now()), }) } yield* updateToolCall(value.toolCallId, (match) => ({ ...match, tool: value.toolName, state: { ...match.state, status: "running", input: value.input, time: { start: Date.now() }, }, metadata: match.metadata?.providerExecuted ? { ...value.providerMetadata, providerExecuted: true } : value.providerMetadata, })) const parts = MessageV2.parts(ctx.assistantMessage.id) const recentParts = parts.slice(-DOOM_LOOP_THRESHOLD) if ( recentParts.length !== DOOM_LOOP_THRESHOLD || !recentParts.every( (part) => part.type === "tool" && part.tool === value.toolName && part.state.status !== "pending" && JSON.stringify(part.state.input) === JSON.stringify(value.input), ) ) { return } const agent = yield* agents.get(ctx.assistantMessage.agent) yield* permission.ask({ permission: "doom_loop", patterns: [value.toolName], sessionID: ctx.assistantMessage.sessionID, metadata: { tool: value.toolName, input: value.input }, always: [value.toolName], ruleset: agent.permission, }) return } case "tool-result": { const toolCall = yield* readToolCall(value.toolCallId) const toolAttachments: MessageV2.FilePart[] = ( Array.isArray(value.output.attachments) ? value.output.attachments : [] ).filter( (attachment: unknown): attachment is MessageV2.FilePart => isRecord(attachment) && attachment.type === "file" && typeof attachment.mime === "string" && typeof attachment.url === "string", ) const normalized = yield* Effect.forEach(toolAttachments, (attachment) => attachment.mime.startsWith("image/") ? image.normalize(attachment).pipe(Effect.exit) : Effect.succeed(Exit.succeed(attachment)), ) const omitted = normalized.filter(Exit.isFailure).length const attachments = normalized.filter(Exit.isSuccess).map((item) => item.value) const output = { ...value.output, output: omitted === 0 ? value.output.output : `${value.output.output}\n\n[${omitted} image${omitted === 1 ? "" : "s"} omitted: could not be resized below the inline image size limit.]`, attachments: attachments?.length ? attachments : undefined, } // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Tool.Success.Sync, { sessionID: ctx.sessionID, callID: value.toolCallId, structured: output.metadata, content: [ { type: "text", text: output.output, }, ...(output.attachments?.map((item: MessageV2.FilePart) => ({ type: "file", uri: item.url, mime: item.mime, name: item.filename, })) ?? []), ], provider: { executed: toolCall?.part.metadata?.providerExecuted === true, }, timestamp: DateTime.makeUnsafe(Date.now()), }) } yield* completeToolCall(value.toolCallId, output) return } case "tool-error": { const toolCall = yield* readToolCall(value.toolCallId) // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Tool.Failed.Sync, { sessionID: ctx.sessionID, callID: value.toolCallId, error: { type: "unknown", message: errorMessage(value.error), }, provider: { executed: toolCall?.part.metadata?.providerExecuted === true, }, timestamp: DateTime.makeUnsafe(Date.now()), }) } yield* failToolCall(value.toolCallId, value.error) return } case "error": throw value.error case "start-step": if (!ctx.snapshot) ctx.snapshot = yield* snapshot.track() if (!ctx.assistantMessage.summary) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Step.Started.Sync, { sessionID: ctx.sessionID, agent: input.assistantMessage.agent, model: { id: Modelv2.ID.make(ctx.model.id), providerID: Modelv2.ProviderID.make(ctx.model.providerID), variant: Modelv2.VariantID.make(input.assistantMessage.variant ?? "default"), }, snapshot: ctx.snapshot, timestamp: DateTime.makeUnsafe(Date.now()), }) } } yield* session.updatePart({ id: PartID.ascending(), messageID: ctx.assistantMessage.id, sessionID: ctx.sessionID, snapshot: ctx.snapshot, type: "step-start", }) return case "finish-step": { const completedSnapshot = yield* snapshot.track() const usage = Session.getUsage({ model: ctx.model, usage: value.usage, metadata: value.providerMetadata, }) if (!ctx.assistantMessage.summary) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Step.Ended.Sync, { sessionID: ctx.sessionID, finish: value.finishReason, cost: usage.cost, tokens: usage.tokens, snapshot: completedSnapshot, timestamp: DateTime.makeUnsafe(Date.now()), }) } } ctx.assistantMessage.finish = value.finishReason ctx.assistantMessage.cost += usage.cost ctx.assistantMessage.tokens = usage.tokens yield* session.updatePart({ id: PartID.ascending(), reason: value.finishReason, snapshot: completedSnapshot, messageID: ctx.assistantMessage.id, sessionID: ctx.assistantMessage.sessionID, type: "step-finish", tokens: usage.tokens, cost: usage.cost, }) yield* session.updateMessage(ctx.assistantMessage) if (ctx.snapshot) { const patch = yield* snapshot.patch(ctx.snapshot) if (patch.files.length) { yield* session.updatePart({ id: PartID.ascending(), messageID: ctx.assistantMessage.id, sessionID: ctx.sessionID, type: "patch", hash: patch.hash, files: patch.files, }) } ctx.snapshot = undefined } yield* summary .summarize({ sessionID: ctx.sessionID, messageID: ctx.assistantMessage.parentID, }) .pipe(Effect.ignore, Effect.forkIn(scope)) if ( !ctx.assistantMessage.summary && isOverflow({ cfg: yield* config.get(), tokens: usage.tokens, model: ctx.model }) ) { ctx.needsCompaction = true } return } case "text-start": if (!ctx.assistantMessage.summary) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Text.Started.Sync, { sessionID: ctx.sessionID, timestamp: DateTime.makeUnsafe(Date.now()), }) } } ctx.currentText = { id: PartID.ascending(), messageID: ctx.assistantMessage.id, sessionID: ctx.assistantMessage.sessionID, type: "text", text: "", time: { start: Date.now() }, metadata: value.providerMetadata, } yield* session.updatePart(ctx.currentText) return case "text-delta": if (!ctx.currentText) return ctx.currentText.text += value.text if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata yield* session.updatePartDelta({ sessionID: ctx.currentText.sessionID, messageID: ctx.currentText.messageID, partID: ctx.currentText.id, field: "text", delta: value.text, }) return case "text-end": if (!ctx.currentText) return // oxlint-disable-next-line no-self-assign -- reactivity trigger ctx.currentText.text = ctx.currentText.text ctx.currentText.text = (yield* plugin.trigger( "experimental.text.complete", { sessionID: ctx.sessionID, messageID: ctx.assistantMessage.id, partID: ctx.currentText.id, }, { text: ctx.currentText.text }, )).text if (!ctx.assistantMessage.summary) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Text.Ended.Sync, { sessionID: ctx.sessionID, text: ctx.currentText.text, timestamp: DateTime.makeUnsafe(Date.now()), }) } } { const end = Date.now() ctx.currentText.time = { start: ctx.currentText.time?.start ?? end, end } } if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata yield* session.updatePart(ctx.currentText) ctx.currentText = undefined return case "finish": return default: slog.info("unhandled", { event: value.type, value }) return } }) const cleanup = Effect.fn("SessionProcessor.cleanup")(function* () { if (ctx.snapshot) { const patch = yield* snapshot.patch(ctx.snapshot) if (patch.files.length) { yield* session.updatePart({ id: PartID.ascending(), messageID: ctx.assistantMessage.id, sessionID: ctx.sessionID, type: "patch", hash: patch.hash, files: patch.files, }) } ctx.snapshot = undefined } if (ctx.currentText) { const end = Date.now() ctx.currentText.time = { start: ctx.currentText.time?.start ?? end, end } yield* session.updatePart(ctx.currentText) ctx.currentText = undefined } for (const part of Object.values(ctx.reasoningMap)) { const end = Date.now() yield* session.updatePart({ ...part, time: { start: part.time.start ?? end, end }, }) } ctx.reasoningMap = {} yield* Effect.forEach( Object.values(ctx.toolcalls), (call) => Deferred.await(call.done).pipe(Effect.timeout("250 millis"), Effect.ignore), { concurrency: "unbounded" }, ) for (const toolCallID of Object.keys(ctx.toolcalls)) { const match = yield* readToolCall(toolCallID) if (!match) continue const part = match.part const end = Date.now() const metadata = "metadata" in part.state && isRecord(part.state.metadata) ? part.state.metadata : {} yield* session.updatePart({ ...part, state: { ...part.state, status: "error", error: "Tool execution aborted", metadata: { ...metadata, interrupted: true }, time: { start: "time" in part.state ? part.state.time.start : end, end }, }, }) } ctx.toolcalls = {} ctx.assistantMessage.time.completed = Date.now() yield* session.updateMessage(ctx.assistantMessage) }) const halt = Effect.fn("SessionProcessor.halt")(function* (e: unknown) { slog.error("process", { error: errorMessage(e), stack: e instanceof Error ? e.stack : undefined }) const error = parse(e) if (MessageV2.ContextOverflowError.isInstance(error)) { ctx.needsCompaction = true yield* bus.publish(Session.Event.Error, { sessionID: ctx.sessionID, error }) return } if (!ctx.assistantMessage.summary) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM) { yield* sync.run(SessionEvent.Step.Failed.Sync, { sessionID: ctx.sessionID, error: { type: "unknown", message: errorMessage(e), }, timestamp: DateTime.makeUnsafe(Date.now()), }) } } ctx.assistantMessage.error = error yield* bus.publish(Session.Event.Error, { sessionID: ctx.assistantMessage.sessionID, error: ctx.assistantMessage.error, }) yield* status.set(ctx.sessionID, { type: "idle" }) }) const process = Effect.fn("SessionProcessor.process")(function* (streamInput: LLM.StreamInput) { slog.info("process") ctx.needsCompaction = false ctx.shouldBreak = (yield* config.get()).experimental?.continue_loop_on_deny !== true return yield* Effect.gen(function* () { yield* Effect.gen(function* () { ctx.currentText = undefined ctx.reasoningMap = {} const stream = llm.stream(streamInput) yield* stream.pipe( Stream.tap((event) => handleEvent(event)), Stream.takeUntil(() => ctx.needsCompaction), Stream.runDrain, ) }).pipe( Effect.onInterrupt(() => Effect.gen(function* () { aborted = true if (!ctx.assistantMessage.error) { yield* halt(new DOMException("Aborted", "AbortError")) } }), ), Effect.catchCauseIf( (cause) => !Cause.hasInterruptsOnly(cause), (cause) => Effect.fail(Cause.squash(cause)), ), Effect.retry( SessionRetry.policy({ provider: input.model.providerID, parse, set: (info) => { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. const event = Flag.OPENCODE_EXPERIMENTAL_EVENT_SYSTEM ? sync.run(SessionEvent.Retried.Sync, { sessionID: ctx.sessionID, attempt: info.attempt, error: { message: info.message, isRetryable: true, }, timestamp: DateTime.makeUnsafe(Date.now()), }) : Effect.void return event.pipe( Effect.andThen( status.set(ctx.sessionID, { type: "retry", attempt: info.attempt, message: info.message, action: info.action, next: info.next, }), ), ) }, }), ), Effect.catch(halt), Effect.ensuring(cleanup()), ) if (ctx.needsCompaction) return "compact" if (ctx.blocked || ctx.assistantMessage.error) return "stop" return "continue" }) }) return { get message() { return ctx.assistantMessage }, updateToolCall, completeToolCall, process, } satisfies Handle }) return Service.of({ create }) }), ) export const defaultLayer = Layer.suspend(() => layer.pipe( Layer.provide(Session.defaultLayer), Layer.provide(Snapshot.defaultLayer), Layer.provide(Agent.defaultLayer), Layer.provide(LLM.defaultLayer), Layer.provide(Permission.defaultLayer), Layer.provide(Plugin.defaultLayer), Layer.provide(SessionSummary.defaultLayer), Layer.provide(SessionStatus.defaultLayer), Layer.provide(Image.defaultLayer), Layer.provide(Bus.layer), Layer.provide(Config.defaultLayer), Layer.provide(SyncEvent.defaultLayer), ), ) export * as SessionProcessor from "./processor"