diff --git a/packages/opencode/src/session/compaction.ts b/packages/opencode/src/session/compaction.ts index f9ee565490..69c7ee8fc3 100644 --- a/packages/opencode/src/session/compaction.ts +++ b/packages/opencode/src/session/compaction.ts @@ -207,7 +207,7 @@ export namespace SessionCompaction { created: Date.now(), }, })) as MessageV2.Assistant - const processor = SessionProcessor.create({ + const processor = await SessionProcessor.create({ assistantMessage: msg, sessionID: input.sessionID, model, diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 30fbb1bb26..5dea925512 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -1,36 +1,55 @@ -import { MessageV2 } from "./message-v2" +import { Cause, Effect, Exit, Layer, ServiceMap } from "effect" +import * as Stream from "effect/Stream" +import { Agent } from "@/agent/agent" +import { Bus } from "@/bus" +import { makeRuntime } from "@/effect/run-service" +import { Config } from "@/config/config" +import { Permission } from "@/permission" +import { Plugin } from "@/plugin" +import { Snapshot } from "@/snapshot" import { Log } from "@/util/log" import { Session } from "." -import { Agent } from "@/agent/agent" -import { Snapshot } from "@/snapshot" -import { SessionSummary } from "./summary" -import { Bus } from "@/bus" +import { LLM } from "./llm" +import { MessageV2 } from "./message-v2" +import { PartID } from "./schema" +import type { SessionID } from "./schema" +import { SessionCompaction } from "./compaction" import { SessionRetry } from "./retry" import { SessionStatus } from "./status" -import { Plugin } from "@/plugin" +import { SessionSummary } from "./summary" import type { Provider } from "@/provider/provider" -import { LLM } from "./llm" -import { Config } from "@/config/config" -import { SessionCompaction } from "./compaction" -import { Permission } from "@/permission" import { Question } from "@/question" -import { PartID } from "./schema" -import type { SessionID, MessageID } from "./schema" -import { Cause, Effect, Exit } from "effect" -import * as Stream from "effect/Stream" export namespace SessionProcessor { const DOOM_LOOP_THRESHOLD = 3 const log = Log.create({ service: "session.processor" }) - export type Info = Awaited> - export type Result = Awaited> + export type Result = "compact" | "stop" | "continue" - interface ProcessorContext { + export interface Handle { + readonly message: MessageV2.Assistant + readonly partFromToolCall: (toolCallID: string) => MessageV2.ToolPart | undefined + readonly process: (streamInput: LLM.StreamInput) => Effect.Effect + } + + export interface Info { + readonly message: MessageV2.Assistant + readonly partFromToolCall: (toolCallID: string) => MessageV2.ToolPart | undefined + readonly process: (streamInput: LLM.StreamInput) => Promise + } + + type Input = { assistantMessage: MessageV2.Assistant sessionID: SessionID model: Provider.Model abort: AbortSignal + } + + export interface Interface { + readonly create: (input: Input) => Effect.Effect + } + + interface ProcessorContext extends Input { toolcalls: Record shouldBreak: boolean snapshot: string | undefined @@ -43,353 +62,42 @@ export namespace SessionProcessor { type StreamResult = Awaited> type StreamEvent = StreamResult["fullStream"] extends AsyncIterable ? T : never - function handleEvent(value: StreamEvent, ctx: ProcessorContext) { - return Effect.promise(async () => { - switch (value.type) { - case "start": - await SessionStatus.set(ctx.sessionID, { type: "busy" }) - break + export class Service extends ServiceMap.Service()("@opencode/SessionProcessor") {} - case "reasoning-start": - if (value.id in ctx.reasoningMap) break - const reasoningPart = { - id: PartID.ascending(), - messageID: ctx.assistantMessage.id, - sessionID: ctx.assistantMessage.sessionID, - type: "reasoning" as const, - text: "", - time: { start: Date.now() }, - metadata: value.providerMetadata, - } - ctx.reasoningMap[value.id] = reasoningPart - await Session.updatePart(reasoningPart) - break - - case "reasoning-delta": - if (value.id in ctx.reasoningMap) { - const part = ctx.reasoningMap[value.id] - part.text += value.text - if (value.providerMetadata) part.metadata = value.providerMetadata - await Session.updatePartDelta({ - sessionID: part.sessionID, - messageID: part.messageID, - partID: part.id, - field: "text", - delta: value.text, - }) - } - break - - case "reasoning-end": - if (value.id in ctx.reasoningMap) { - const part = ctx.reasoningMap[value.id] - part.text = part.text.trimEnd() - part.time = { ...part.time, end: Date.now() } - if (value.providerMetadata) part.metadata = value.providerMetadata - await Session.updatePart(part) - delete ctx.reasoningMap[value.id] - } - break - - case "tool-input-start": { - const part = await Session.updatePart({ - id: ctx.toolcalls[value.id]?.id ?? PartID.ascending(), - messageID: ctx.assistantMessage.id, - sessionID: ctx.assistantMessage.sessionID, - type: "tool", - tool: value.toolName, - callID: value.id, - state: { status: "pending", input: {}, raw: "" }, - }) - ctx.toolcalls[value.id] = part as MessageV2.ToolPart - break - } - - case "tool-input-delta": - break - - case "tool-input-end": - break - - case "tool-call": { - const match = ctx.toolcalls[value.toolCallId] - if (match) { - const part = await Session.updatePart({ - ...match, - tool: value.toolName, - state: { status: "running", input: value.input, time: { start: Date.now() } }, - metadata: value.providerMetadata, - }) - ctx.toolcalls[value.toolCallId] = part as MessageV2.ToolPart - - const parts = await MessageV2.parts(ctx.assistantMessage.id) - const recentParts = parts.slice(-DOOM_LOOP_THRESHOLD) - - if ( - recentParts.length === DOOM_LOOP_THRESHOLD && - recentParts.every( - (p) => - p.type === "tool" && - p.tool === value.toolName && - p.state.status !== "pending" && - JSON.stringify(p.state.input) === JSON.stringify(value.input), - ) - ) { - const agent = await Agent.get(ctx.assistantMessage.agent) - await 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, - }) - } - } - break - } - - case "tool-result": { - const match = ctx.toolcalls[value.toolCallId] - if (match && match.state.status === "running") { - await Session.updatePart({ - ...match, - state: { - status: "completed", - input: value.input ?? match.state.input, - output: value.output.output, - metadata: value.output.metadata, - title: value.output.title, - time: { start: match.state.time.start, end: Date.now() }, - attachments: value.output.attachments, - }, - }) - delete ctx.toolcalls[value.toolCallId] - } - break - } - - case "tool-error": { - const match = ctx.toolcalls[value.toolCallId] - if (match && match.state.status === "running") { - await Session.updatePart({ - ...match, - state: { - status: "error", - input: value.input ?? match.state.input, - error: value.error instanceof Error ? value.error.message : String(value.error), - time: { start: match.state.time.start, end: Date.now() }, - }, - }) - if (value.error instanceof Permission.RejectedError || value.error instanceof Question.RejectedError) { - ctx.blocked = ctx.shouldBreak - } - delete ctx.toolcalls[value.toolCallId] - } - break - } - - case "error": - throw value.error - - case "start-step": - ctx.snapshot = await Snapshot.track() - await Session.updatePart({ - id: PartID.ascending(), - messageID: ctx.assistantMessage.id, - sessionID: ctx.sessionID, - snapshot: ctx.snapshot, - type: "step-start", - }) - break - - case "finish-step": { - const usage = Session.getUsage({ - model: ctx.model, - usage: value.usage, - metadata: value.providerMetadata, - }) - ctx.assistantMessage.finish = value.finishReason - ctx.assistantMessage.cost += usage.cost - ctx.assistantMessage.tokens = usage.tokens - await Session.updatePart({ - id: PartID.ascending(), - reason: value.finishReason, - snapshot: await Snapshot.track(), - messageID: ctx.assistantMessage.id, - sessionID: ctx.assistantMessage.sessionID, - type: "step-finish", - tokens: usage.tokens, - cost: usage.cost, - }) - await Session.updateMessage(ctx.assistantMessage) - if (ctx.snapshot) { - const patch = await Snapshot.patch(ctx.snapshot) - if (patch.files.length) { - await Session.updatePart({ - id: PartID.ascending(), - messageID: ctx.assistantMessage.id, - sessionID: ctx.sessionID, - type: "patch", - hash: patch.hash, - files: patch.files, - }) - } - ctx.snapshot = undefined - } - SessionSummary.summarize({ - sessionID: ctx.sessionID, - messageID: ctx.assistantMessage.parentID, - }) - if ( - !ctx.assistantMessage.summary && - (await SessionCompaction.isOverflow({ tokens: usage.tokens, model: ctx.model })) - ) { - ctx.needsCompaction = true - } - break - } - - case "text-start": - ctx.currentText = { - id: PartID.ascending(), - messageID: ctx.assistantMessage.id, - sessionID: ctx.assistantMessage.sessionID, - type: "text", - text: "", - time: { start: Date.now() }, - metadata: value.providerMetadata, - } - await Session.updatePart(ctx.currentText) - break - - case "text-delta": - if (ctx.currentText) { - ctx.currentText.text += value.text - if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata - await Session.updatePartDelta({ - sessionID: ctx.currentText.sessionID, - messageID: ctx.currentText.messageID, - partID: ctx.currentText.id, - field: "text", - delta: value.text, - }) - } - break - - case "text-end": - if (ctx.currentText) { - ctx.currentText.text = ctx.currentText.text.trimEnd() - const textOutput = await Plugin.trigger( - "experimental.text.complete", - { - sessionID: ctx.sessionID, - messageID: ctx.assistantMessage.id, - partID: ctx.currentText.id, - }, - { text: ctx.currentText.text }, - ) - ctx.currentText.text = textOutput.text - ctx.currentText.time = { start: Date.now(), end: Date.now() } - if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata - await Session.updatePart(ctx.currentText) - } - ctx.currentText = undefined - break - - case "finish": - break - - default: - log.info("unhandled", { ...value }) - break - } - }) - } - - function cleanupEffect(ctx: ProcessorContext) { - return Effect.promise(async () => { - if (ctx.snapshot) { - const patch = await Snapshot.patch(ctx.snapshot) - if (patch.files.length) { - await Session.updatePart({ - id: PartID.ascending(), - messageID: ctx.assistantMessage.id, - sessionID: ctx.sessionID, - type: "patch", - hash: patch.hash, - files: patch.files, - }) - } - ctx.snapshot = undefined - } - const parts = await MessageV2.parts(ctx.assistantMessage.id) - for (const part of parts) { - if (part.type === "tool" && part.state.status !== "completed" && part.state.status !== "error") { - await Session.updatePart({ - ...part, - state: { - ...part.state, - status: "error", - error: "Tool execution aborted", - time: { start: Date.now(), end: Date.now() }, - }, - }) - } - } - ctx.assistantMessage.time.completed = Date.now() - await Session.updateMessage(ctx.assistantMessage) - }) - } - - export function create(input: { - assistantMessage: MessageV2.Assistant - sessionID: SessionID - model: Provider.Model - abort: AbortSignal - }) { - const toolcalls: Record = {} - let snapshot: string | undefined - let blocked = false - let needsCompaction = false - - const result = { - get message() { - return input.assistantMessage - }, - partFromToolCall(toolCallID: string) { - return toolcalls[toolCallID] - }, - async process(streamInput: LLM.StreamInput): Promise<"compact" | "stop" | "continue"> { - log.info("process") - needsCompaction = false - const shouldBreak = (await Config.get()).experimental?.continue_loop_on_deny !== true + export const layer: Layer.Layer< + Service, + never, + | Session.Service + | Config.Service + | Bus.Service + | Snapshot.Service + | Agent.Service + | Permission.Service + | Plugin.Service + | SessionStatus.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 permission = yield* Permission.Service + const plugin = yield* Plugin.Service + const status = yield* SessionStatus.Service + const create = Effect.fn("SessionProcessor.create")(function* (input: Input) { const ctx: ProcessorContext = { assistantMessage: input.assistantMessage, sessionID: input.sessionID, model: input.model, abort: input.abort, - toolcalls, - shouldBreak, - get snapshot() { - return snapshot - }, - set snapshot(v) { - snapshot = v - }, - get blocked() { - return blocked - }, - set blocked(v) { - blocked = v - }, - get needsCompaction() { - return needsCompaction - }, - set needsCompaction(v) { - needsCompaction = v - }, + toolcalls: {}, + shouldBreak: false, + snapshot: undefined, + blocked: false, + needsCompaction: false, currentText: undefined, reasoningMap: {}, } @@ -400,63 +108,391 @@ export namespace SessionProcessor { aborted: input.abort.aborted, }) - const consumeStream = Effect.gen(function* () { - ctx.currentText = undefined - ctx.reasoningMap = {} - const stream = yield* Effect.promise(() => LLM.stream(streamInput)) + const handleEvent = Effect.fn("SessionProcessor.handleEvent")(function* (value: StreamEvent) { + switch (value.type) { + case "start": + yield* status.set(ctx.sessionID, { type: "busy" }) + return - yield* Stream.fromAsyncIterable(stream.fullStream, (e) => e).pipe( - Stream.runForEachWhile((event) => - Effect.gen(function* () { - input.abort.throwIfAborted() - yield* handleEvent(event, ctx) - return !needsCompaction - }), - ), - ) - }) + case "reasoning-start": + if (value.id in ctx.reasoningMap) return + 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 - const halt = (e: unknown): Effect.Effect => - Effect.gen(function* () { - log.error("process", { error: e, stack: JSON.stringify((e as any)?.stack) }) - const error = MessageV2.fromError(e, { - providerID: input.model.providerID, - aborted: input.abort.aborted, - }) - if (MessageV2.ContextOverflowError.isInstance(error)) { - needsCompaction = true - Bus.publish(Session.Event.Error, { sessionID: input.sessionID, error }) + 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 + ctx.reasoningMap[value.id].text = ctx.reasoningMap[value.id].text.trimEnd() + 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": + ctx.toolcalls[value.id] = (yield* session.updatePart({ + id: ctx.toolcalls[value.id]?.id ?? PartID.ascending(), + messageID: ctx.assistantMessage.id, + sessionID: ctx.assistantMessage.sessionID, + type: "tool", + tool: value.toolName, + callID: value.id, + state: { status: "pending", input: {}, raw: "" }, + })) as MessageV2.ToolPart + return + + case "tool-input-delta": + return + + case "tool-input-end": + return + + case "tool-call": { + const match = ctx.toolcalls[value.toolCallId] + if (!match) return + ctx.toolcalls[value.toolCallId] = (yield* session.updatePart({ + ...match, + tool: value.toolName, + state: { status: "running", input: value.input, time: { start: Date.now() } }, + metadata: value.providerMetadata, + })) as MessageV2.ToolPart + + const parts = yield* Effect.promise(() => 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 } - input.assistantMessage.error = error - Bus.publish(Session.Event.Error, { - sessionID: input.assistantMessage.sessionID, - error: input.assistantMessage.error, + + case "tool-result": { + const match = ctx.toolcalls[value.toolCallId] + if (!match || match.state.status !== "running") return + yield* session.updatePart({ + ...match, + state: { + status: "completed", + input: value.input ?? match.state.input, + output: value.output.output, + metadata: value.output.metadata, + title: value.output.title, + time: { start: match.state.time.start, end: Date.now() }, + attachments: value.output.attachments, + }, + }) + delete ctx.toolcalls[value.toolCallId] + return + } + + case "tool-error": { + const match = ctx.toolcalls[value.toolCallId] + if (!match || match.state.status !== "running") return + yield* session.updatePart({ + ...match, + state: { + status: "error", + input: value.input ?? match.state.input, + error: value.error instanceof Error ? value.error.message : String(value.error), + time: { start: match.state.time.start, end: Date.now() }, + }, + }) + if (value.error instanceof Permission.RejectedError || value.error instanceof Question.RejectedError) { + ctx.blocked = ctx.shouldBreak + } + delete ctx.toolcalls[value.toolCallId] + return + } + + case "error": + throw value.error + + case "start-step": + ctx.snapshot = yield* snapshot.track() + yield* session.updatePart({ + id: PartID.ascending(), + messageID: ctx.assistantMessage.id, + sessionID: ctx.sessionID, + snapshot: ctx.snapshot, + type: "step-start", + }) + return + + case "finish-step": { + const usage = Session.getUsage({ + model: ctx.model, + usage: value.usage, + metadata: value.providerMetadata, + }) + 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: yield* snapshot.track(), + 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* Effect.sync(() => { + void SessionSummary.summarize({ + sessionID: ctx.sessionID, + messageID: ctx.assistantMessage.parentID, + }) + }) + if ( + !ctx.assistantMessage.summary && + (yield* Effect.promise(() => SessionCompaction.isOverflow({ tokens: usage.tokens, model: ctx.model }))) + ) { + ctx.needsCompaction = true + } + return + } + + case "text-start": + 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 + ctx.currentText.text = ctx.currentText.text.trimEnd() + 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 + ctx.currentText.time = { start: Date.now(), end: Date.now() } + if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata + yield* session.updatePart(ctx.currentText) + ctx.currentText = undefined + return + + case "finish": + return + + default: + log.info("unhandled", { ...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 + } + const parts = yield* Effect.promise(() => MessageV2.parts(ctx.assistantMessage.id)) + for (const part of parts) { + if (part.type !== "tool" || part.state.status === "completed" || part.state.status === "error") continue + yield* session.updatePart({ + ...part, + state: { + ...part.state, + status: "error", + error: "Tool execution aborted", + time: { start: Date.now(), end: Date.now() }, + }, }) - yield* Effect.promise(() => SessionStatus.set(input.sessionID, { type: "idle" })) + } + ctx.assistantMessage.time.completed = Date.now() + yield* session.updateMessage(ctx.assistantMessage) + }) + + const halt = Effect.fn("SessionProcessor.halt")(function* (e: unknown) { + log.error("process", { error: e, stack: JSON.stringify((e as any)?.stack) }) + const error = parse(e) + if (MessageV2.ContextOverflowError.isInstance(error)) { + ctx.needsCompaction = true + yield* bus.publish(Session.Event.Error, { sessionID: ctx.sessionID, error }) + return + } + 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 loop: Effect.Effect = consumeStream.pipe( - Effect.catchCause((cause) => Effect.fail(Cause.squash(cause))), - Effect.retry(SessionRetry.policy({ sessionID: input.sessionID, parse })), - Effect.ensuring(cleanupEffect(ctx)), - ) + const process = Effect.fn("SessionProcessor.process")(function* (streamInput: LLM.StreamInput) { + log.info("process") + ctx.needsCompaction = false + ctx.shouldBreak = (yield* config.get()).experimental?.continue_loop_on_deny !== true - const exit = await Effect.runPromiseExit(loop, { signal: input.abort }) + yield* Effect.gen(function* () { + ctx.currentText = undefined + ctx.reasoningMap = {} + const stream = yield* Effect.promise(() => LLM.stream(streamInput)) + yield* Stream.fromAsyncIterable(stream.fullStream, (e) => e).pipe( + Stream.runForEachWhile((event) => + Effect.gen(function* () { + input.abort.throwIfAborted() + yield* handleEvent(event) + return !ctx.needsCompaction + }), + ), + ) + }).pipe( + Effect.catchCause((cause) => Effect.fail(Cause.squash(cause))), + Effect.retry(SessionRetry.policy({ sessionID: ctx.sessionID, parse })), + Effect.catchCauseIf( + (cause) => !Cause.hasInterruptsOnly(cause), + (cause) => halt(Cause.squash(cause)), + ), + Effect.catchCause(() => Effect.interrupt), + Effect.onInterrupt(() => halt(new DOMException("Aborted", "AbortError"))), + Effect.ensuring(cleanup()), + ) + + if (ctx.needsCompaction) return "compact" + if (ctx.blocked || ctx.assistantMessage.error) return "stop" + return "continue" + }) + + return { + get message() { + return ctx.assistantMessage + }, + partFromToolCall(toolCallID: string) { + return ctx.toolcalls[toolCallID] + }, + process, + } satisfies Handle + }) + + return Service.of({ create }) + }), + ) + + export const defaultLayer = Layer.unwrap( + Effect.sync(() => + layer.pipe( + Layer.provide(Session.defaultLayer), + Layer.provide(Snapshot.defaultLayer), + Layer.provide(Agent.defaultLayer), + Layer.provide(Permission.layer), + Layer.provide(Plugin.defaultLayer), + Layer.provide(SessionStatus.layer.pipe(Layer.provide(Bus.layer))), + Layer.provide(Bus.layer), + Layer.provide(Config.defaultLayer), + ), + ), + ) + + const { runPromise } = makeRuntime(Service, defaultLayer) + + export async function create(input: Input): Promise { + const hit = await runPromise((svc) => svc.create(input)) + return { + get message() { + return hit.message + }, + partFromToolCall(toolCallID: string) { + return hit.partFromToolCall(toolCallID) + }, + async process(streamInput: LLM.StreamInput) { + const exit = await Effect.runPromiseExit(hit.process(streamInput), { signal: input.abort }) if (Exit.isFailure(exit)) { - const err = - Cause.hasInterrupts(exit.cause) && input.abort.aborted - ? new DOMException("Aborted", "AbortError") - : Cause.squash(exit.cause) - await Effect.runPromise(halt(err)) + if (Cause.hasInterrupts(exit.cause) && input.abort.aborted) return "stop" + throw Cause.squash(exit.cause) } - - if (needsCompaction) return "compact" - if (blocked || input.assistantMessage.error) return "stop" - return "continue" + return exit.value }, } - return result } } diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index dd74b83f50..acc9f63595 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -594,7 +594,7 @@ export namespace SessionPrompt { session, }) - const processor = SessionProcessor.create({ + const processor = await SessionProcessor.create({ assistantMessage: (await Session.updateMessage({ id: MessageID.ascending(), parentID: lastUser.id, diff --git a/packages/opencode/src/session/retry.ts b/packages/opencode/src/session/retry.ts index d2df7f5772..37069a0016 100644 --- a/packages/opencode/src/session/retry.ts +++ b/packages/opencode/src/session/retry.ts @@ -11,7 +11,7 @@ export namespace SessionRetry { export const RETRY_INITIAL_DELAY = 2000 export const RETRY_BACKOFF_FACTOR = 2 export const RETRY_MAX_DELAY_NO_HEADERS = 30_000 // 30 seconds - export const RETRY_MAX_DELAY = 2_147_483_647 + export const RETRY_MAX_DELAY = 2_147_483_647 // max 32-bit signed integer for setTimeout function cap(ms: number) { return Math.min(ms, RETRY_MAX_DELAY)