From 83d07b7e169041e95651e8a7a05eaac0c4f198a3 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sun, 10 May 2026 12:17:10 -0400 Subject: [PATCH 01/13] refactor(session): consume native LLM events --- bun.lock | 1 + packages/opencode/package.json | 1 + packages/opencode/src/session/llm-ai-sdk.ts | 223 ++++++++++++ packages/opencode/src/session/llm.ts | 13 +- packages/opencode/src/session/processor.ts | 318 +++++++++++------- packages/opencode/src/session/prompt.ts | 3 +- packages/opencode/src/session/session.ts | 12 +- .../opencode/test/session/compaction.test.ts | 275 +++------------ 8 files changed, 485 insertions(+), 361 deletions(-) create mode 100644 packages/opencode/src/session/llm-ai-sdk.ts diff --git a/bun.lock b/bun.lock index 4268e5fb7d..58a9aa8fc8 100644 --- a/bun.lock +++ b/bun.lock @@ -396,6 +396,7 @@ "@octokit/graphql": "9.0.2", "@octokit/rest": "catalog:", "@openauthjs/openauth": "catalog:", + "@opencode-ai/llm": "workspace:*", "@opencode-ai/plugin": "workspace:*", "@opencode-ai/script": "workspace:*", "@opencode-ai/sdk": "workspace:*", diff --git a/packages/opencode/package.json b/packages/opencode/package.json index e9b811fc5e..3f63a85bf1 100644 --- a/packages/opencode/package.json +++ b/packages/opencode/package.json @@ -105,6 +105,7 @@ "@octokit/graphql": "9.0.2", "@octokit/rest": "catalog:", "@openauthjs/openauth": "catalog:", + "@opencode-ai/llm": "workspace:*", "@opencode-ai/plugin": "workspace:*", "@opencode-ai/script": "workspace:*", "@opencode-ai/sdk": "workspace:*", diff --git a/packages/opencode/src/session/llm-ai-sdk.ts b/packages/opencode/src/session/llm-ai-sdk.ts new file mode 100644 index 0000000000..9d5407f975 --- /dev/null +++ b/packages/opencode/src/session/llm-ai-sdk.ts @@ -0,0 +1,223 @@ +import { ContentBlockID, FinishReason, LLMEvent, ProviderMetadata, ToolCallID, ToolResultValue, Usage } from "@opencode-ai/llm" +import { Effect, Schema } from "effect" +import { type streamText } from "ai" +import { errorMessage } from "@/util/error" + +type Result = Awaited> +type AISDKEvent = Result["fullStream"] extends AsyncIterable ? T : never + +export function adapterState() { + return { + step: 0, + text: 0, + reasoning: 0, + currentTextID: undefined as ContentBlockID | undefined, + currentReasoningID: undefined as ContentBlockID | undefined, + toolNames: {} as Record, + } +} + +const contentBlockID = (value: string) => ContentBlockID.make(value) +const toolCallID = (value: string) => ToolCallID.make(value) + +function finishReason(value: string | undefined): FinishReason { + return Schema.is(FinishReason)(value) ? value : "unknown" +} + +function providerMetadata(value: unknown): ProviderMetadata | undefined { + return Schema.is(ProviderMetadata)(value) ? value : undefined +} + +function usage(value: unknown): Usage | undefined { + if (!value || typeof value !== "object") return undefined + const item = value as { + inputTokens?: number + outputTokens?: number + totalTokens?: number + reasoningTokens?: number + cachedInputTokens?: number + inputTokenDetails?: { cacheReadTokens?: number; cacheWriteTokens?: number } + outputTokenDetails?: { reasoningTokens?: number } + } + const result = Object.fromEntries( + Object.entries({ + inputTokens: item.inputTokens, + outputTokens: item.outputTokens, + totalTokens: item.totalTokens, + reasoningTokens: item.outputTokenDetails?.reasoningTokens ?? item.reasoningTokens, + cacheReadInputTokens: item.inputTokenDetails?.cacheReadTokens ?? item.cachedInputTokens, + cacheWriteInputTokens: item.inputTokenDetails?.cacheWriteTokens, + }).filter((entry) => entry[1] !== undefined), + ) + return new Usage(result) +} + +export function toLLMEvents( + state: ReturnType, + event: AISDKEvent, +): Effect.Effect, unknown> { + switch (event.type) { + case "start": + return Effect.succeed([]) + + case "start-step": + return Effect.succeed([LLMEvent.stepStart({ index: state.step })]) + + case "finish-step": + return Effect.sync(() => [ + LLMEvent.stepFinish({ + index: state.step++, + reason: finishReason(event.finishReason), + usage: usage(event.usage), + providerMetadata: providerMetadata(event.providerMetadata), + }), + ]) + + case "finish": + return Effect.sync(() => { + state.toolNames = {} + return [ + LLMEvent.requestFinish({ + reason: finishReason(event.finishReason), + usage: usage(event.totalUsage), + }), + ] + }) + + case "text-start": + return Effect.sync(() => { + state.currentTextID = contentBlockID(event.id ?? `text-${state.text++}`) + return [ + LLMEvent.textStart({ + id: state.currentTextID, + providerMetadata: providerMetadata(event.providerMetadata), + }), + ] + }) + + case "text-delta": + return Effect.succeed([ + LLMEvent.textDelta({ + id: event.id ? contentBlockID(event.id) : (state.currentTextID ?? contentBlockID(`text-${state.text++}`)), + text: event.text, + }), + ]) + + case "text-end": + return Effect.succeed([ + LLMEvent.textEnd({ + id: event.id ? contentBlockID(event.id) : (state.currentTextID ?? contentBlockID(`text-${state.text++}`)), + providerMetadata: providerMetadata(event.providerMetadata), + }), + ]) + + case "reasoning-start": + return Effect.sync(() => { + state.currentReasoningID = contentBlockID(event.id) + return [ + LLMEvent.reasoningStart({ + id: state.currentReasoningID, + providerMetadata: providerMetadata(event.providerMetadata), + }), + ] + }) + + case "reasoning-delta": + return Effect.succeed([ + LLMEvent.reasoningDelta({ + id: event.id ? contentBlockID(event.id) : (state.currentReasoningID ?? contentBlockID(`reasoning-${state.reasoning++}`)), + text: event.text, + }), + ]) + + case "reasoning-end": + return Effect.sync(() => { + const id = contentBlockID(event.id) + state.currentReasoningID = undefined + return [ + LLMEvent.reasoningEnd({ + id, + providerMetadata: providerMetadata(event.providerMetadata), + }), + ] + }) + + case "tool-input-start": + return Effect.sync(() => { + state.toolNames[event.id] = event.toolName + return [ + LLMEvent.toolInputStart({ + id: toolCallID(event.id), + name: event.toolName, + providerMetadata: providerMetadata(event.providerMetadata), + }), + ] + }) + + case "tool-input-delta": + return Effect.succeed([ + LLMEvent.toolInputDelta({ + id: toolCallID(event.id), + name: state.toolNames[event.id] ?? "unknown", + text: event.delta ?? "", + }), + ]) + + case "tool-input-end": + return Effect.succeed([ + LLMEvent.toolInputEnd({ + id: toolCallID(event.id), + name: state.toolNames[event.id] ?? "unknown", + }), + ]) + + case "tool-call": + return Effect.sync(() => { + state.toolNames[event.toolCallId] = event.toolName + return [ + LLMEvent.toolCall({ + id: toolCallID(event.toolCallId), + name: event.toolName, + input: event.input, + providerExecuted: "providerExecuted" in event ? event.providerExecuted : undefined, + providerMetadata: providerMetadata(event.providerMetadata), + }), + ] + }) + + case "tool-result": + return Effect.sync(() => { + const name = state.toolNames[event.toolCallId] ?? "unknown" + delete state.toolNames[event.toolCallId] + return [ + LLMEvent.toolResult({ + id: toolCallID(event.toolCallId), + name, + result: ToolResultValue.make(event.output), + providerExecuted: "providerExecuted" in event ? event.providerExecuted : undefined, + }), + ] + }) + + case "tool-error": + return Effect.sync(() => { + const name = state.toolNames[event.toolCallId] ?? "unknown" + delete state.toolNames[event.toolCallId] + return [ + LLMEvent.toolError({ + id: toolCallID(event.toolCallId), + name, + message: errorMessage(event.error), + }), + ] + }) + + case "error": + return Effect.fail(event.error) + + default: + return Effect.succeed([]) + } +} + +export * as LLMAISDK from "./llm-ai-sdk" diff --git a/packages/opencode/src/session/llm.ts b/packages/opencode/src/session/llm.ts index c7990d1b35..28e05aba8f 100644 --- a/packages/opencode/src/session/llm.ts +++ b/packages/opencode/src/session/llm.ts @@ -3,6 +3,7 @@ import * as Log from "@opencode-ai/core/util/log" import { Context, Effect, Layer, Record } from "effect" import * as Stream from "effect/Stream" import { streamText, wrapLanguageModel, type ModelMessage, type Tool, tool, jsonSchema } from "ai" +import type { LLMEvent } from "@opencode-ai/llm" import { mergeDeep } from "remeda" import { GitLabWorkflowLanguageModel } from "gitlab-ai-provider" import { ProviderTransform } from "@/provider/transform" @@ -24,10 +25,10 @@ import { InstallationVersion } from "@opencode-ai/core/installation/version" import { EffectBridge } from "@/effect/bridge" import * as Option from "effect/Option" import * as OtelTracer from "@effect/opentelemetry/Tracer" +import { LLMAISDK } from "./llm-ai-sdk" const log = Log.create({ service: "llm" }) export const OUTPUT_TOKEN_MAX = ProviderTransform.OUTPUT_TOKEN_MAX -type Result = Awaited> // Avoid re-instantiating remeda's deep merge types in this hot LLM path; the runtime behavior is still mergeDeep. const mergeOptions = (target: Record, source: Record | undefined): Record => @@ -52,10 +53,8 @@ export type StreamRequest = StreamInput & { abort: AbortSignal } -export type Event = Result["fullStream"] extends AsyncIterable ? T : never - export interface Interface { - readonly stream: (input: StreamInput) => Stream.Stream + readonly stream: (input: StreamInput) => Stream.Stream } export class Service extends Context.Service()("@opencode/LLM") {} @@ -427,7 +426,11 @@ const live: Layer.Layer< const result = yield* run({ ...input, abort: ctrl.signal }) - return Stream.fromAsyncIterable(result.fullStream, (e) => (e instanceof Error ? e : new Error(String(e)))) + const state = LLMAISDK.adapterState() + return Stream.fromAsyncIterable(result.fullStream, (e) => (e instanceof Error ? e : new Error(String(e)))).pipe( + Stream.mapEffect((event) => LLMAISDK.toLLMEvents(state, event)), + Stream.flatMap((events) => Stream.fromIterable(events)), + ) }), ), ) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 579c4cc42c..2d78dc70e8 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -1,4 +1,4 @@ -import { Cause, Deferred, Effect, Exit, Layer, Context, Scope } from "effect" +import { Cause, Deferred, Effect, Exit, Layer, Context, Scope, Schema } from "effect" import * as Stream from "effect/Stream" import { Agent } from "@/agent/agent" import { Bus } from "@/bus" @@ -26,14 +26,13 @@ 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" +import { Usage, type LLMEvent } from "@opencode-ai/llm" 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: ( @@ -67,6 +66,7 @@ type ToolCall = { messageID: MessageV2.ToolPart["messageID"] sessionID: MessageV2.ToolPart["sessionID"] done: Deferred.Deferred + inputEnded: boolean } interface ProcessorContext extends Input { @@ -79,7 +79,7 @@ interface ProcessorContext extends Input { reasoningMap: Record } -type StreamEvent = Event +type StreamEvent = LLMEvent export class Service extends Context.Service()("@opencode/SessionProcessor") {} @@ -223,9 +223,85 @@ export const layer: Layer.Layer< return true }) + const finishReasoning = Effect.fn("SessionProcessor.finishReasoning")(function* (reasoningID: string) { + if (!(reasoningID 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, + text: ctx.reasoningMap[reasoningID].text, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } + // oxlint-disable-next-line no-self-assign -- reactivity trigger + ctx.reasoningMap[reasoningID].text = ctx.reasoningMap[reasoningID].text + ctx.reasoningMap[reasoningID].time = { ...ctx.reasoningMap[reasoningID].time, end: Date.now() } + yield* session.updatePart(ctx.reasoningMap[reasoningID]) + delete ctx.reasoningMap[reasoningID] + }) + + const ensureToolCall = Effect.fn("SessionProcessor.ensureToolCall")(function* (input: { + id: string + name: string + providerExecuted?: boolean + }) { + const existing = yield* readToolCall(input.id) + if (existing) return existing + // 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: input.id, + name: input.name, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } + const part = yield* session.updatePart({ + id: PartID.ascending(), + messageID: ctx.assistantMessage.id, + sessionID: ctx.assistantMessage.sessionID, + type: "tool", + tool: input.name, + callID: input.id, + state: { status: "pending", input: {}, raw: "" }, + metadata: input.providerExecuted ? { providerExecuted: true } : undefined, + } satisfies MessageV2.ToolPart) + ctx.toolcalls[input.id] = { + done: yield* Deferred.make(), + partID: part.id, + messageID: part.messageID, + sessionID: part.sessionID, + inputEnded: false, + } + return { call: ctx.toolcalls[input.id], part } + }) + + const isFilePart = Schema.is(MessageV2.FilePart) + + const toolResultOutput = (value: Extract) => { + if (isRecord(value.result.value) && typeof value.result.value.output === "string") { + return { + title: typeof value.result.value.title === "string" ? value.result.value.title : value.name, + metadata: isRecord(value.result.value.metadata) ? value.result.value.metadata : {}, + output: value.result.value.output, + attachments: Array.isArray(value.result.value.attachments) + ? value.result.value.attachments.filter(isFilePart) + : undefined, + } + } + return { + title: value.name, + metadata: value.result.type === "json" && isRecord(value.result.value) ? value.result.value : {}, + output: typeof value.result.value === "string" ? value.result.value : (JSON.stringify(value.result.value) ?? ""), + } + } + + const toolInput = (value: unknown): Record => (isRecord(value) ? value : { value }) + const handleEvent = Effect.fnUntraced(function* (value: StreamEvent) { switch (value.type) { - case "start": + case "request-start": yield* status.set(ctx.sessionID, { type: "busy" }) return @@ -251,116 +327,132 @@ export const layer: Layer.Layer< 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 + case "reasoning-delta": { + const reasoningID = value.id ?? "reasoning" + if (!(reasoningID in ctx.reasoningMap)) { + // 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, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } + ctx.reasoningMap[reasoningID] = { + id: PartID.ascending(), + messageID: ctx.assistantMessage.id, + sessionID: ctx.assistantMessage.sessionID, + type: "reasoning", + text: "", + time: { start: Date.now() }, + } + yield* session.updatePart(ctx.reasoningMap[reasoningID]) + } + ctx.reasoningMap[reasoningID].text += value.text yield* session.updatePartDelta({ - sessionID: ctx.reasoningMap[value.id].sessionID, - messageID: ctx.reasoningMap[value.id].messageID, - partID: ctx.reasoningMap[value.id].id, + sessionID: ctx.reasoningMap[reasoningID].sessionID, + messageID: ctx.reasoningMap[reasoningID].messageID, + partID: ctx.reasoningMap[reasoningID].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()), - }) + if (value.providerMetadata && value.id in ctx.reasoningMap) { + ctx.reasoningMap[value.id].metadata = value.providerMetadata } - // 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] + yield* finishReasoning(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, + throw new Error(`Tool call not allowed while generating summary: ${value.name}`) } + yield* ensureToolCall(value) return - case "tool-input-delta": + case "tool-input-delta": { + if (ctx.assistantMessage.summary) { + throw new Error(`Tool call not allowed while generating summary: ${value.name}`) + } + yield* ensureToolCall(value) + if (value.text) { + yield* updateToolCall(value.id, (match) => ({ + ...match, + state: + match.state.status === "pending" + ? { ...match.state, raw: match.state.raw + value.text } + : match.state, + })) + } return + } case "tool-input-end": { + const toolCall = yield* ensureToolCall(value) // 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: "", + text: toolCall.part.state.status === "pending" ? toolCall.part.state.raw : "", timestamp: DateTime.makeUnsafe(Date.now()), }) } + ctx.toolcalls[value.id] = { ...toolCall.call, inputEnded: true } return } case "tool-call": { if (ctx.assistantMessage.summary) { - throw new Error(`Tool call not allowed while generating summary: ${value.toolName}`) + throw new Error(`Tool call not allowed while generating summary: ${value.name}`) + } + const toolCall = yield* ensureToolCall(value) + const input = toolInput(value.input) + const raw = toolCall.part.state.status === "pending" ? toolCall.part.state.raw : "" + if (!toolCall.call.inputEnded) { + // 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: raw, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } } - 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, + callID: value.id, + tool: value.name, + input, provider: { - executed: toolCall?.part.metadata?.providerExecuted === true, + executed: toolCall.part.metadata?.providerExecuted === true, ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), }, timestamp: DateTime.makeUnsafe(Date.now()), }) } - yield* updateToolCall(value.toolCallId, (match) => ({ + yield* updateToolCall(value.id, (match) => ({ ...match, - tool: value.toolName, - state: { - ...match.state, - status: "running", - input: value.input, - time: { start: Date.now() }, + tool: value.name, + state: + match.state.status === "running" + ? { ...match.state, input } + : { + status: "running", + input, + time: { start: Date.now() }, + }, + metadata: { + ...match.metadata, + ...value.providerMetadata, + ...(match.metadata?.providerExecuted ? { providerExecuted: true } : {}), }, - metadata: match.metadata?.providerExecuted - ? { ...value.providerMetadata, providerExecuted: true } - : value.providerMetadata, })) const parts = MessageV2.parts(ctx.assistantMessage.id) @@ -371,9 +463,9 @@ export const layer: Layer.Layer< !recentParts.every( (part) => part.type === "tool" && - part.tool === value.toolName && + part.tool === value.name && part.state.status !== "pending" && - JSON.stringify(part.state.input) === JSON.stringify(value.input), + JSON.stringify(part.state.input) === JSON.stringify(input), ) ) { return @@ -382,27 +474,19 @@ export const layer: Layer.Layer< const agent = yield* agents.get(ctx.assistantMessage.agent) yield* permission.ask({ permission: "doom_loop", - patterns: [value.toolName], + patterns: [value.name], sessionID: ctx.assistantMessage.sessionID, - metadata: { tool: value.toolName, input: value.input }, - always: [value.toolName], + metadata: { tool: value.name, input }, + always: [value.name], 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) => + const toolCall = yield* readToolCall(value.id) + const rawOutput = toolResultOutput(value) + const normalized = yield* Effect.forEach(rawOutput.attachments ?? [], (attachment) => attachment.mime.startsWith("image/") ? image.normalize(attachment).pipe(Effect.exit) : Effect.succeed(Exit.succeed(attachment)), @@ -410,18 +494,18 @@ export const layer: Layer.Layer< const omitted = normalized.filter(Exit.isFailure).length const attachments = normalized.filter(Exit.isSuccess).map((item) => item.value) const output = { - ...value.output, + ...rawOutput, 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, + ? rawOutput.output + : `${rawOutput.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, + callID: value.id, structured: output.metadata, content: [ { @@ -429,32 +513,32 @@ export const layer: Layer.Layer< text: output.output, }, ...(output.attachments?.map((item: MessageV2.FilePart) => ({ - type: "file", + type: "file" as const, uri: item.url, mime: item.mime, name: item.filename, })) ?? []), ], provider: { - executed: toolCall?.part.metadata?.providerExecuted === true, + executed: value.providerExecuted === true || toolCall?.part.metadata?.providerExecuted === true, }, timestamp: DateTime.makeUnsafe(Date.now()), }) } - yield* completeToolCall(value.toolCallId, output) + yield* completeToolCall(value.id, output) return } case "tool-error": { - const toolCall = yield* readToolCall(value.toolCallId) + const toolCall = yield* readToolCall(value.id) // 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, + callID: value.id, error: { type: "unknown", - message: errorMessage(value.error), + message: value.message, }, provider: { executed: toolCall?.part.metadata?.providerExecuted === true, @@ -462,14 +546,14 @@ export const layer: Layer.Layer< timestamp: DateTime.makeUnsafe(Date.now()), }) } - yield* failToolCall(value.toolCallId, value.error) + yield* failToolCall(value.id, new Error(value.message)) return } - case "error": - throw value.error + case "provider-error": + throw new Error(value.message) - case "start-step": + case "step-start": if (!ctx.snapshot) ctx.snapshot = yield* snapshot.track() if (!ctx.assistantMessage.summary) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. @@ -496,19 +580,20 @@ export const layer: Layer.Layer< }) return - case "finish-step": { + case "step-finish": { const completedSnapshot = yield* snapshot.track() - const usage = Session.getUsage({ - model: ctx.model, - usage: value.usage, - metadata: value.providerMetadata, - }) + yield* Effect.forEach(Object.keys(ctx.reasoningMap), finishReasoning) + const usage = Session.getUsage({ + model: ctx.model, + usage: value.usage ?? new 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, + finish: value.reason, cost: usage.cost, tokens: usage.tokens, snapshot: completedSnapshot, @@ -516,12 +601,12 @@ export const layer: Layer.Layer< }) } } - ctx.assistantMessage.finish = value.finishReason + ctx.assistantMessage.finish = value.reason ctx.assistantMessage.cost += usage.cost ctx.assistantMessage.tokens = usage.tokens yield* session.updatePart({ id: PartID.ascending(), - reason: value.finishReason, + reason: value.reason, snapshot: completedSnapshot, messageID: ctx.assistantMessage.id, sessionID: ctx.assistantMessage.sessionID, @@ -584,7 +669,6 @@ export const layer: Layer.Layer< 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, @@ -626,12 +710,9 @@ export const layer: Layer.Layer< ctx.currentText = undefined return - case "finish": + case "request-finish": return - default: - slog.info("unhandled", { event: value.type, value }) - return } }) @@ -733,6 +814,7 @@ export const layer: Layer.Layer< yield* Effect.gen(function* () { ctx.currentText = undefined ctx.reasoningMap = {} + yield* status.set(ctx.sessionID, { type: "busy" }) const stream = llm.stream(streamInput) yield* stream.pipe( @@ -816,12 +898,12 @@ export const defaultLayer = Layer.suspend(() => Layer.provide(LLM.defaultLayer), Layer.provide(Permission.defaultLayer), Layer.provide(Plugin.defaultLayer), + Layer.provide(Image.defaultLayer), Layer.provide(SessionSummary.defaultLayer), Layer.provide(SessionStatus.defaultLayer), - Layer.provide(Image.defaultLayer), + Layer.provide(SyncEvent.defaultLayer), Layer.provide(Bus.layer), Layer.provide(Config.defaultLayer), - Layer.provide(SyncEvent.defaultLayer), ), ) diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index 15246dac39..6fd0c5049c 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -60,6 +60,7 @@ import * as DateTime from "effect/DateTime" import { eq } from "@/storage/db" import * as Database from "@/storage/db" import { SessionTable } from "./session.sql" +import { LLMEvent } from "@opencode-ai/llm" // @ts-ignore globalThis.AI_SDK_LOG_WARNINGS = false @@ -359,7 +360,7 @@ export const layer = Layer.effect( messages: [{ role: "user", content: "Generate a title for this conversation:\n" }, ...msgs], }) .pipe( - Stream.filter((e): e is Extract => e.type === "text-delta"), + Stream.filter(LLMEvent.is.textDelta), Stream.map((e) => e.text), Stream.mkString, Effect.orDie, diff --git a/packages/opencode/src/session/session.ts b/packages/opencode/src/session/session.ts index eff027579a..4bcb4e5caa 100644 --- a/packages/opencode/src/session/session.ts +++ b/packages/opencode/src/session/session.ts @@ -3,7 +3,7 @@ import path from "path" import { BusEvent } from "@/bus/bus-event" import { Bus } from "@/bus" import { Decimal } from "decimal.js" -import { type ProviderMetadata, type LanguageModelUsage } from "ai" +import type { ProviderMetadata, Usage } from "@opencode-ai/llm" import { Flag } from "@opencode-ai/core/flag/flag" import { InstallationVersion } from "@opencode-ai/core/installation/version" @@ -373,21 +373,19 @@ export function plan(input: { slug: string; time: { created: number } }, instanc return path.join(base, [input.time.created, input.slug].join("-") + ".md") } -export const getUsage = (input: { model: Provider.Model; usage: LanguageModelUsage; metadata?: ProviderMetadata }) => { +export const getUsage = (input: { model: Provider.Model; usage: Usage; metadata?: ProviderMetadata }) => { const safe = (value: number) => { if (!Number.isFinite(value)) return 0 return Math.max(0, value) } const inputTokens = safe(input.usage.inputTokens ?? 0) const outputTokens = safe(input.usage.outputTokens ?? 0) - const reasoningTokens = safe(input.usage.outputTokenDetails?.reasoningTokens ?? input.usage.reasoningTokens ?? 0) + const reasoningTokens = safe(input.usage.reasoningTokens ?? 0) - const cacheReadInputTokens = safe( - input.usage.inputTokenDetails?.cacheReadTokens ?? input.usage.cachedInputTokens ?? 0, - ) + const cacheReadInputTokens = safe(input.usage.cacheReadInputTokens ?? 0) const cacheWriteInputTokens = safe( Number( - input.usage.inputTokenDetails?.cacheWriteTokens ?? + input.usage.cacheWriteInputTokens ?? input.metadata?.["anthropic"]?.["cacheCreationInputTokens"] ?? // google-vertex-anthropic returns metadata under "vertex" key // (AnthropicMessagesLanguageModel custom provider key from 'vertex.anthropic.messages') diff --git a/packages/opencode/test/session/compaction.test.ts b/packages/opencode/test/session/compaction.test.ts index c7f349d5ce..fa77b22a68 100644 --- a/packages/opencode/test/session/compaction.test.ts +++ b/packages/opencode/test/session/compaction.test.ts @@ -28,6 +28,7 @@ import { testEffect } from "../lib/effect" import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" import { TestConfig } from "../fixture/config" import { SyncEvent } from "@/sync" +import { LLMEvent, Usage } from "@opencode-ai/llm" void Log.init({ print: false }) @@ -45,6 +46,10 @@ const ref = { modelID: ModelID.make("test-model"), } +const usage = (input: ConstructorParameters[0]) => new Usage(input) + +const basicUsage = () => usage({ inputTokens: 1, outputTokens: 1, totalTokens: 2 }) + afterEach(() => { mock.restore() }) @@ -289,11 +294,11 @@ function readCompactionPart(sessionID: SessionID) { function llm() { const queue: Array< - Stream.Stream | ((input: LLM.StreamInput) => Stream.Stream) + Stream.Stream | ((input: LLM.StreamInput) => Stream.Stream) > = [] return { - push(stream: Stream.Stream | ((input: LLM.StreamInput) => Stream.Stream)) { + push(stream: Stream.Stream | ((input: LLM.StreamInput) => Stream.Stream)) { queue.push(stream) }, layer: Layer.succeed( @@ -312,54 +317,22 @@ function llm() { function reply( text: string, capture?: (input: LLM.StreamInput) => void, -): (input: LLM.StreamInput) => Stream.Stream { +): (input: LLM.StreamInput) => Stream.Stream { return (input) => { capture?.(input) return Stream.make( - { type: "start" } satisfies LLM.Event, - { type: "text-start", id: "txt-0" } satisfies LLM.Event, - { type: "text-delta", id: "txt-0", delta: text, text } as LLM.Event, - { type: "text-end", id: "txt-0" } satisfies LLM.Event, - { - type: "finish-step", - finishReason: "stop", - rawFinishReason: "stop", - response: { id: "res", modelId: "test-model", timestamp: new Date() }, - providerMetadata: undefined, - usage: { - inputTokens: 1, - outputTokens: 1, - totalTokens: 2, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, - } satisfies LLM.Event, - { - type: "finish", - finishReason: "stop", - rawFinishReason: "stop", - totalUsage: { - inputTokens: 1, - outputTokens: 1, - totalTokens: 2, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, - } satisfies LLM.Event, + LLMEvent.textStart({ id: "txt-0" }), + LLMEvent.textDelta({ id: "txt-0", text }), + LLMEvent.textEnd({ id: "txt-0" }), + LLMEvent.stepFinish({ + index: 0, + reason: "stop", + usage: basicUsage(), + }), + LLMEvent.requestFinish({ + reason: "stop", + usage: basicUsage(), + }), ) } } @@ -1198,7 +1171,7 @@ describe("session.compaction.process", () => { Stream.fromAsyncIterable( { async *[Symbol.asyncIterator]() { - yield { type: "start" } as LLM.Event + yield LLMEvent.stepStart({ index: 0 }) throw new APICallError({ message: "boom", url: "https://example.com/v1/chat/completions", @@ -1290,49 +1263,16 @@ describe("session.compaction.process", () => { const stub = llm() stub.push( Stream.make( - { type: "start" } satisfies LLM.Event, - { type: "tool-input-start", id: "call-1", toolName: "_noop" } satisfies LLM.Event, - { type: "tool-call", toolCallId: "call-1", toolName: "_noop", input: {} } satisfies LLM.Event, - { - type: "finish-step", - finishReason: "tool-calls", - rawFinishReason: "tool_calls", - response: { id: "res", modelId: "test-model", timestamp: new Date() }, - providerMetadata: undefined, - usage: { - inputTokens: 1, - outputTokens: 1, - totalTokens: 2, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, - } satisfies LLM.Event, - { - type: "finish", - finishReason: "tool-calls", - rawFinishReason: "tool_calls", - totalUsage: { - inputTokens: 1, - outputTokens: 1, - totalTokens: 2, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, - } satisfies LLM.Event, + LLMEvent.toolCall({ id: "call-1", name: "_noop", input: {} }), + LLMEvent.stepFinish({ + index: 0, + reason: "tool-calls", + usage: basicUsage(), + }), + LLMEvent.requestFinish({ + reason: "tool-calls", + usage: basicUsage(), + }), ), ) return Effect.gen(function* () { @@ -1543,20 +1483,7 @@ describe("SessionNs.getUsage", () => { const model = createModel({ context: 100_000, output: 32_000 }) const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 1000, - outputTokens: 500, - totalTokens: 1500, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, + usage: usage({ inputTokens: 1000, outputTokens: 500, totalTokens: 1500 }), }) expect(result.tokens.input).toBe(1000) @@ -1570,20 +1497,7 @@ describe("SessionNs.getUsage", () => { const model = createModel({ context: 100_000, output: 32_000 }) const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 1000, - outputTokens: 500, - totalTokens: 1500, - inputTokenDetails: { - noCacheTokens: 800, - cacheReadTokens: 200, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, + usage: usage({ inputTokens: 1000, outputTokens: 500, totalTokens: 1500, cacheReadInputTokens: 200 }), }) expect(result.tokens.input).toBe(800) @@ -1594,20 +1508,7 @@ describe("SessionNs.getUsage", () => { const model = createModel({ context: 100_000, output: 32_000 }) const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 1000, - outputTokens: 500, - totalTokens: 1500, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, + usage: usage({ inputTokens: 1000, outputTokens: 500, totalTokens: 1500 }), metadata: { anthropic: { cacheCreationInputTokens: 300, @@ -1623,20 +1524,7 @@ describe("SessionNs.getUsage", () => { // AI SDK v6 normalizes inputTokens to include cached tokens for all providers const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 1000, - outputTokens: 500, - totalTokens: 1500, - inputTokenDetails: { - noCacheTokens: 800, - cacheReadTokens: 200, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, + usage: usage({ inputTokens: 1000, outputTokens: 500, totalTokens: 1500, cacheReadInputTokens: 200 }), metadata: { anthropic: {}, }, @@ -1650,20 +1538,7 @@ describe("SessionNs.getUsage", () => { const model = createModel({ context: 100_000, output: 32_000 }) const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 1000, - outputTokens: 500, - totalTokens: 1500, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: 400, - reasoningTokens: 100, - }, - }, + usage: usage({ inputTokens: 1000, outputTokens: 500, reasoningTokens: 100, totalTokens: 1500 }), }) expect(result.tokens.input).toBe(1000) @@ -1684,20 +1559,7 @@ describe("SessionNs.getUsage", () => { }) const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 0, - outputTokens: 1_000_000, - totalTokens: 1_000_000, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: 750_000, - reasoningTokens: 250_000, - }, - }, + usage: usage({ inputTokens: 0, outputTokens: 1_000_000, reasoningTokens: 250_000, totalTokens: 1_000_000 }), }) expect(result.tokens.output).toBe(750_000) @@ -1709,20 +1571,7 @@ describe("SessionNs.getUsage", () => { const model = createModel({ context: 100_000, output: 32_000 }) const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 0, - outputTokens: 0, - totalTokens: 0, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, + usage: usage({ inputTokens: 0, outputTokens: 0, totalTokens: 0 }), }) expect(result.tokens.input).toBe(0) @@ -1745,20 +1594,7 @@ describe("SessionNs.getUsage", () => { }) const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 1_000_000, - outputTokens: 100_000, - totalTokens: 1_100_000, - inputTokenDetails: { - noCacheTokens: undefined, - cacheReadTokens: undefined, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, + usage: usage({ inputTokens: 1_000_000, outputTokens: 100_000, totalTokens: 1_100_000 }), }) expect(result.cost).toBe(3 + 1.5) @@ -1769,24 +1605,16 @@ describe("SessionNs.getUsage", () => { (npm) => { const model = createModel({ context: 100_000, output: 32_000, npm }) // AI SDK v6: inputTokens includes cached tokens for all providers - const usage = { + const item = usage({ inputTokens: 1000, outputTokens: 500, totalTokens: 1500, - inputTokenDetails: { - noCacheTokens: 800, - cacheReadTokens: 200, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - } + cacheReadInputTokens: 200, + }) if (npm === "@ai-sdk/amazon-bedrock") { const result = SessionNs.getUsage({ model, - usage, + usage: item, metadata: { bedrock: { usage: { @@ -1807,7 +1635,7 @@ describe("SessionNs.getUsage", () => { const result = SessionNs.getUsage({ model, - usage, + usage: item, metadata: { anthropic: { cacheCreationInputTokens: 300, @@ -1828,20 +1656,7 @@ describe("SessionNs.getUsage", () => { const model = createModel({ context: 100_000, output: 32_000, npm: "@ai-sdk/google-vertex/anthropic" }) const result = SessionNs.getUsage({ model, - usage: { - inputTokens: 1000, - outputTokens: 500, - totalTokens: 1500, - inputTokenDetails: { - noCacheTokens: 800, - cacheReadTokens: 200, - cacheWriteTokens: undefined, - }, - outputTokenDetails: { - textTokens: undefined, - reasoningTokens: undefined, - }, - }, + usage: usage({ inputTokens: 1000, outputTokens: 500, totalTokens: 1500, cacheReadInputTokens: 200 }), metadata: { vertex: { cacheCreationInputTokens: 300, From 648c5cd1b6f1756c751c1cf28a904b55f1bd8d17 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sun, 10 May 2026 12:47:50 -0400 Subject: [PATCH 02/13] fix(session): align OpenAI response stream fixtures --- packages/llm/src/schema/events.ts | 9 ++++++-- packages/opencode/test/session/llm.test.ts | 24 ++++++++++++++++++++++ 2 files changed, 31 insertions(+), 2 deletions(-) diff --git a/packages/llm/src/schema/events.ts b/packages/llm/src/schema/events.ts index 6e6bb1541b..b2109f3693 100644 --- a/packages/llm/src/schema/events.ts +++ b/packages/llm/src/schema/events.ts @@ -222,6 +222,9 @@ const llmEventTagged = Schema.Union([ ]).pipe(Schema.toTaggedUnion("type")) type WithID = Omit & { readonly id: ID | string } +type WithUsage = Omit & { + readonly usage?: Usage | ConstructorParameters[0] +} const responseID = (value: ResponseID | string) => ResponseID.make(value) const contentBlockID = (value: ContentBlockID | string) => ContentBlockID.make(value) @@ -252,8 +255,10 @@ export const LLMEvent = Object.assign(llmEventTagged, { toolCall: (input: WithID) => ToolCall.make({ ...input, id: toolCallID(input.id) }), toolResult: (input: WithID) => ToolResult.make({ ...input, id: toolCallID(input.id) }), toolError: (input: WithID) => ToolError.make({ ...input, id: toolCallID(input.id) }), - stepFinish: StepFinish.make, - requestFinish: RequestFinish.make, + stepFinish: (input: WithUsage) => + StepFinish.make({ ...input, usage: input.usage instanceof Usage ? input.usage : new Usage(input.usage ?? {}) }), + requestFinish: (input: WithUsage) => + RequestFinish.make({ ...input, usage: input.usage instanceof Usage ? input.usage : new Usage(input.usage ?? {}) }), providerError: ProviderErrorEvent.make, is: { requestStart: llmEventTagged.guards["request-start"], diff --git a/packages/opencode/test/session/llm.test.ts b/packages/opencode/test/session/llm.test.ts index 2879d04812..f078e352ac 100644 --- a/packages/opencode/test/session/llm.test.ts +++ b/packages/opencode/test/session/llm.test.ts @@ -581,6 +581,18 @@ describe("session.llm.stream", () => { service_tier: null, }, }, + { + type: "response.output_item.added", + output_index: 0, + item: { type: "message", id: "item-1", status: "in_progress", role: "assistant", content: [] }, + }, + { + type: "response.content_part.added", + item_id: "item-1", + output_index: 0, + content_index: 0, + part: { type: "output_text", text: "", annotations: [] }, + }, { type: "response.output_text.delta", item_id: "item-1", @@ -694,6 +706,18 @@ describe("session.llm.stream", () => { service_tier: null, }, }, + { + type: "response.output_item.added", + output_index: 0, + item: { type: "message", id: "item-data-url", status: "in_progress", role: "assistant", content: [] }, + }, + { + type: "response.content_part.added", + item_id: "item-data-url", + output_index: 0, + content_index: 0, + part: { type: "output_text", text: "", annotations: [] }, + }, { type: "response.output_text.delta", item_id: "item-data-url", From 4481d46071e21acc2a7c4a5a9d7f2893fae4322f Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sun, 10 May 2026 21:50:04 -0400 Subject: [PATCH 03/13] refactor(session): simplify LLM event adapter --- packages/llm/src/schema/events.ts | 10 ++++- packages/opencode/src/session/llm-ai-sdk.ts | 47 ++++++++++++--------- 2 files changed, 35 insertions(+), 22 deletions(-) diff --git a/packages/llm/src/schema/events.ts b/packages/llm/src/schema/events.ts index b2109f3693..212a01f6a3 100644 --- a/packages/llm/src/schema/events.ts +++ b/packages/llm/src/schema/events.ts @@ -256,9 +256,15 @@ export const LLMEvent = Object.assign(llmEventTagged, { toolResult: (input: WithID) => ToolResult.make({ ...input, id: toolCallID(input.id) }), toolError: (input: WithID) => ToolError.make({ ...input, id: toolCallID(input.id) }), stepFinish: (input: WithUsage) => - StepFinish.make({ ...input, usage: input.usage instanceof Usage ? input.usage : new Usage(input.usage ?? {}) }), + StepFinish.make({ + ...input, + usage: input.usage === undefined ? undefined : input.usage instanceof Usage ? input.usage : new Usage(input.usage), + }), requestFinish: (input: WithUsage) => - RequestFinish.make({ ...input, usage: input.usage instanceof Usage ? input.usage : new Usage(input.usage ?? {}) }), + RequestFinish.make({ + ...input, + usage: input.usage === undefined ? undefined : input.usage instanceof Usage ? input.usage : new Usage(input.usage), + }), providerError: ProviderErrorEvent.make, is: { requestStart: llmEventTagged.guards["request-start"], diff --git a/packages/opencode/src/session/llm-ai-sdk.ts b/packages/opencode/src/session/llm-ai-sdk.ts index 9d5407f975..bba31861e1 100644 --- a/packages/opencode/src/session/llm-ai-sdk.ts +++ b/packages/opencode/src/session/llm-ai-sdk.ts @@ -1,4 +1,4 @@ -import { ContentBlockID, FinishReason, LLMEvent, ProviderMetadata, ToolCallID, ToolResultValue, Usage } from "@opencode-ai/llm" +import { FinishReason, LLMEvent, ProviderMetadata, ToolResultValue } from "@opencode-ai/llm" import { Effect, Schema } from "effect" import { type streamText } from "ai" import { errorMessage } from "@/util/error" @@ -11,15 +11,12 @@ export function adapterState() { step: 0, text: 0, reasoning: 0, - currentTextID: undefined as ContentBlockID | undefined, - currentReasoningID: undefined as ContentBlockID | undefined, + currentTextID: undefined as string | undefined, + currentReasoningID: undefined as string | undefined, toolNames: {} as Record, } } -const contentBlockID = (value: string) => ContentBlockID.make(value) -const toolCallID = (value: string) => ToolCallID.make(value) - function finishReason(value: string | undefined): FinishReason { return Schema.is(FinishReason)(value) ? value : "unknown" } @@ -28,7 +25,7 @@ function providerMetadata(value: unknown): ProviderMetadata | undefined { return Schema.is(ProviderMetadata)(value) ? value : undefined } -function usage(value: unknown): Usage | undefined { +function usage(value: unknown) { if (!value || typeof value !== "object") return undefined const item = value as { inputTokens?: number @@ -49,7 +46,17 @@ function usage(value: unknown): Usage | undefined { cacheWriteInputTokens: item.inputTokenDetails?.cacheWriteTokens, }).filter((entry) => entry[1] !== undefined), ) - return new Usage(result) + return result +} + +function currentTextID(state: ReturnType, id: string | undefined) { + state.currentTextID = id ?? state.currentTextID ?? `text-${state.text++}` + return state.currentTextID +} + +function currentReasoningID(state: ReturnType, id: string | undefined) { + state.currentReasoningID = id ?? state.currentReasoningID ?? `reasoning-${state.reasoning++}` + return state.currentReasoningID } export function toLLMEvents( @@ -86,7 +93,7 @@ export function toLLMEvents( case "text-start": return Effect.sync(() => { - state.currentTextID = contentBlockID(event.id ?? `text-${state.text++}`) + state.currentTextID = currentTextID(state, event.id) return [ LLMEvent.textStart({ id: state.currentTextID, @@ -98,7 +105,7 @@ export function toLLMEvents( case "text-delta": return Effect.succeed([ LLMEvent.textDelta({ - id: event.id ? contentBlockID(event.id) : (state.currentTextID ?? contentBlockID(`text-${state.text++}`)), + id: currentTextID(state, event.id), text: event.text, }), ]) @@ -106,14 +113,14 @@ export function toLLMEvents( case "text-end": return Effect.succeed([ LLMEvent.textEnd({ - id: event.id ? contentBlockID(event.id) : (state.currentTextID ?? contentBlockID(`text-${state.text++}`)), + id: currentTextID(state, event.id), providerMetadata: providerMetadata(event.providerMetadata), }), ]) case "reasoning-start": return Effect.sync(() => { - state.currentReasoningID = contentBlockID(event.id) + state.currentReasoningID = currentReasoningID(state, event.id) return [ LLMEvent.reasoningStart({ id: state.currentReasoningID, @@ -125,14 +132,14 @@ export function toLLMEvents( case "reasoning-delta": return Effect.succeed([ LLMEvent.reasoningDelta({ - id: event.id ? contentBlockID(event.id) : (state.currentReasoningID ?? contentBlockID(`reasoning-${state.reasoning++}`)), + id: currentReasoningID(state, event.id), text: event.text, }), ]) case "reasoning-end": return Effect.sync(() => { - const id = contentBlockID(event.id) + const id = currentReasoningID(state, event.id) state.currentReasoningID = undefined return [ LLMEvent.reasoningEnd({ @@ -147,7 +154,7 @@ export function toLLMEvents( state.toolNames[event.id] = event.toolName return [ LLMEvent.toolInputStart({ - id: toolCallID(event.id), + id: event.id, name: event.toolName, providerMetadata: providerMetadata(event.providerMetadata), }), @@ -157,7 +164,7 @@ export function toLLMEvents( case "tool-input-delta": return Effect.succeed([ LLMEvent.toolInputDelta({ - id: toolCallID(event.id), + id: event.id, name: state.toolNames[event.id] ?? "unknown", text: event.delta ?? "", }), @@ -166,7 +173,7 @@ export function toLLMEvents( case "tool-input-end": return Effect.succeed([ LLMEvent.toolInputEnd({ - id: toolCallID(event.id), + id: event.id, name: state.toolNames[event.id] ?? "unknown", }), ]) @@ -176,7 +183,7 @@ export function toLLMEvents( state.toolNames[event.toolCallId] = event.toolName return [ LLMEvent.toolCall({ - id: toolCallID(event.toolCallId), + id: event.toolCallId, name: event.toolName, input: event.input, providerExecuted: "providerExecuted" in event ? event.providerExecuted : undefined, @@ -191,7 +198,7 @@ export function toLLMEvents( delete state.toolNames[event.toolCallId] return [ LLMEvent.toolResult({ - id: toolCallID(event.toolCallId), + id: event.toolCallId, name, result: ToolResultValue.make(event.output), providerExecuted: "providerExecuted" in event ? event.providerExecuted : undefined, @@ -205,7 +212,7 @@ export function toLLMEvents( delete state.toolNames[event.toolCallId] return [ LLMEvent.toolError({ - id: toolCallID(event.toolCallId), + id: event.toolCallId, name, message: errorMessage(event.error), }), From d79e8ba701e2dad7dc0aeeec10bbf76fe987313b Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 17:30:03 -0400 Subject: [PATCH 04/13] refactor(session): add native LLM request adapter --- packages/opencode/src/session/llm-native.ts | 184 +++++++++++++++ .../opencode/test/session/llm-native.test.ts | 219 ++++++++++++++++++ 2 files changed, 403 insertions(+) create mode 100644 packages/opencode/src/session/llm-native.ts create mode 100644 packages/opencode/test/session/llm-native.test.ts diff --git a/packages/opencode/src/session/llm-native.ts b/packages/opencode/src/session/llm-native.ts new file mode 100644 index 0000000000..6bb2159412 --- /dev/null +++ b/packages/opencode/src/session/llm-native.ts @@ -0,0 +1,184 @@ +import type { JsonSchema, LLMRequest, ProviderMetadata, ToolDefinition } from "@opencode-ai/llm" +import { LLM } from "@opencode-ai/llm" +import type { ModelMessage } from "ai" +import type { Provider } from "@/provider/provider" + +type ToolInput = { + readonly description?: string + readonly inputSchema?: unknown +} + +export type RequestInput = { + readonly model: Provider.Model + readonly system?: readonly string[] + readonly messages: readonly ModelMessage[] + readonly tools?: Record + readonly toolChoice?: "auto" | "required" | "none" + readonly temperature?: number + readonly topP?: number + readonly topK?: number + readonly maxOutputTokens?: number + readonly providerOptions?: LLMRequest["providerOptions"] + readonly headers?: Record +} + +const DEFAULT_BASE_URL: Record = { + "@ai-sdk/openai": "https://api.openai.com/v1", + "@ai-sdk/anthropic": "https://api.anthropic.com/v1", + "@ai-sdk/google": "https://generativelanguage.googleapis.com/v1beta", + "@ai-sdk/amazon-bedrock": "https://bedrock-runtime.us-east-1.amazonaws.com", +} + +const ROUTE: Record = { + "@ai-sdk/openai": "openai-responses", + "@ai-sdk/azure": "azure-openai-responses", + "@ai-sdk/anthropic": "anthropic-messages", + "@ai-sdk/google": "gemini", + "@ai-sdk/amazon-bedrock": "bedrock-converse", + "@ai-sdk/openai-compatible": "openai-compatible-chat", + "@openrouter/ai-sdk-provider": "openai-compatible-chat", +} + +const isRecord = (value: unknown): value is Record => + typeof value === "object" && value !== null && !Array.isArray(value) + +const providerMetadata = (value: unknown): ProviderMetadata | undefined => { + if (!isRecord(value)) return undefined + const result = Object.fromEntries( + Object.entries(value).filter((entry): entry is [string, Record] => isRecord(entry[1])), + ) + return Object.keys(result).length === 0 ? undefined : result +} + +const textPart = (part: Record) => ({ + type: "text" as const, + text: typeof part.text === "string" ? part.text : "", + providerMetadata: providerMetadata(part.providerOptions), +}) + +const mediaPart = (part: Record) => { + if (typeof part.data !== "string" && !(part.data instanceof Uint8Array)) + throw new Error("Native LLM request adapter only supports file parts with string or Uint8Array data") + return { + type: "media" as const, + mediaType: typeof part.mediaType === "string" ? part.mediaType : "application/octet-stream", + data: part.data, + filename: typeof part.filename === "string" ? part.filename : undefined, + } +} + +const toolResult = (part: Record) => { + const output = isRecord(part.output) ? part.output : { type: "json", value: part.output } + const type = output.type === "text" ? "text" : output.type === "error-text" ? "error" : "json" + return LLM.toolResult({ + id: typeof part.toolCallId === "string" ? part.toolCallId : "", + name: typeof part.toolName === "string" ? part.toolName : "", + result: "value" in output ? output.value : output, + resultType: type, + providerExecuted: typeof part.providerExecuted === "boolean" ? part.providerExecuted : undefined, + providerMetadata: providerMetadata(part.providerOptions), + }) +} + +const contentPart = (part: unknown) => { + if (!isRecord(part)) throw new Error("Native LLM request adapter only supports object content parts") + if (part.type === "text") return textPart(part) + if (part.type === "file") return mediaPart(part) + if (part.type === "reasoning") + return { + type: "reasoning" as const, + text: typeof part.text === "string" ? part.text : "", + providerMetadata: providerMetadata(part.providerOptions), + } + if (part.type === "tool-call") + return LLM.toolCall({ + id: typeof part.toolCallId === "string" ? part.toolCallId : "", + name: typeof part.toolName === "string" ? part.toolName : "", + input: part.input, + providerExecuted: typeof part.providerExecuted === "boolean" ? part.providerExecuted : undefined, + providerMetadata: providerMetadata(part.providerOptions), + }) + if (part.type === "tool-result") return toolResult(part) + throw new Error(`Native LLM request adapter does not support ${String(part.type)} content parts`) +} + +const content = (value: ModelMessage["content"]) => + typeof value === "string" ? [LLM.text(value)] : value.map(contentPart) + +const messages = (input: readonly ModelMessage[]) => { + const system = input.flatMap((message) => (message.role === "system" ? [LLM.system(message.content)] : [])) + const messages = input.flatMap((message) => { + if (message.role === "system") return [] + return [ + LLM.message({ + role: message.role, + content: content(message.content), + native: isRecord(message.providerOptions) ? { providerOptions: message.providerOptions } : undefined, + }), + ] + }) + return { system, messages } +} + +const schema = (value: unknown): JsonSchema => { + if (!isRecord(value)) return { type: "object", properties: {} } + if (isRecord(value.jsonSchema)) return value.jsonSchema + return value +} + +const tools = (input: Record | undefined): ToolDefinition[] => + Object.entries(input ?? {}).map(([name, item]) => + LLM.toolDefinition({ + name, + description: item.description ?? "", + inputSchema: schema(item.inputSchema), + }), + ) + +const generation = (input: RequestInput) => { + const result = { + temperature: input.temperature, + topP: input.topP, + topK: input.topK, + maxTokens: input.maxOutputTokens, + } + return Object.values(result).some((value) => value !== undefined) ? result : undefined +} + +const baseURL = (model: Provider.Model) => { + if (model.api.url) return model.api.url + const fallback = DEFAULT_BASE_URL[model.api.npm] + if (fallback) return fallback + throw new Error(`Native LLM request adapter requires a base URL for ${model.providerID}/${model.id}`) +} + +export const model = (model: Provider.Model, headers?: Record) => { + const route = ROUTE[model.api.npm] + if (!route) throw new Error(`Native LLM request adapter does not support provider package ${model.api.npm}`) + return LLM.model({ + id: model.api.id, + provider: model.providerID, + route, + baseURL: baseURL(model), + headers: Object.keys({ ...model.headers, ...headers }).length === 0 ? undefined : { ...model.headers, ...headers }, + limits: { + context: model.limit.context, + output: model.limit.output, + }, + }) +} + +export const request = (input: RequestInput) => { + const converted = messages(input.messages) + return LLM.request({ + model: model(input.model, input.headers), + system: [...(input.system ?? []).map(LLM.system), ...converted.system], + messages: converted.messages, + tools: tools(input.tools), + toolChoice: input.toolChoice, + generation: generation(input), + providerOptions: input.providerOptions, + }) +} + +export * as LLMNative from "./llm-native" diff --git a/packages/opencode/test/session/llm-native.test.ts b/packages/opencode/test/session/llm-native.test.ts new file mode 100644 index 0000000000..9a8e003d1f --- /dev/null +++ b/packages/opencode/test/session/llm-native.test.ts @@ -0,0 +1,219 @@ +import { describe, expect, test } from "bun:test" +import { jsonSchema, tool, type ModelMessage } from "ai" +import { LLMNative } from "@/session/llm-native" +import type { Provider } from "@/provider/provider" +import { ModelID, ProviderID } from "@/provider/schema" + +const baseModel: Provider.Model = { + id: ModelID.make("gpt-5-mini"), + providerID: ProviderID.make("openai"), + api: { + id: "gpt-5-mini", + url: "https://api.openai.com/v1", + npm: "@ai-sdk/openai", + }, + name: "GPT-5 Mini", + capabilities: { + temperature: true, + reasoning: true, + attachment: true, + toolcall: true, + input: { + text: true, + audio: false, + image: true, + video: false, + pdf: false, + }, + output: { + text: true, + audio: false, + image: false, + video: false, + pdf: false, + }, + interleaved: false, + }, + cost: { + input: 0, + output: 0, + cache: { + read: 0, + write: 0, + }, + }, + limit: { + context: 128_000, + input: 128_000, + output: 32_000, + }, + status: "active", + options: {}, + headers: { + "x-model": "model-header", + }, + release_date: "2026-01-01", +} + +describe("session.llm-native.request", () => { + test("maps normalized stream inputs to a native LLM request", () => { + const messages: ModelMessage[] = [ + { + role: "system", + content: "system from messages", + }, + { + role: "user", + content: [ + { type: "text", text: "hello", providerOptions: { openai: { cacheControl: { type: "ephemeral" } } } }, + { type: "file", mediaType: "image/png", filename: "img.png", data: "data:image/png;base64,Zm9v" }, + ], + }, + { + role: "assistant", + content: [ + { type: "reasoning", text: "thinking", providerOptions: { openai: { encryptedContent: "secret" } } }, + { type: "text", text: "I'll run it" }, + { + type: "tool-call", + toolCallId: "call-1", + toolName: "bash", + input: { command: "ls" }, + providerOptions: { openai: { itemId: "item-1" } }, + }, + ], + }, + { + role: "tool", + content: [ + { + type: "tool-result", + toolCallId: "call-1", + toolName: "bash", + output: { type: "text", value: "ok" }, + providerOptions: { openai: { outputId: "output-1" } }, + }, + ], + }, + ] + + const request = LLMNative.request({ + model: baseModel, + system: ["agent system"], + messages, + tools: { + bash: tool({ + description: "Run a shell command", + inputSchema: jsonSchema({ + type: "object", + properties: { + command: { type: "string" }, + }, + required: ["command"], + }), + }), + }, + toolChoice: "required", + temperature: 0.2, + topP: 0.9, + topK: 40, + maxOutputTokens: 1024, + providerOptions: { openai: { store: false } }, + headers: { "x-request": "request-header" }, + }) + + expect(request.model).toMatchObject({ + id: "gpt-5-mini", + provider: "openai", + route: "openai-responses", + baseURL: "https://api.openai.com/v1", + headers: { + "x-model": "model-header", + "x-request": "request-header", + }, + limits: { + context: 128_000, + output: 32_000, + }, + }) + expect(request.system).toEqual([ + { type: "text", text: "agent system" }, + { type: "text", text: "system from messages" }, + ]) + expect(request.generation).toMatchObject({ + temperature: 0.2, + topP: 0.9, + topK: 40, + maxTokens: 1024, + }) + expect(request.providerOptions).toEqual({ openai: { store: false } }) + expect(request.toolChoice).toMatchObject({ type: "required" }) + expect(request.tools).toMatchObject([ + { + name: "bash", + description: "Run a shell command", + inputSchema: { + type: "object", + properties: { + command: { type: "string" }, + }, + required: ["command"], + }, + }, + ]) + expect(request.messages).toMatchObject([ + { + role: "user", + content: [ + { type: "text", text: "hello", providerMetadata: { openai: { cacheControl: { type: "ephemeral" } } } }, + { type: "media", mediaType: "image/png", filename: "img.png", data: "data:image/png;base64,Zm9v" }, + ], + }, + { + role: "assistant", + content: [ + { type: "reasoning", text: "thinking", providerMetadata: { openai: { encryptedContent: "secret" } } }, + { type: "text", text: "I'll run it" }, + { + type: "tool-call", + id: "call-1", + name: "bash", + input: { command: "ls" }, + providerMetadata: { openai: { itemId: "item-1" } }, + }, + ], + }, + { + role: "tool", + content: [ + { + type: "tool-result", + id: "call-1", + name: "bash", + result: { type: "text", value: "ok" }, + providerMetadata: { openai: { outputId: "output-1" } }, + }, + ], + }, + ]) + }) + + test("selects native routes from existing provider packages", () => { + expect(LLMNative.model({ ...baseModel, api: { ...baseModel.api, npm: "@ai-sdk/anthropic" } }).route).toBe( + "anthropic-messages", + ) + expect(LLMNative.model({ ...baseModel, api: { ...baseModel.api, npm: "@ai-sdk/google" } }).route).toBe("gemini") + expect(LLMNative.model({ ...baseModel, api: { ...baseModel.api, npm: "@ai-sdk/openai-compatible" } }).route).toBe( + "openai-compatible-chat", + ) + }) + + test("fails fast for unsupported provider packages", () => { + expect(() => + LLMNative.request({ + model: { ...baseModel, api: { ...baseModel.api, npm: "unknown-provider" } }, + messages: [], + }), + ).toThrow("Native LLM request adapter does not support provider package unknown-provider") + }) +}) From 3e8fc071eda5c01a2fc0c7322b1c50ebd6c51edc Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 17:32:17 -0400 Subject: [PATCH 05/13] fix(session): target native OpenRouter route --- packages/opencode/src/session/llm-native.ts | 3 +- .../opencode/test/session/llm-native.test.ts | 29 ++++++++++++++----- 2 files changed, 24 insertions(+), 8 deletions(-) diff --git a/packages/opencode/src/session/llm-native.ts b/packages/opencode/src/session/llm-native.ts index 6bb2159412..1c359dcf9d 100644 --- a/packages/opencode/src/session/llm-native.ts +++ b/packages/opencode/src/session/llm-native.ts @@ -27,6 +27,7 @@ const DEFAULT_BASE_URL: Record = { "@ai-sdk/anthropic": "https://api.anthropic.com/v1", "@ai-sdk/google": "https://generativelanguage.googleapis.com/v1beta", "@ai-sdk/amazon-bedrock": "https://bedrock-runtime.us-east-1.amazonaws.com", + "@openrouter/ai-sdk-provider": "https://openrouter.ai/api/v1", } const ROUTE: Record = { @@ -36,7 +37,7 @@ const ROUTE: Record = { "@ai-sdk/google": "gemini", "@ai-sdk/amazon-bedrock": "bedrock-converse", "@ai-sdk/openai-compatible": "openai-compatible-chat", - "@openrouter/ai-sdk-provider": "openai-compatible-chat", + "@openrouter/ai-sdk-provider": "openrouter", } const isRecord = (value: unknown): value is Record => diff --git a/packages/opencode/test/session/llm-native.test.ts b/packages/opencode/test/session/llm-native.test.ts index 9a8e003d1f..81d4a13119 100644 --- a/packages/opencode/test/session/llm-native.test.ts +++ b/packages/opencode/test/session/llm-native.test.ts @@ -199,13 +199,28 @@ describe("session.llm-native.request", () => { }) test("selects native routes from existing provider packages", () => { - expect(LLMNative.model({ ...baseModel, api: { ...baseModel.api, npm: "@ai-sdk/anthropic" } }).route).toBe( - "anthropic-messages", - ) - expect(LLMNative.model({ ...baseModel, api: { ...baseModel.api, npm: "@ai-sdk/google" } }).route).toBe("gemini") - expect(LLMNative.model({ ...baseModel, api: { ...baseModel.api, npm: "@ai-sdk/openai-compatible" } }).route).toBe( - "openai-compatible-chat", - ) + expect( + LLMNative.model({ ...baseModel, api: { ...baseModel.api, url: "", npm: "@ai-sdk/anthropic" } }), + ).toMatchObject({ + route: "anthropic-messages", + baseURL: "https://api.anthropic.com/v1", + }) + expect(LLMNative.model({ ...baseModel, api: { ...baseModel.api, url: "", npm: "@ai-sdk/google" } })).toMatchObject({ + route: "gemini", + baseURL: "https://generativelanguage.googleapis.com/v1beta", + }) + expect( + LLMNative.model({ ...baseModel, api: { ...baseModel.api, npm: "@ai-sdk/openai-compatible" } }), + ).toMatchObject({ + route: "openai-compatible-chat", + baseURL: "https://api.openai.com/v1", + }) + expect( + LLMNative.model({ ...baseModel, api: { ...baseModel.api, url: "", npm: "@openrouter/ai-sdk-provider" } }), + ).toMatchObject({ + route: "openrouter", + baseURL: "https://openrouter.ai/api/v1", + }) }) test("fails fast for unsupported provider packages", () => { From a8a58fd61ec54e9383e29dcbafcbbe89adc15ef9 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 18:04:02 -0400 Subject: [PATCH 06/13] fix(session): use current native LLM constructors --- packages/opencode/src/session/llm-native.ts | 19 ++++++++++--------- 1 file changed, 10 insertions(+), 9 deletions(-) diff --git a/packages/opencode/src/session/llm-native.ts b/packages/opencode/src/session/llm-native.ts index 1c359dcf9d..74be5988ee 100644 --- a/packages/opencode/src/session/llm-native.ts +++ b/packages/opencode/src/session/llm-native.ts @@ -1,5 +1,6 @@ -import type { JsonSchema, LLMRequest, ProviderMetadata, ToolDefinition } from "@opencode-ai/llm" -import { LLM } from "@opencode-ai/llm" +import type { JsonSchema, LLMRequest, ProviderMetadata } from "@opencode-ai/llm" +import { LLM, Message, SystemPart, ToolCallPart, ToolDefinition, ToolResultPart } from "@opencode-ai/llm" +import "@opencode-ai/llm/providers" import type { ModelMessage } from "ai" import type { Provider } from "@/provider/provider" @@ -71,7 +72,7 @@ const mediaPart = (part: Record) => { const toolResult = (part: Record) => { const output = isRecord(part.output) ? part.output : { type: "json", value: part.output } const type = output.type === "text" ? "text" : output.type === "error-text" ? "error" : "json" - return LLM.toolResult({ + return ToolResultPart.make({ id: typeof part.toolCallId === "string" ? part.toolCallId : "", name: typeof part.toolName === "string" ? part.toolName : "", result: "value" in output ? output.value : output, @@ -92,7 +93,7 @@ const contentPart = (part: unknown) => { providerMetadata: providerMetadata(part.providerOptions), } if (part.type === "tool-call") - return LLM.toolCall({ + return ToolCallPart.make({ id: typeof part.toolCallId === "string" ? part.toolCallId : "", name: typeof part.toolName === "string" ? part.toolName : "", input: part.input, @@ -104,14 +105,14 @@ const contentPart = (part: unknown) => { } const content = (value: ModelMessage["content"]) => - typeof value === "string" ? [LLM.text(value)] : value.map(contentPart) + typeof value === "string" ? [{ type: "text" as const, text: value }] : value.map(contentPart) const messages = (input: readonly ModelMessage[]) => { - const system = input.flatMap((message) => (message.role === "system" ? [LLM.system(message.content)] : [])) + const system = input.flatMap((message) => (message.role === "system" ? [SystemPart.make(message.content)] : [])) const messages = input.flatMap((message) => { if (message.role === "system") return [] return [ - LLM.message({ + Message.make({ role: message.role, content: content(message.content), native: isRecord(message.providerOptions) ? { providerOptions: message.providerOptions } : undefined, @@ -129,7 +130,7 @@ const schema = (value: unknown): JsonSchema => { const tools = (input: Record | undefined): ToolDefinition[] => Object.entries(input ?? {}).map(([name, item]) => - LLM.toolDefinition({ + ToolDefinition.make({ name, description: item.description ?? "", inputSchema: schema(item.inputSchema), @@ -173,7 +174,7 @@ export const request = (input: RequestInput) => { const converted = messages(input.messages) return LLM.request({ model: model(input.model, input.headers), - system: [...(input.system ?? []).map(LLM.system), ...converted.system], + system: [...(input.system ?? []).map(SystemPart.make), ...converted.system], messages: converted.messages, tools: tools(input.tools), toolChoice: input.toolChoice, From 9bbb688459a5782e0418cebe61a929087333eb93 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 17:38:45 -0400 Subject: [PATCH 07/13] test(session): compile native LLM requests --- .../opencode/test/session/llm-native.test.ts | 28 +++++++++++++++++++ 1 file changed, 28 insertions(+) diff --git a/packages/opencode/test/session/llm-native.test.ts b/packages/opencode/test/session/llm-native.test.ts index 81d4a13119..40aa71df4d 100644 --- a/packages/opencode/test/session/llm-native.test.ts +++ b/packages/opencode/test/session/llm-native.test.ts @@ -1,5 +1,7 @@ import { describe, expect, test } from "bun:test" +import { LLMClient, RequestExecutor } from "@opencode-ai/llm/route" import { jsonSchema, tool, type ModelMessage } from "ai" +import { Effect } from "effect" import { LLMNative } from "@/session/llm-native" import type { Provider } from "@/provider/provider" import { ModelID, ProviderID } from "@/provider/schema" @@ -231,4 +233,30 @@ describe("session.llm-native.request", () => { }), ).toThrow("Native LLM request adapter does not support provider package unknown-provider") }) + + test("compiles through the native OpenAI Responses route", async () => { + const prepared = await Effect.runPromise( + LLMClient.prepare( + LLMNative.request({ + model: baseModel, + messages: [{ role: "user", content: "hello" }], + providerOptions: { openai: { store: false } }, + maxOutputTokens: 512, + headers: { "x-request": "request-header" }, + }), + ).pipe(Effect.provide(LLMClient.layer), Effect.provide(RequestExecutor.defaultLayer)), + ) + + expect(prepared).toMatchObject({ + route: "openai-responses", + protocol: "openai-responses", + body: { + model: "gpt-5-mini", + input: [{ role: "user", content: [{ type: "input_text", text: "hello" }] }], + max_output_tokens: 512, + store: false, + stream: true, + }, + }) + }) }) From a878036f63a510fcc6057058a8d8b4d5b2169bb2 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 17:44:09 -0400 Subject: [PATCH 08/13] feat(session): add native OpenAI runtime opt-in --- packages/opencode/src/session/llm-native.ts | 10 +- packages/opencode/src/session/llm.ts | 198 ++++++++++++-------- packages/opencode/test/session/llm.test.ts | 116 +++++++++++- 3 files changed, 243 insertions(+), 81 deletions(-) diff --git a/packages/opencode/src/session/llm-native.ts b/packages/opencode/src/session/llm-native.ts index 74be5988ee..817f29832d 100644 --- a/packages/opencode/src/session/llm-native.ts +++ b/packages/opencode/src/session/llm-native.ts @@ -11,6 +11,8 @@ type ToolInput = { export type RequestInput = { readonly model: Provider.Model + readonly apiKey?: string + readonly baseURL?: string readonly system?: readonly string[] readonly messages: readonly ModelMessage[] readonly tools?: Record @@ -154,14 +156,16 @@ const baseURL = (model: Provider.Model) => { throw new Error(`Native LLM request adapter requires a base URL for ${model.providerID}/${model.id}`) } -export const model = (model: Provider.Model, headers?: Record) => { +export const model = (input: Provider.Model | RequestInput, headers?: Record) => { + const model = "model" in input ? input.model : input const route = ROUTE[model.api.npm] if (!route) throw new Error(`Native LLM request adapter does not support provider package ${model.api.npm}`) return LLM.model({ id: model.api.id, provider: model.providerID, route, - baseURL: baseURL(model), + baseURL: "model" in input && input.baseURL ? input.baseURL : baseURL(model), + apiKey: "model" in input ? input.apiKey : undefined, headers: Object.keys({ ...model.headers, ...headers }).length === 0 ? undefined : { ...model.headers, ...headers }, limits: { context: model.limit.context, @@ -173,7 +177,7 @@ export const model = (model: Provider.Model, headers?: Record) = export const request = (input: RequestInput) => { const converted = messages(input.messages) return LLM.request({ - model: model(input.model, input.headers), + model: model(input, input.headers), system: [...(input.system ?? []).map(SystemPart.make), ...converted.system], messages: converted.messages, tools: tools(input.tools), diff --git a/packages/opencode/src/session/llm.ts b/packages/opencode/src/session/llm.ts index 28e05aba8f..3460a04c02 100644 --- a/packages/opencode/src/session/llm.ts +++ b/packages/opencode/src/session/llm.ts @@ -4,6 +4,7 @@ import { Context, Effect, Layer, Record } from "effect" import * as Stream from "effect/Stream" import { streamText, wrapLanguageModel, type ModelMessage, type Tool, tool, jsonSchema } from "ai" import type { LLMEvent } from "@opencode-ai/llm" +import { LLMClient, RequestExecutor } from "@opencode-ai/llm/route" import { mergeDeep } from "remeda" import { GitLabWorkflowLanguageModel } from "gitlab-ai-provider" import { ProviderTransform } from "@/provider/transform" @@ -20,12 +21,12 @@ import { Bus } from "@/bus" import { Wildcard } from "@/util/wildcard" import { SessionID } from "@/session/schema" import { Auth } from "@/auth" -import { Installation } from "@/installation" import { InstallationVersion } from "@opencode-ai/core/installation/version" import { EffectBridge } from "@/effect/bridge" import * as Option from "effect/Option" import * as OtelTracer from "@effect/opentelemetry/Tracer" import { LLMAISDK } from "./llm-ai-sdk" +import { LLMNative } from "./llm-native" const log = Log.create({ service: "llm" }) export const OUTPUT_TOKEN_MAX = ProviderTransform.OUTPUT_TOKEN_MAX @@ -34,6 +35,8 @@ export const OUTPUT_TOKEN_MAX = ProviderTransform.OUTPUT_TOKEN_MAX const mergeOptions = (target: Record, source: Record | undefined): Record => mergeDeep(target, source ?? {}) as Record +const runtime = () => (process.env.OPENCODE_LLM_RUNTIME === "native" ? "native" : "ai-sdk") + export type StreamInput = { user: MessageV2.User sessionID: string @@ -333,86 +336,123 @@ const live: Layer.Layer< ? (yield* InstanceState.context).project.id : undefined - return streamText({ - onError(error) { - l.error("stream error", { - error, - }) - }, - async experimental_repairToolCall(failed) { - const lower = failed.toolCall.toolName.toLowerCase() - if (lower !== failed.toolCall.toolName && sortedTools[lower]) { - l.info("repairing tool call", { - tool: failed.toolCall.toolName, - repaired: lower, + const requestHeaders = { + ...(input.model.providerID.startsWith("opencode") + ? { + ...(opencodeProjectID ? { "x-opencode-project": opencodeProjectID } : {}), + "x-opencode-session": input.sessionID, + "x-opencode-request": input.user.id, + "x-opencode-client": Flag.OPENCODE_CLIENT, + "User-Agent": `opencode/${InstallationVersion}`, + } + : { + "x-session-affinity": input.sessionID, + ...(input.parentSessionID ? { "x-parent-session-id": input.parentSessionID } : {}), + "User-Agent": `opencode/${InstallationVersion}`, + }), + ...input.model.headers, + ...headers, + } + + if (runtime() === "native") { + if (input.model.providerID !== "openai" || input.model.api.npm !== "@ai-sdk/openai") { + return yield* Effect.fail(new Error("Native LLM runtime currently only supports OpenAI models")) + } + if (Object.keys(sortedTools).length > 0) { + return yield* Effect.fail(new Error("Native LLM runtime does not support tools yet")) + } + const apiKey = + info?.type === "api" ? info.key : typeof item.options.apiKey === "string" ? item.options.apiKey : undefined + if (!apiKey) return yield* Effect.fail(new Error("Native LLM runtime requires API key auth for OpenAI")) + const baseURL = typeof item.options.baseURL === "string" ? item.options.baseURL : undefined + return { + type: "native" as const, + stream: LLMClient.stream( + LLMNative.request({ + model: input.model, + apiKey, + baseURL, + system: isOpenaiOauth ? system : [], + messages: ProviderTransform.message(messages, input.model, options), + toolChoice: input.toolChoice, + temperature: params.temperature, + topP: params.topP, + topK: params.topK, + maxOutputTokens: params.maxOutputTokens, + providerOptions: ProviderTransform.providerOptions(input.model, params.options), + headers: requestHeaders, + }), + ).pipe(Stream.provide(LLMClient.layer), Stream.provide(RequestExecutor.defaultLayer)), + } + } + + return { + type: "ai-sdk" as const, + result: streamText({ + onError(error) { + l.error("stream error", { + error, }) + }, + async experimental_repairToolCall(failed) { + const lower = failed.toolCall.toolName.toLowerCase() + if (lower !== failed.toolCall.toolName && sortedTools[lower]) { + l.info("repairing tool call", { + tool: failed.toolCall.toolName, + repaired: lower, + }) + return { + ...failed.toolCall, + toolName: lower, + } + } return { ...failed.toolCall, - toolName: lower, - } - } - return { - ...failed.toolCall, - input: JSON.stringify({ - tool: failed.toolCall.toolName, - error: failed.error.message, - }), - toolName: "invalid", - } - }, - temperature: params.temperature, - topP: params.topP, - topK: params.topK, - providerOptions: ProviderTransform.providerOptions(input.model, params.options), - activeTools: Object.keys(sortedTools).filter((x) => x !== "invalid"), - tools: sortedTools, - toolChoice: input.toolChoice, - maxOutputTokens: params.maxOutputTokens, - abortSignal: input.abort, - headers: { - ...(input.model.providerID.startsWith("opencode") - ? { - "x-opencode-project": opencodeProjectID, - "x-opencode-session": input.sessionID, - "x-opencode-request": input.user.id, - "x-opencode-client": Flag.OPENCODE_CLIENT, - "User-Agent": `opencode/${InstallationVersion}`, - } - : { - "x-session-affinity": input.sessionID, - ...(input.parentSessionID ? { "x-parent-session-id": input.parentSessionID } : {}), - "User-Agent": `opencode/${InstallationVersion}`, + input: JSON.stringify({ + tool: failed.toolCall.toolName, + error: failed.error.message, }), - ...input.model.headers, - ...headers, - }, - maxRetries: input.retries ?? 0, - messages, - model: wrapLanguageModel({ - model: language, - middleware: [ - { - specificationVersion: "v3" as const, - async transformParams(args) { - if (args.type === "stream") { - // @ts-expect-error - args.params.prompt = ProviderTransform.message(args.params.prompt, input.model, options) - } - return args.params - }, - }, - ], - }), - experimental_telemetry: { - isEnabled: cfg.experimental?.openTelemetry, - functionId: "session.llm", - tracer: telemetryTracer, - metadata: { - userId: cfg.username ?? "unknown", - sessionId: input.sessionID, + toolName: "invalid", + } }, - }, - }) + temperature: params.temperature, + topP: params.topP, + topK: params.topK, + providerOptions: ProviderTransform.providerOptions(input.model, params.options), + activeTools: Object.keys(sortedTools).filter((x) => x !== "invalid"), + tools: sortedTools, + toolChoice: input.toolChoice, + maxOutputTokens: params.maxOutputTokens, + abortSignal: input.abort, + headers: requestHeaders, + maxRetries: input.retries ?? 0, + messages, + model: wrapLanguageModel({ + model: language, + middleware: [ + { + specificationVersion: "v3" as const, + async transformParams(args) { + if (args.type === "stream") { + // @ts-expect-error + args.params.prompt = ProviderTransform.message(args.params.prompt, input.model, options) + } + return args.params + }, + }, + ], + }), + experimental_telemetry: { + isEnabled: cfg.experimental?.openTelemetry, + functionId: "session.llm", + tracer: telemetryTracer, + metadata: { + userId: cfg.username ?? "unknown", + sessionId: input.sessionID, + }, + }, + }), + } }) const stream: Interface["stream"] = (input) => @@ -426,8 +466,12 @@ const live: Layer.Layer< const result = yield* run({ ...input, abort: ctrl.signal }) + if (result.type === "native") return result.stream + const state = LLMAISDK.adapterState() - return Stream.fromAsyncIterable(result.fullStream, (e) => (e instanceof Error ? e : new Error(String(e)))).pipe( + return Stream.fromAsyncIterable(result.result.fullStream, (e) => + e instanceof Error ? e : new Error(String(e)), + ).pipe( Stream.mapEffect((event) => LLMAISDK.toLLMEvents(state, event)), Stream.flatMap((events) => Stream.fromIterable(events)), ) diff --git a/packages/opencode/test/session/llm.test.ts b/packages/opencode/test/session/llm.test.ts index f078e352ac..16a693dfa0 100644 --- a/packages/opencode/test/session/llm.test.ts +++ b/packages/opencode/test/session/llm.test.ts @@ -5,7 +5,6 @@ import { Cause, Effect, Exit, Stream } from "effect" import z from "zod" import { makeRuntime } from "../../src/effect/run-service" import { LLM } from "../../src/session/llm" -import { Instance } from "../../src/project/instance" import { WithInstance } from "../../src/project/with-instance" import { Provider } from "@/provider/provider" import { ProviderTransform } from "@/provider/transform" @@ -688,6 +687,121 @@ describe("session.llm.stream", () => { }) }) + test("streams OpenAI through native runtime when opted in", async () => { + const server = state.server + if (!server) { + throw new Error("Server not initialized") + } + + const source = await loadFixture("openai", "gpt-5.2") + const model = source.model + const chunks = [ + { + type: "response.created", + response: { + id: "resp-native", + }, + }, + { + type: "response.output_item.added", + item: { type: "message", id: "item-native", status: "in_progress" }, + }, + { + type: "response.output_text.delta", + item_id: "item-native", + delta: "Hello native", + }, + { + type: "response.completed", + response: { + incomplete_details: null, + usage: { + input_tokens: 1, + input_tokens_details: null, + output_tokens: 1, + output_tokens_details: null, + }, + }, + }, + ] + const request = waitRequest("/responses", createEventResponse(chunks, true)) + + await using tmp = await tmpdir({ + init: async (dir) => { + await Bun.write( + path.join(dir, "opencode.json"), + JSON.stringify({ + $schema: "https://opencode.ai/config.json", + enabled_providers: ["openai"], + provider: { + openai: { + name: "OpenAI", + env: ["OPENAI_API_KEY"], + npm: "@ai-sdk/openai", + api: "https://api.openai.com/v1", + models: { + [model.id]: model, + }, + options: { + apiKey: "test-openai-key", + baseURL: `${server.url.origin}/v1`, + }, + }, + }, + }), + ) + }, + }) + + await WithInstance.provide({ + directory: tmp.path, + fn: async () => { + const previous = process.env.OPENCODE_LLM_RUNTIME + process.env.OPENCODE_LLM_RUNTIME = "native" + try { + const resolved = await getModel(ProviderID.openai, ModelID.make(model.id)) + const sessionID = SessionID.make("session-test-native") + const agent = { + name: "test", + mode: "primary", + options: {}, + permission: [{ permission: "*", pattern: "*", action: "allow" }], + temperature: 0.2, + } satisfies Agent.Info + + await drain({ + user: { + id: MessageID.make("msg_user-native"), + sessionID, + role: "user", + time: { created: Date.now() }, + agent: agent.name, + model: { providerID: ProviderID.make("openai"), modelID: resolved.id, variant: "high" }, + } satisfies MessageV2.User, + sessionID, + model: resolved, + agent, + system: ["You are a helpful assistant."], + messages: [{ role: "user", content: "Hello" }], + tools: {}, + }) + } finally { + if (previous === undefined) delete process.env.OPENCODE_LLM_RUNTIME + else process.env.OPENCODE_LLM_RUNTIME = previous + } + + const capture = await request + expect(capture.url.pathname.endsWith("/responses")).toBe(true) + expect(capture.headers.get("Authorization")).toBe("Bearer test-openai-key") + expect(capture.body.model).toBe(model.id) + expect(capture.body.stream).toBe(true) + expect((capture.body.reasoning as { effort?: string } | undefined)?.effort).toBe("high") + expect(JSON.stringify(capture.body.input)).toContain("You are a helpful assistant.") + expect(capture.body.input).toContainEqual({ role: "user", content: [{ type: "input_text", text: "Hello" }] }) + }, + }) + }) + test("accepts user image attachments as data URLs for OpenAI models", async () => { const server = state.server if (!server) { From 03b60f7623946ad80f698319accc3b874930a63b Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 19:42:38 -0400 Subject: [PATCH 09/13] fix(session): finish native request streams --- packages/opencode/src/session/processor.ts | 143 ++++++++++++--------- 1 file changed, 80 insertions(+), 63 deletions(-) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 2d78dc70e8..a207d4db74 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -293,12 +293,78 @@ export const layer: Layer.Layer< return { title: value.name, metadata: value.result.type === "json" && isRecord(value.result.value) ? value.result.value : {}, - output: typeof value.result.value === "string" ? value.result.value : (JSON.stringify(value.result.value) ?? ""), + output: + typeof value.result.value === "string" ? value.result.value : (JSON.stringify(value.result.value) ?? ""), } } const toolInput = (value: unknown): Record => (isRecord(value) ? value : { value }) + const finishStep = Effect.fn("SessionProcessor.finishStep")(function* ( + value: Extract, + ) { + const completedSnapshot = yield* snapshot.track() + yield* Effect.forEach(Object.keys(ctx.reasoningMap), finishReasoning) + const usage = Session.getUsage({ + model: ctx.model, + usage: value.usage ?? new 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.reason, + cost: usage.cost, + tokens: usage.tokens, + snapshot: completedSnapshot, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } + } + ctx.assistantMessage.finish = value.reason + ctx.assistantMessage.cost += usage.cost + ctx.assistantMessage.tokens = usage.tokens + yield* session.updatePart({ + id: PartID.ascending(), + reason: value.reason, + 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 + } + }) + const handleEvent = Effect.fnUntraced(function* (value: StreamEvent) { switch (value.type) { case "request-start": @@ -581,66 +647,7 @@ export const layer: Layer.Layer< return case "step-finish": { - const completedSnapshot = yield* snapshot.track() - yield* Effect.forEach(Object.keys(ctx.reasoningMap), finishReasoning) - const usage = Session.getUsage({ - model: ctx.model, - usage: value.usage ?? new 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.reason, - cost: usage.cost, - tokens: usage.tokens, - snapshot: completedSnapshot, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } - } - ctx.assistantMessage.finish = value.reason - ctx.assistantMessage.cost += usage.cost - ctx.assistantMessage.tokens = usage.tokens - yield* session.updatePart({ - id: PartID.ascending(), - reason: value.reason, - 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 - } + yield* finishStep(value) return } @@ -667,7 +674,17 @@ export const layer: Layer.Layer< return case "text-delta": - if (!ctx.currentText) return + if (!ctx.currentText) { + ctx.currentText = { + id: PartID.ascending(), + messageID: ctx.assistantMessage.id, + sessionID: ctx.assistantMessage.sessionID, + type: "text", + text: "", + time: { start: Date.now() }, + } + yield* session.updatePart(ctx.currentText) + } ctx.currentText.text += value.text yield* session.updatePartDelta({ sessionID: ctx.currentText.sessionID, @@ -711,8 +728,8 @@ export const layer: Layer.Layer< return case "request-finish": + if (!ctx.assistantMessage.finish) yield* finishStep(value) return - } }) From 4c75b463e78905207d1503d6049f31846c4b2c72 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 20:11:16 -0400 Subject: [PATCH 10/13] fix(session): finalize native text before request finish --- packages/opencode/src/session/processor.ts | 63 ++++++++++++---------- 1 file changed, 34 insertions(+), 29 deletions(-) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index a207d4db74..1290cc545f 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -365,6 +365,38 @@ export const layer: Layer.Layer< } }) + const finishText = Effect.fn("SessionProcessor.finishText")(function* ( + providerMetadata?: Extract["providerMetadata"], + ) { + 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 (providerMetadata) ctx.currentText.metadata = providerMetadata + yield* session.updatePart(ctx.currentText) + ctx.currentText = undefined + }) + const handleEvent = Effect.fnUntraced(function* (value: StreamEvent) { switch (value.type) { case "request-start": @@ -696,38 +728,11 @@ export const layer: Layer.Layer< 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 + yield* finishText(value.providerMetadata) return case "request-finish": + yield* finishText() if (!ctx.assistantMessage.finish) yield* finishStep(value) return } From a6e8ec4b35d7894c85bdb9f45627bf5102f03631 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 21:33:52 -0400 Subject: [PATCH 11/13] fix(session): rely on native LLM lifecycle events --- packages/opencode/src/session/processor.ts | 206 +++++++++------------ 1 file changed, 92 insertions(+), 114 deletions(-) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 1290cc545f..2d78dc70e8 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -293,110 +293,12 @@ export const layer: Layer.Layer< return { title: value.name, metadata: value.result.type === "json" && isRecord(value.result.value) ? value.result.value : {}, - output: - typeof value.result.value === "string" ? value.result.value : (JSON.stringify(value.result.value) ?? ""), + output: typeof value.result.value === "string" ? value.result.value : (JSON.stringify(value.result.value) ?? ""), } } const toolInput = (value: unknown): Record => (isRecord(value) ? value : { value }) - const finishStep = Effect.fn("SessionProcessor.finishStep")(function* ( - value: Extract, - ) { - const completedSnapshot = yield* snapshot.track() - yield* Effect.forEach(Object.keys(ctx.reasoningMap), finishReasoning) - const usage = Session.getUsage({ - model: ctx.model, - usage: value.usage ?? new 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.reason, - cost: usage.cost, - tokens: usage.tokens, - snapshot: completedSnapshot, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } - } - ctx.assistantMessage.finish = value.reason - ctx.assistantMessage.cost += usage.cost - ctx.assistantMessage.tokens = usage.tokens - yield* session.updatePart({ - id: PartID.ascending(), - reason: value.reason, - 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 - } - }) - - const finishText = Effect.fn("SessionProcessor.finishText")(function* ( - providerMetadata?: Extract["providerMetadata"], - ) { - 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 (providerMetadata) ctx.currentText.metadata = providerMetadata - yield* session.updatePart(ctx.currentText) - ctx.currentText = undefined - }) - const handleEvent = Effect.fnUntraced(function* (value: StreamEvent) { switch (value.type) { case "request-start": @@ -679,7 +581,66 @@ export const layer: Layer.Layer< return case "step-finish": { - yield* finishStep(value) + const completedSnapshot = yield* snapshot.track() + yield* Effect.forEach(Object.keys(ctx.reasoningMap), finishReasoning) + const usage = Session.getUsage({ + model: ctx.model, + usage: value.usage ?? new 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.reason, + cost: usage.cost, + tokens: usage.tokens, + snapshot: completedSnapshot, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } + } + ctx.assistantMessage.finish = value.reason + ctx.assistantMessage.cost += usage.cost + ctx.assistantMessage.tokens = usage.tokens + yield* session.updatePart({ + id: PartID.ascending(), + reason: value.reason, + 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 } @@ -706,17 +667,7 @@ export const layer: Layer.Layer< return case "text-delta": - if (!ctx.currentText) { - ctx.currentText = { - id: PartID.ascending(), - messageID: ctx.assistantMessage.id, - sessionID: ctx.assistantMessage.sessionID, - type: "text", - text: "", - time: { start: Date.now() }, - } - yield* session.updatePart(ctx.currentText) - } + if (!ctx.currentText) return ctx.currentText.text += value.text yield* session.updatePartDelta({ sessionID: ctx.currentText.sessionID, @@ -728,13 +679,40 @@ export const layer: Layer.Layer< return case "text-end": - yield* finishText(value.providerMetadata) + 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 "request-finish": - yield* finishText() - if (!ctx.assistantMessage.finish) yield* finishStep(value) return + } }) From a2b0ef98237380d9403a5430529d0054cca822bf Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 21:43:24 -0400 Subject: [PATCH 12/13] feat(session): execute tools in native LLM runtime --- packages/llm/src/tool-runtime.ts | 8 +- packages/llm/src/tool.ts | 11 +- packages/opencode/src/session/llm.ts | 76 ++++++++---- packages/opencode/test/session/llm.test.ts | 128 +++++++++++++++++++++ 4 files changed, 195 insertions(+), 28 deletions(-) diff --git a/packages/llm/src/tool-runtime.ts b/packages/llm/src/tool-runtime.ts index f464525827..a875d2e438 100644 --- a/packages/llm/src/tool-runtime.ts +++ b/packages/llm/src/tool-runtime.ts @@ -200,17 +200,17 @@ const dispatch = (tools: Tools, call: ToolCallPart): Effect.Effect Effect.succeed({ type: "error" as const, value: failure.message } satisfies ToolResultValue), ), ) } -const decodeAndExecute = (tool: AnyTool, input: unknown): Effect.Effect => - tool._decode(input).pipe( +const decodeAndExecute = (tool: AnyTool, call: ToolCallPart): Effect.Effect => + tool._decode(call.input).pipe( Effect.mapError((error) => new ToolFailure({ message: `Invalid tool input: ${error.message}` })), - Effect.flatMap((decoded) => tool.execute!(decoded)), + Effect.flatMap((decoded) => tool.execute!(decoded, { id: call.id, name: call.name })), Effect.flatMap((value) => tool._encode(value).pipe( Effect.mapError( diff --git a/packages/llm/src/tool.ts b/packages/llm/src/tool.ts index 311c8798b6..14e22688ae 100644 --- a/packages/llm/src/tool.ts +++ b/packages/llm/src/tool.ts @@ -11,6 +11,7 @@ export type ToolSchema = Schema.Codec export type ToolExecute, Success extends ToolSchema> = ( params: Schema.Schema.Type, + context?: { readonly id: string; readonly name: string }, ) => Effect.Effect, ToolFailure> /** @@ -61,7 +62,10 @@ type TypedToolConfig = { type DynamicToolConfig = { readonly description: string readonly jsonSchema: JsonSchema.JsonSchema - readonly execute?: (params: unknown) => Effect.Effect + readonly execute?: ( + params: unknown, + context?: { readonly id: string; readonly name: string }, + ) => Effect.Effect } /** @@ -110,7 +114,10 @@ export function make, Success extends ToolSch export function make(config: { readonly description: string readonly jsonSchema: JsonSchema.JsonSchema - readonly execute: (params: unknown) => Effect.Effect + readonly execute: ( + params: unknown, + context?: { readonly id: string; readonly name: string }, + ) => Effect.Effect }): AnyExecutableTool export function make(config: { readonly description: string diff --git a/packages/opencode/src/session/llm.ts b/packages/opencode/src/session/llm.ts index 3460a04c02..32fb4cf880 100644 --- a/packages/opencode/src/session/llm.ts +++ b/packages/opencode/src/session/llm.ts @@ -2,8 +2,8 @@ import { Provider } from "@/provider/provider" import * as Log from "@opencode-ai/core/util/log" import { Context, Effect, Layer, Record } from "effect" import * as Stream from "effect/Stream" -import { streamText, wrapLanguageModel, type ModelMessage, type Tool, tool, jsonSchema } from "ai" -import type { LLMEvent } from "@opencode-ai/llm" +import { streamText, wrapLanguageModel, type ModelMessage, type Tool, tool as aiTool, jsonSchema, asSchema } from "ai" +import { tool as nativeTool, ToolFailure, type JsonSchema, type LLMEvent } from "@opencode-ai/llm" import { LLMClient, RequestExecutor } from "@opencode-ai/llm/route" import { mergeDeep } from "remeda" import { GitLabWorkflowLanguageModel } from "gitlab-ai-provider" @@ -18,6 +18,7 @@ import { Flag } from "@opencode-ai/core/flag/flag" import { Permission } from "@/permission" import { PermissionID } from "@/permission/schema" import { Bus } from "@/bus" +import { errorMessage } from "@/util/error" import { Wildcard } from "@/util/wildcard" import { SessionID } from "@/session/schema" import { Auth } from "@/auth" @@ -216,7 +217,7 @@ const live: Layer.Layer< Object.keys(tools).length === 0 && hasToolCalls(input.messages) ) { - tools["_noop"] = tool({ + tools["_noop"] = aiTool({ description: "Do not call this tool. It exists only for API compatibility and must never be invoked.", inputSchema: jsonSchema({ type: "object", @@ -358,31 +359,31 @@ const live: Layer.Layer< if (input.model.providerID !== "openai" || input.model.api.npm !== "@ai-sdk/openai") { return yield* Effect.fail(new Error("Native LLM runtime currently only supports OpenAI models")) } - if (Object.keys(sortedTools).length > 0) { - return yield* Effect.fail(new Error("Native LLM runtime does not support tools yet")) - } const apiKey = info?.type === "api" ? info.key : typeof item.options.apiKey === "string" ? item.options.apiKey : undefined if (!apiKey) return yield* Effect.fail(new Error("Native LLM runtime requires API key auth for OpenAI")) const baseURL = typeof item.options.baseURL === "string" ? item.options.baseURL : undefined + const request = LLMNative.request({ + model: input.model, + apiKey, + baseURL, + system: isOpenaiOauth ? system : [], + messages: ProviderTransform.message(messages, input.model, options), + tools: sortedTools, + toolChoice: input.toolChoice, + temperature: params.temperature, + topP: params.topP, + topK: params.topK, + maxOutputTokens: params.maxOutputTokens, + providerOptions: ProviderTransform.providerOptions(input.model, params.options), + headers: requestHeaders, + }) return { type: "native" as const, - stream: LLMClient.stream( - LLMNative.request({ - model: input.model, - apiKey, - baseURL, - system: isOpenaiOauth ? system : [], - messages: ProviderTransform.message(messages, input.model, options), - toolChoice: input.toolChoice, - temperature: params.temperature, - topP: params.topP, - topK: params.topK, - maxOutputTokens: params.maxOutputTokens, - providerOptions: ProviderTransform.providerOptions(input.model, params.options), - headers: requestHeaders, - }), - ).pipe(Stream.provide(LLMClient.layer), Stream.provide(RequestExecutor.defaultLayer)), + stream: LLMClient.stream({ request, tools: nativeTools(sortedTools, input) }).pipe( + Stream.provide(LLMClient.layer), + Stream.provide(RequestExecutor.defaultLayer), + ), } } @@ -502,6 +503,37 @@ function resolveTools(input: Pick input.user.tools?.[k] !== false && !disabled.has(k)) } +function nativeSchema(value: unknown): JsonSchema { + if (!value || typeof value !== "object") return { type: "object", properties: {} } + if ("jsonSchema" in value && value.jsonSchema && typeof value.jsonSchema === "object") + return value.jsonSchema as JsonSchema + return asSchema(value as Parameters[0]).jsonSchema as JsonSchema +} + +function nativeTools(tools: Record, input: StreamRequest) { + return Object.fromEntries( + Object.entries(tools).map(([name, item]) => [ + name, + nativeTool({ + description: item.description ?? "", + jsonSchema: nativeSchema(item.inputSchema), + execute: (args: unknown, ctx?: { readonly id: string; readonly name: string }) => + Effect.tryPromise({ + try: () => { + if (!item.execute) throw new Error(`Tool has no execute handler: ${name}`) + return item.execute(args, { + toolCallId: ctx?.id ?? name, + messages: input.messages, + abortSignal: input.abort, + }) + }, + catch: (error) => new ToolFailure({ message: errorMessage(error) }), + }), + }), + ]), + ) +} + // Check if messages contain any tool-call content // Used to determine if a dummy tool should be added for LiteLLM proxy compatibility export function hasToolCalls(messages: ModelMessage[]): boolean { diff --git a/packages/opencode/test/session/llm.test.ts b/packages/opencode/test/session/llm.test.ts index 16a693dfa0..2054d05343 100644 --- a/packages/opencode/test/session/llm.test.ts +++ b/packages/opencode/test/session/llm.test.ts @@ -802,6 +802,134 @@ describe("session.llm.stream", () => { }) }) + test("executes OpenAI tool calls through native runtime", async () => { + const server = state.server + if (!server) { + throw new Error("Server not initialized") + } + + const source = await loadFixture("openai", "gpt-5.2") + const model = source.model + const chunks = [ + { + type: "response.output_item.added", + item: { type: "function_call", id: "item-native-tool", call_id: "call-native-tool", name: "lookup" }, + }, + { + type: "response.function_call_arguments.delta", + item_id: "item-native-tool", + delta: '{"query":"weather"}', + }, + { + type: "response.output_item.done", + item: { + type: "function_call", + id: "item-native-tool", + call_id: "call-native-tool", + name: "lookup", + arguments: '{"query":"weather"}', + }, + }, + { + type: "response.completed", + response: { incomplete_details: null, usage: { input_tokens: 1, output_tokens: 1 } }, + }, + ] + const request = waitRequest("/responses", createEventResponse(chunks, true)) + let executed: unknown + + await using tmp = await tmpdir({ + init: async (dir) => { + await Bun.write( + path.join(dir, "opencode.json"), + JSON.stringify({ + $schema: "https://opencode.ai/config.json", + enabled_providers: ["openai"], + provider: { + openai: { + name: "OpenAI", + env: ["OPENAI_API_KEY"], + npm: "@ai-sdk/openai", + api: "https://api.openai.com/v1", + models: { + [model.id]: model, + }, + options: { + apiKey: "test-openai-key", + baseURL: `${server.url.origin}/v1`, + }, + }, + }, + }), + ) + }, + }) + + await WithInstance.provide({ + directory: tmp.path, + fn: async () => { + const previous = process.env.OPENCODE_LLM_RUNTIME + process.env.OPENCODE_LLM_RUNTIME = "native" + try { + const resolved = await getModel(ProviderID.openai, ModelID.make(model.id)) + const sessionID = SessionID.make("session-test-native-tool") + const agent = { + name: "test", + mode: "primary", + options: {}, + permission: [{ permission: "*", pattern: "*", action: "allow" }], + } satisfies Agent.Info + + await drain({ + user: { + id: MessageID.make("msg_user-native-tool"), + sessionID, + role: "user", + time: { created: Date.now() }, + agent: agent.name, + model: { providerID: ProviderID.make("openai"), modelID: resolved.id }, + } satisfies MessageV2.User, + sessionID, + model: resolved, + agent, + system: [], + messages: [{ role: "user", content: "Use lookup" }], + tools: { + lookup: tool({ + description: "Lookup data", + inputSchema: z.object({ query: z.string() }), + execute: async (args, options) => { + executed = { args, toolCallId: options.toolCallId } + return { output: "looked up" } + }, + }), + }, + }) + } finally { + if (previous === undefined) delete process.env.OPENCODE_LLM_RUNTIME + else process.env.OPENCODE_LLM_RUNTIME = previous + } + + const capture = await request + expect(capture.body.tools).toEqual([ + { + type: "function", + name: "lookup", + description: "Lookup data", + parameters: { + type: "object", + properties: { query: { type: "string" } }, + required: ["query"], + additionalProperties: false, + $schema: "http://json-schema.org/draft-07/schema#", + }, + }, + ]) + expect(executed).toEqual({ args: { query: "weather" }, toolCallId: "call-native-tool" }) + }, + }) + }) + test("accepts user image attachments as data URLs for OpenAI models", async () => { const server = state.server if (!server) { From d37bc3e71fd6083a7fc9cd2991f16851c7ff186f Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 11 May 2026 22:24:09 -0400 Subject: [PATCH 13/13] refactor(session): inject native LLM client --- packages/opencode/src/session/llm.ts | 10 +- packages/opencode/test/session/llm.test.ts | 211 +++++++++++++++------ 2 files changed, 162 insertions(+), 59 deletions(-) diff --git a/packages/opencode/src/session/llm.ts b/packages/opencode/src/session/llm.ts index 32fb4cf880..e7cdf8d18b 100644 --- a/packages/opencode/src/session/llm.ts +++ b/packages/opencode/src/session/llm.ts @@ -5,6 +5,7 @@ import * as Stream from "effect/Stream" import { streamText, wrapLanguageModel, type ModelMessage, type Tool, tool as aiTool, jsonSchema, asSchema } from "ai" import { tool as nativeTool, ToolFailure, type JsonSchema, type LLMEvent } from "@opencode-ai/llm" import { LLMClient, RequestExecutor } from "@opencode-ai/llm/route" +import type { LLMClientService } from "@opencode-ai/llm/route" import { mergeDeep } from "remeda" import { GitLabWorkflowLanguageModel } from "gitlab-ai-provider" import { ProviderTransform } from "@/provider/transform" @@ -66,7 +67,7 @@ export class Service extends Context.Service()("@opencode/LL const live: Layer.Layer< Service, never, - Auth.Service | Config.Service | Provider.Service | Plugin.Service | Permission.Service + Auth.Service | Config.Service | Provider.Service | Plugin.Service | Permission.Service | LLMClientService > = Layer.effect( Service, Effect.gen(function* () { @@ -75,6 +76,7 @@ const live: Layer.Layer< const provider = yield* Provider.Service const plugin = yield* Plugin.Service const perm = yield* Permission.Service + const llmClient = yield* LLMClient.Service const run = Effect.fn("LLM.run")(function* (input: StreamRequest) { const l = log @@ -380,10 +382,7 @@ const live: Layer.Layer< }) return { type: "native" as const, - stream: LLMClient.stream({ request, tools: nativeTools(sortedTools, input) }).pipe( - Stream.provide(LLMClient.layer), - Stream.provide(RequestExecutor.defaultLayer), - ), + stream: llmClient.stream({ request, tools: nativeTools(sortedTools, input) }), } } @@ -492,6 +491,7 @@ export const defaultLayer = Layer.suspend(() => Layer.provide(Config.defaultLayer), Layer.provide(Provider.defaultLayer), Layer.provide(Plugin.defaultLayer), + Layer.provide(LLMClient.layer.pipe(Layer.provide(RequestExecutor.defaultLayer))), ), ) diff --git a/packages/opencode/test/session/llm.test.ts b/packages/opencode/test/session/llm.test.ts index 2054d05343..1ca282bb06 100644 --- a/packages/opencode/test/session/llm.test.ts +++ b/packages/opencode/test/session/llm.test.ts @@ -1,14 +1,19 @@ import { afterAll, beforeAll, beforeEach, describe, expect, test } from "bun:test" import path from "path" import { tool, type ModelMessage } from "ai" -import { Cause, Effect, Exit, Stream } from "effect" +import { Cause, Effect, Exit, Layer, Stream } from "effect" +import { HttpClientRequest, HttpClientResponse } from "effect/unstable/http" import z from "zod" -import { makeRuntime } from "../../src/effect/run-service" +import { attach, makeRuntime } from "../../src/effect/run-service" import { LLM } from "../../src/session/llm" +import { LLMClient, RequestExecutor } from "@opencode-ai/llm/route" import { WithInstance } from "../../src/project/with-instance" +import { Auth } from "@/auth" +import { Config } from "@/config/config" import { Provider } from "@/provider/provider" import { ProviderTransform } from "@/provider/transform" import { ModelsDev } from "@/provider/models" +import { Plugin } from "@/plugin" import { ProviderID, ModelID } from "../../src/provider/schema" import { Filesystem } from "@/util/filesystem" import { tmpdir } from "../fixture/fixture" @@ -17,6 +22,29 @@ import { MessageV2 } from "../../src/session/message-v2" import { SessionID, MessageID } from "../../src/session/schema" import { AppRuntime } from "../../src/effect/app-runtime" +const openAIConfig = (model: ModelsDev.Provider["models"][string], baseURL: string): Partial => { + const { experimental: _experimental, ...configModel } = model + type ConfigModel = NonNullable[string]["models"]>[string] + return { + enabled_providers: ["openai"], + provider: { + openai: { + name: "OpenAI", + env: ["OPENAI_API_KEY"], + npm: "@ai-sdk/openai", + api: "https://api.openai.com/v1", + models: { + [model.id]: JSON.parse(JSON.stringify(configModel)) as ConfigModel, + }, + options: { + apiKey: "test-openai-key", + baseURL, + }, + }, + }, + } +} + async function getModel(providerID: ProviderID, modelID: ModelID) { return AppRuntime.runPromise( Effect.gen(function* () { @@ -32,6 +60,22 @@ async function drain(input: LLM.StreamInput) { return llm.runPromise((svc) => svc.stream(input).pipe(Stream.runDrain)) } +async function drainWith(layer: Layer.Layer, input: LLM.StreamInput) { + return Effect.runPromise( + attach(LLM.Service.use((svc) => svc.stream(input).pipe(Stream.runDrain))).pipe(Effect.provide(layer)), + ) +} + +function llmLayerWithExecutor(executor: Layer.Layer) { + return LLM.layer.pipe( + Layer.provide(Auth.defaultLayer), + Layer.provide(Config.defaultLayer), + Layer.provide(Provider.defaultLayer), + Layer.provide(Plugin.defaultLayer), + Layer.provide(LLMClient.layer.pipe(Layer.provide(executor))), + ) +} + describe("session.llm.hasToolCalls", () => { test("returns false for empty messages array", () => { expect(LLM.hasToolCalls([])).toBe(false) @@ -614,32 +658,7 @@ describe("session.llm.stream", () => { ] const request = waitRequest("/responses", createEventResponse(responseChunks, true)) - await using tmp = await tmpdir({ - init: async (dir) => { - await Bun.write( - path.join(dir, "opencode.json"), - JSON.stringify({ - $schema: "https://opencode.ai/config.json", - enabled_providers: ["openai"], - provider: { - openai: { - name: "OpenAI", - env: ["OPENAI_API_KEY"], - npm: "@ai-sdk/openai", - api: "https://api.openai.com/v1", - models: { - [model.id]: model, - }, - options: { - apiKey: "test-openai-key", - baseURL: `${server.url.origin}/v1`, - }, - }, - }, - }), - ) - }, - }) + await using tmp = await tmpdir({ config: openAIConfig(model, `${server.url.origin}/v1`) }) await WithInstance.provide({ directory: tmp.path, @@ -726,32 +745,7 @@ describe("session.llm.stream", () => { ] const request = waitRequest("/responses", createEventResponse(chunks, true)) - await using tmp = await tmpdir({ - init: async (dir) => { - await Bun.write( - path.join(dir, "opencode.json"), - JSON.stringify({ - $schema: "https://opencode.ai/config.json", - enabled_providers: ["openai"], - provider: { - openai: { - name: "OpenAI", - env: ["OPENAI_API_KEY"], - npm: "@ai-sdk/openai", - api: "https://api.openai.com/v1", - models: { - [model.id]: model, - }, - options: { - apiKey: "test-openai-key", - baseURL: `${server.url.origin}/v1`, - }, - }, - }, - }), - ) - }, - }) + await using tmp = await tmpdir({ config: openAIConfig(model, `${server.url.origin}/v1`) }) await WithInstance.provide({ directory: tmp.path, @@ -802,6 +796,115 @@ describe("session.llm.stream", () => { }) }) + test("uses injected native request executor for tool calls", async () => { + const source = await loadFixture("openai", "gpt-5.2") + const model = source.model + const chunks = [ + { + type: "response.output_item.added", + item: { type: "function_call", id: "item-injected-tool", call_id: "call-injected-tool", name: "lookup" }, + }, + { + type: "response.function_call_arguments.delta", + item_id: "item-injected-tool", + delta: '{"query":"weather"}', + }, + { + type: "response.output_item.done", + item: { + type: "function_call", + id: "item-injected-tool", + call_id: "call-injected-tool", + name: "lookup", + arguments: '{"query":"weather"}', + }, + }, + { + type: "response.completed", + response: { incomplete_details: null, usage: { input_tokens: 1, output_tokens: 1 } }, + }, + ] + let captured: Record | undefined + let executed: unknown + const executor = Layer.succeed( + RequestExecutor.Service, + RequestExecutor.Service.of({ + execute: (request) => + Effect.gen(function* () { + const web = yield* HttpClientRequest.toWeb(request).pipe(Effect.orDie) + captured = (yield* Effect.promise(() => web.json())) as Record + return HttpClientResponse.fromWeb(request, createEventResponse(chunks, true)) + }), + }), + ) + + await using tmp = await tmpdir({ config: openAIConfig(model, "https://injected-openai.test/v1") }) + + await WithInstance.provide({ + directory: tmp.path, + fn: async () => { + const previous = process.env.OPENCODE_LLM_RUNTIME + process.env.OPENCODE_LLM_RUNTIME = "native" + try { + const resolved = await getModel(ProviderID.openai, ModelID.make(model.id)) + const sessionID = SessionID.make("session-test-native-injected-tool") + const agent = { + name: "test", + mode: "primary", + options: {}, + permission: [{ permission: "*", pattern: "*", action: "allow" }], + } satisfies Agent.Info + + await drainWith(llmLayerWithExecutor(executor), { + user: { + id: MessageID.make("msg_user-native-injected-tool"), + sessionID, + role: "user", + time: { created: Date.now() }, + agent: agent.name, + model: { providerID: ProviderID.make("openai"), modelID: resolved.id }, + } satisfies MessageV2.User, + sessionID, + model: resolved, + agent, + system: [], + messages: [{ role: "user", content: "Use lookup" }], + tools: { + lookup: tool({ + description: "Lookup data", + inputSchema: z.object({ query: z.string() }), + execute: async (args, options) => { + executed = { args, toolCallId: options.toolCallId } + return { output: "looked up" } + }, + }), + }, + }) + } finally { + if (previous === undefined) delete process.env.OPENCODE_LLM_RUNTIME + else process.env.OPENCODE_LLM_RUNTIME = previous + } + + expect(captured?.model).toBe(model.id) + expect(captured?.tools).toEqual([ + { + type: "function", + name: "lookup", + description: "Lookup data", + parameters: { + type: "object", + properties: { query: { type: "string" } }, + required: ["query"], + additionalProperties: false, + $schema: "http://json-schema.org/draft-07/schema#", + }, + }, + ]) + expect(executed).toEqual({ args: { query: "weather" }, toolCallId: "call-injected-tool" }) + }, + }) + }) + test("executes OpenAI tool calls through native runtime", async () => { const server = state.server if (!server) {