Compare commits

...

3 Commits

Author SHA1 Message Date
Aiden Cline 3e3910bc84 fix(ai): derive Anthropic tool finish reason 2026-08-05 01:06:17 +00:00
Aiden Cline 4ed343706d fix: retry empty incomplete streams (#40535) 2026-08-04 19:38:56 -05:00
opencode-agent[bot] 6d2eb2240a fix(app): finish tool call ID rename (#40539)
Co-authored-by: Aiden Cline <rekram1-node@users.noreply.github.com>
2026-08-04 19:26:45 -05:00
16 changed files with 156 additions and 38 deletions
@@ -303,6 +303,7 @@ type AnthropicEvent = Schema.Schema.Type<typeof AnthropicEvent>
interface ParserState {
readonly tools: ToolStream.State<number>
readonly hasLocalToolCalls: boolean
readonly reasoningSignatures: Readonly<Record<number, string>>
readonly usage?: Usage
readonly pendingFinish?: {
@@ -668,7 +669,12 @@ const fromRequest = Effect.fn("AnthropicMessages.fromRequest")(function* (reques
// =============================================================================
// Stream Parsing
// =============================================================================
const mapFinishReason = (reason: string | null | undefined): FinishReason => {
const mapFinishReason = (reason: string | null | undefined, hasLocalToolCalls: boolean): FinishReason => {
if (
hasLocalToolCalls &&
(reason === undefined || reason === null || reason === "end_turn" || reason === "stop_sequence")
)
return "tool-calls"
if (reason === "end_turn" || reason === "stop_sequence" || reason === "pause_turn") return "stop"
if (reason === "max_tokens" || reason === "model_context_window_exceeded") return "length"
if (reason === "tool_use") return "tool-calls"
@@ -931,7 +937,17 @@ const onContentBlockStop = Effect.fn("AnthropicMessages.onContentBlockStop")(fun
events.push(...resultEvents)
const reasoningSignatures = { ...state.reasoningSignatures }
delete reasoningSignatures[event.index]
return [{ ...state, lifecycle, tools: result.tools, reasoningSignatures }, events] satisfies StepResult
return [
{
...state,
hasLocalToolCalls:
state.hasLocalToolCalls || resultEvents.some((item) => item.type === "tool-call" && !item.providerExecuted),
lifecycle,
tools: result.tools,
reasoningSignatures,
},
events,
] satisfies StepResult
})
const onMessageDelta = (state: ParserState, event: AnthropicEvent): StepResult => {
@@ -942,7 +958,7 @@ const onMessageDelta = (state: ParserState, event: AnthropicEvent): StepResult =
usage,
pendingFinish: {
reason: {
normalized: mapFinishReason(event.delta?.stop_reason),
normalized: mapFinishReason(event.delta?.stop_reason, state.hasLocalToolCalls),
raw: event.delta?.stop_reason ?? undefined,
},
providerMetadata:
@@ -1013,6 +1029,7 @@ export const protocol = Protocol.make({
event: Protocol.jsonEvent(AnthropicEvent),
initial: () => ({
tools: ToolStream.empty<number>(),
hasLocalToolCalls: false,
reasoningSignatures: {},
lifecycle: Lifecycle.initial(),
}),
+14 -5
View File
@@ -20,6 +20,7 @@ import {
LanguageModel,
LanguageModelLimits,
LLMEvent,
InvalidProviderOutputReason,
ProviderID,
mergeGenerationOptions,
mergeHttpOptions,
@@ -231,6 +232,17 @@ const streamError = (route: string, message: string, cause: Cause.Cause<unknown>
return ProviderShared.eventError(route, message, Cause.pretty(cause))
}
const incompleteStreamError = (route: string) =>
new AIError({
module: "LLMClient",
method: "stream",
reason: new InvalidProviderOutputReason({
classification: "incomplete-stream",
message: "The provider response ended unexpectedly.",
route,
}),
})
const requireTerminalEvent = (route: string) => (events: Stream.Stream<LLMEvent, AIError>) =>
Stream.suspend(() => {
let terminal = false
@@ -247,7 +259,7 @@ const requireTerminalEvent = (route: string) => (events: Stream.Stream<LLMEvent,
Effect.suspend(() =>
terminal
? Effect.void
: Effect.fail(ProviderShared.eventError(route, "Provider stream ended without a terminal finish event")),
: Effect.fail(incompleteStreamError(route)),
),
),
)
@@ -416,10 +428,7 @@ const generateWith = (stream: Interface["stream"]) =>
const state = yield* stream(request, options).pipe(Stream.runFold(LLMResponse.empty, LLMResponse.reduce))
const response = LLMResponse.complete(state)
if (response) return response
return yield* ProviderShared.eventError(
`${request.model.provider}/${request.model.route.id}`,
"Provider stream ended without a terminal finish event",
)
return yield* incompleteStreamError(`${request.model.provider}/${request.model.route.id}`)
})
export function stream(request: LLMRequest, options?: StreamOptions): Stream.Stream<LLMEvent, AIError, Service> {
+1
View File
@@ -105,6 +105,7 @@ export class InvalidProviderOutputReason extends Schema.Class<InvalidProviderOut
)({
_tag: Schema.tag("InvalidProviderOutput"),
message: Schema.String,
classification: Schema.optional(Schema.Literals(["incomplete-stream"])),
route: Schema.optional(Schema.String),
raw: Schema.optional(Schema.String),
providerMetadata: Schema.optional(ProviderMetadata),
+2 -2
View File
@@ -133,8 +133,8 @@ describe("llm route", () => {
Effect.gen(function* () {
const error = yield* (yield* LLMClient.Service).stream(request).pipe(Stream.runDrain, Effect.flip)
expect(error.reason).toMatchObject({ _tag: "InvalidProviderOutput" })
expect(error.message).toContain("Provider stream ended without a terminal finish event")
expect(error.reason).toMatchObject({ _tag: "InvalidProviderOutput", classification: "incomplete-stream" })
expect(error.message).toContain("The provider response ended unexpectedly.")
}),
)
@@ -538,7 +538,8 @@ describe("Anthropic Messages route", () => {
expect(error.reason).toMatchObject({
_tag: "InvalidProviderOutput",
message: "Provider stream ended without a terminal finish event",
classification: "incomplete-stream",
message: "The provider response ended unexpectedly.",
})
}),
)
@@ -914,6 +915,30 @@ describe("Anthropic Messages route", () => {
}),
)
it.effect("maps a local tool call with end_turn as tool-calls", () =>
Effect.gen(function* () {
const body = sseEvents(
{ type: "message_start", message: { usage: { input_tokens: 5 } } },
{ type: "content_block_start", index: 0, content_block: { type: "tool_use", id: "call_1", name: "lookup" } },
{
type: "content_block_delta",
index: 0,
delta: { type: "input_json_delta", partial_json: '{"query":"weather"}' },
},
{ type: "content_block_stop", index: 0 },
{ type: "message_delta", delta: { stop_reason: "end_turn" }, usage: { output_tokens: 1 } },
{ type: "message_stop" },
)
const response = yield* LLMClient.generate(
LLMRequest.update(request, {
tools: [ToolDefinition.make({ name: "lookup", description: "Lookup data", inputSchema: { type: "object" } })],
}),
).pipe(Effect.provide(fixedResponse(body)))
expect(response.finishReason).toEqual({ normalized: "tool-calls", raw: "end_turn" })
}),
)
it.effect("keeps malformed server tool input terminal", () =>
Effect.gen(function* () {
const body = sseEvents(
@@ -1136,9 +1136,12 @@ describe("OpenAI Chat route", () => {
{ type: "tool-input-delta", id: "call_1", name: "lookup", text: ':"weather"}' },
])
expect(events.filter(LLMEvent.is.toolCall)).toEqual([])
expect(streamError.reason).toMatchObject({ _tag: "InvalidProviderOutput" })
expect(streamError.message).toContain("Provider stream ended without a terminal finish event")
expect(error.message).toContain("Provider stream ended without a terminal finish event")
expect(streamError.reason).toMatchObject({
_tag: "InvalidProviderOutput",
classification: "incomplete-stream",
})
expect(streamError.message).toContain("The provider response ended unexpectedly.")
expect(error.message).toContain("The provider response ended unexpectedly.")
}),
)
@@ -53,7 +53,7 @@ describe("normalizePermissionRequest", () => {
resources: ["README.md"],
save: ["*.md"],
metadata: { path: "README.md" },
source: { type: "tool", messageID: "message-1", callID: "call-1" },
source: { type: "tool", messageID: "message-1", id: "call-1" },
}),
).toEqual({
id: "permission-1",
@@ -48,7 +48,7 @@ export function normalizePermissionRequest(input: PermissionRequest | LegacyPerm
always: input.save ?? [],
metadata: input.metadata ?? {},
tool:
input.source?.type === "tool" ? { messageID: input.source.messageID, callID: input.source.callID } : undefined,
input.source?.type === "tool" ? { messageID: input.source.messageID, callID: input.source.id } : undefined,
}
}
+34 -2
View File
@@ -21,12 +21,24 @@ describe("adaptServerEvent", () => {
id: "evt_1",
created: 1,
type: "permission.asked",
data: { id: "perm_1", sessionID: "ses_1", action: "read", resources: ["src/**"] },
data: {
id: "perm_1",
sessionID: "ses_1",
action: "read",
resources: ["src/**"],
source: { type: "tool", messageID: "msg_1", id: "call_1" },
},
} as OpenCodeEvent
expect(adaptServerEvent(current)).toMatchObject({
type: "permission.asked",
properties: { id: "perm_1", sessionID: "ses_1", permission: "read", patterns: ["src/**"] },
properties: {
id: "perm_1",
sessionID: "ses_1",
permission: "read",
patterns: ["src/**"],
tool: { messageID: "msg_1", callID: "call_1" },
},
current,
})
})
@@ -70,6 +82,26 @@ describe("coalesceServerEvents", () => {
expect(result[0]?.payload.current).toMatchObject({ id: "evt_2", data: { delta: "hello world" } })
})
test("coalesces current tool input deltas by tool ID", () => {
const current = (eventID: string, id: string, delta: string) =>
adaptServerEvent({
id: eventID,
created: 1,
type: "session.tool.input.delta",
location: { directory: "/repo" },
data: { sessionID: "ses", assistantMessageID: "msg", id, delta },
} as OpenCodeEvent)
const result = coalesceServerEvents([
{ directory: "/repo", payload: current("evt_1", "call_1", "{") },
{ directory: "/repo", payload: current("evt_2", "call_1", "}") },
{ directory: "/repo", payload: current("evt_3", "call_2", "[]") },
])
expect(result).toHaveLength(2)
expect(result[0]?.payload.current).toMatchObject({ id: "evt_2", data: { id: "call_1", delta: "{}" } })
expect(result[1]?.payload.current).toMatchObject({ id: "evt_3", data: { id: "call_2", delta: "[]" } })
})
test("preserves event boundaries and distinct fields", () => {
const status = {
directory: "/repo",
+2 -2
View File
@@ -39,7 +39,7 @@ export function adaptServerEvent(event: OpenCodeEvent): ServerEvent {
metadata: event.data.metadata ?? {},
tool:
event.data.source?.type === "tool"
? { messageID: event.data.source.messageID, callID: event.data.source.callID }
? { messageID: event.data.source.messageID, callID: event.data.source.id }
: undefined,
},
current: event,
@@ -142,7 +142,7 @@ function currentDelta(event: OpenCodeEvent | undefined): CurrentDelta | undefine
function currentDeltaKey(event: CurrentDelta) {
if (event.type === "session.tool.input.delta")
return `${event.type}:${event.data.sessionID}:${event.data.assistantMessageID}:${event.data.callID}`
return `${event.type}:${event.data.sessionID}:${event.data.assistantMessageID}:${event.data.id}`
if (event.type === "session.compaction.delta") return `${event.type}:${event.data.sessionID}`
return `${event.type}:${event.data.sessionID}:${event.data.assistantMessageID}:${event.data.ordinal}`
}
@@ -92,19 +92,19 @@ describe("v2 session reducer", () => {
...base,
id: "evt_tool_start",
type: "session.tool.input.started",
data: { sessionID: "ses_1", assistantMessageID: "msg_assistant", callID: "call_1", name: "bash" },
data: { sessionID: "ses_1", assistantMessageID: "msg_assistant", id: "call_1", name: "bash" },
})
apply({
...base,
id: "evt_tool_delta",
type: "session.tool.input.delta",
data: { sessionID: "ses_1", assistantMessageID: "msg_assistant", callID: "call_1", delta: "{}" },
data: { sessionID: "ses_1", assistantMessageID: "msg_assistant", id: "call_1", delta: "{}" },
})
apply({
...base,
id: "evt_tool_called",
type: "session.tool.called",
data: { sessionID: "ses_1", assistantMessageID: "msg_assistant", callID: "call_1", input: {}, executed: true },
data: { sessionID: "ses_1", assistantMessageID: "msg_assistant", id: "call_1", input: {}, executed: true },
})
apply({
...base,
@@ -113,7 +113,7 @@ describe("v2 session reducer", () => {
data: {
sessionID: "ses_1",
assistantMessageID: "msg_assistant",
callID: "call_1",
id: "call_1",
metadata: {},
content: [{ type: "text", text: "done" }],
executed: true,
@@ -241,13 +241,13 @@ export function createV2SessionReducer() {
case "session.tool.input.started":
return updateAssistant(source, event.data.assistantMessageID, sessionID, (item) => ({
...item,
content: item.content.some((content) => content.type === "tool" && content.id === event.data.callID)
content: item.content.some((content) => content.type === "tool" && content.id === event.data.id)
? item.content
: [
...item.content,
{
type: "tool",
id: event.data.callID,
id: event.data.id,
name: event.data.name,
state: { status: "streaming", input: "" },
time: { created: event.created },
@@ -255,17 +255,17 @@ export function createV2SessionReducer() {
],
}))
case "session.tool.input.delta":
return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) =>
return updateTool(source, event.data.assistantMessageID, event.data.id, sessionID, (tool) =>
tool.state.status === "streaming"
? { ...tool, state: { ...tool.state, input: tool.state.input + event.data.delta } }
: tool,
)
case "session.tool.input.ended":
return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) =>
return updateTool(source, event.data.assistantMessageID, event.data.id, sessionID, (tool) =>
tool.state.status === "streaming" ? { ...tool, state: { ...tool.state, input: event.data.text } } : tool,
)
case "session.tool.called":
return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) => ({
return updateTool(source, event.data.assistantMessageID, event.data.id, sessionID, (tool) => ({
...tool,
executed: event.data.executed,
providerState: event.data.state,
@@ -274,7 +274,7 @@ export function createV2SessionReducer() {
time: { ...tool.time, ran: event.created },
}))
case "session.tool.progress":
return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) =>
return updateTool(source, event.data.assistantMessageID, event.data.id, sessionID, (tool) =>
tool.state.status === "running"
? {
...tool,
@@ -284,7 +284,7 @@ export function createV2SessionReducer() {
: tool,
)
case "session.tool.success":
return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) => {
return updateTool(source, event.data.assistantMessageID, event.data.id, sessionID, (tool) => {
if (tool.state.status !== "running") return tool
return {
...tool,
@@ -302,7 +302,7 @@ export function createV2SessionReducer() {
}
})
case "session.tool.failed":
return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) => {
return updateTool(source, event.data.assistantMessageID, event.data.id, sessionID, (tool) => {
if (tool.state.status !== "streaming" && tool.state.status !== "running") return tool
return {
...tool,
+1 -1
View File
@@ -470,7 +470,7 @@ export async function runNonInteractivePrompt(input: Input) {
if (event.type === "session.step.failed") {
if (
input.compatibility === "v1" &&
event.data.error.message === "Provider stream ended without a terminal finish event"
event.data.error.message === "The provider response ended unexpectedly."
) {
pendingStep = undefined
v1InvalidOutput = true
+2 -2
View File
@@ -503,8 +503,8 @@ describe("runNonInteractivePrompt", () => {
turn: (messageID) => [
prompted(messageID),
stepStarted(),
stepFailed("Provider stream ended without a terminal finish event"),
executionFailed("Provider stream ended without a terminal finish event"),
stepFailed("The provider response ended unexpectedly."),
executionFailed("The provider response ended unexpectedly."),
],
})
+2 -1
View File
@@ -20,10 +20,11 @@ export function isRetryable(error: AIError) {
case "ProviderInternal":
case "Transport":
return true
case "InvalidProviderOutput":
return error.reason.classification === "incomplete-stream"
case "Authentication":
case "QuotaExceeded":
case "ContentPolicy":
case "InvalidProviderOutput":
case "InvalidRequest":
case "NoRoute":
case "UnknownProvider":
+32 -2
View File
@@ -513,6 +513,16 @@ const providerUnavailable = () =>
reason: new TransportReason({ message: "Provider unavailable" }),
})
const incompleteStream = () =>
new AIError({
module: "test",
method: "stream",
reason: new InvalidProviderOutputReason({
classification: "incomplete-stream",
message: "The provider response ended unexpectedly.",
}),
})
const invalidRequest = () =>
new AIError({
module: "test",
@@ -3949,6 +3959,26 @@ describe("SessionRunnerLLM", () => {
}),
)
it.effect("retries an incomplete stream before output", () =>
Effect.gen(function* () {
const session = yield* setup
yield* admit(session, "Retry incomplete stream")
yield* TestLLM.push(Stream.fail(incompleteStream()))
yield* TestLLM.push(TestLLM.text("Recovered", "incomplete-stream-success"))
const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
yield* TestLLM.wait(1)
yield* TestClock.adjust("2 seconds")
yield* Fiber.join(run)
expect(requests).toHaveLength(2)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user" },
{ type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] },
])
}),
)
it.effect("uses a larger provider retry-after delay", () =>
Effect.gen(function* () {
const session = yield* setup
@@ -3969,7 +3999,7 @@ describe("SessionRunnerLLM", () => {
it.effect("does not retry eligible failures after observable output", () =>
Effect.gen(function* () {
const session = yield* setup
const failure = rateLimited()
const failure = incompleteStream()
yield* TestLLM.push(
TestLLM.failAfter(
failure,
@@ -3987,7 +4017,7 @@ describe("SessionRunnerLLM", () => {
{
type: "assistant",
finish: "error",
error: { type: "provider.rate-limit" },
error: { type: "provider.invalid-output" },
content: [{ type: "text", text: "Partial" }],
},
])