mirror of
https://github.com/anomalyco/opencode.git
synced 2026-10-04 22:46:17 +00:00
Compare commits
1
Commits
v2
...
tool-stream-sync
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d175644ce1 |
No files matched your search
@@ -1300,32 +1300,29 @@ const onContentBlockStart = (
|
||||
return [{ ...state, lifecycle: Lifecycle.stepStart(state.lifecycle, events) }, [...events, result]]
|
||||
}
|
||||
|
||||
const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(function* (
|
||||
const onContentBlockDelta = (
|
||||
state: ParserState,
|
||||
event: AnthropicEvent & { readonly delta: AnthropicStreamDelta },
|
||||
) {
|
||||
): StepResult | AIError => {
|
||||
const delta = event.delta
|
||||
|
||||
if (delta.type === "compaction_delta") {
|
||||
if (event.index === undefined || !(event.index in state.compactions) || delta.content === undefined)
|
||||
return yield* ProviderShared.eventError(ADAPTER, "Compaction delta is missing its block or content")
|
||||
return [
|
||||
{ ...state, compactions: { ...state.compactions, [event.index]: delta.content } },
|
||||
NO_EVENTS,
|
||||
] satisfies StepResult
|
||||
return ProviderShared.eventError(ADAPTER, "Compaction delta is missing its block or content")
|
||||
return [{ ...state, compactions: { ...state.compactions, [event.index]: delta.content } }, NO_EVENTS]
|
||||
}
|
||||
|
||||
if (delta.type === "text_delta" && delta.text) {
|
||||
if (!state.lifecycle.text.has(`text-${event.index ?? 0}`)) return [state, NO_EVENTS] satisfies StepResult
|
||||
if (!state.lifecycle.text.has(`text-${event.index ?? 0}`)) return [state, NO_EVENTS]
|
||||
const events: LLMEvent[] = []
|
||||
return [
|
||||
{ ...state, lifecycle: Lifecycle.textDelta(state.lifecycle, events, `text-${event.index ?? 0}`, delta.text) },
|
||||
events,
|
||||
] satisfies StepResult
|
||||
]
|
||||
}
|
||||
|
||||
if (delta.type === "thinking_delta" && delta.thinking) {
|
||||
if (!state.lifecycle.reasoning.has(`reasoning-${event.index ?? 0}`)) return [state, NO_EVENTS] satisfies StepResult
|
||||
if (!state.lifecycle.reasoning.has(`reasoning-${event.index ?? 0}`)) return [state, NO_EVENTS]
|
||||
const events: LLMEvent[] = []
|
||||
return [
|
||||
{
|
||||
@@ -1333,24 +1330,24 @@ const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(f
|
||||
lifecycle: Lifecycle.reasoningDelta(state.lifecycle, events, `reasoning-${event.index ?? 0}`, delta.thinking),
|
||||
},
|
||||
events,
|
||||
] satisfies StepResult
|
||||
]
|
||||
}
|
||||
|
||||
if (delta.type === "signature_delta" && delta.signature) {
|
||||
const index = event.index ?? 0
|
||||
if (!state.lifecycle.reasoning.has(`reasoning-${index}`)) return [state, NO_EVENTS] satisfies StepResult
|
||||
if (!state.lifecycle.reasoning.has(`reasoning-${index}`)) return [state, NO_EVENTS]
|
||||
return [
|
||||
{
|
||||
...state,
|
||||
reasoningSignatures: { ...state.reasoningSignatures, [index]: delta.signature },
|
||||
},
|
||||
NO_EVENTS,
|
||||
] satisfies StepResult
|
||||
]
|
||||
}
|
||||
|
||||
if (delta.type === "input_json_delta" && event.index !== undefined) {
|
||||
if (!delta.partial_json) return [state, NO_EVENTS] satisfies StepResult
|
||||
if (!state.tools[event.index]) return [state, NO_EVENTS] satisfies StepResult
|
||||
if (!delta.partial_json) return [state, NO_EVENTS]
|
||||
if (!state.tools[event.index]) return [state, NO_EVENTS]
|
||||
const result = ToolStream.appendExisting(
|
||||
ADAPTER,
|
||||
state.tools,
|
||||
@@ -1358,21 +1355,18 @@ const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(f
|
||||
delta.partial_json,
|
||||
"Anthropic Messages tool argument delta is missing its tool call",
|
||||
)
|
||||
if (ToolStream.isError(result)) return yield* result
|
||||
if (ToolStream.isError(result)) return result
|
||||
const events: LLMEvent[] = []
|
||||
const lifecycle = result.events.length ? Lifecycle.stepStart(state.lifecycle, events) : state.lifecycle
|
||||
events.push(...result.events)
|
||||
return [{ ...state, lifecycle, tools: result.tools }, events] satisfies StepResult
|
||||
return [{ ...state, lifecycle, tools: result.tools }, events]
|
||||
}
|
||||
|
||||
return [state, NO_EVENTS] satisfies StepResult
|
||||
})
|
||||
return [state, NO_EVENTS]
|
||||
}
|
||||
|
||||
const onContentBlockStop = Effect.fn("AnthropicMessages.onContentBlockStop")(function* (
|
||||
state: ParserState,
|
||||
event: AnthropicEvent,
|
||||
) {
|
||||
if (event.index === undefined) return [state, NO_EVENTS] satisfies StepResult
|
||||
const onContentBlockStop = (state: ParserState, event: AnthropicEvent): StepResult | AIError => {
|
||||
if (event.index === undefined) return [state, NO_EVENTS]
|
||||
if (event.index in state.compactions) {
|
||||
const { [event.index]: content, ...compactions } = state.compactions
|
||||
const events: LLMEvent[] = []
|
||||
@@ -1383,9 +1377,10 @@ const onContentBlockStop = Effect.fn("AnthropicMessages.onContentBlockStop")(fun
|
||||
text: content,
|
||||
}),
|
||||
)
|
||||
return [{ ...state, compactions, lifecycle }, events] satisfies StepResult
|
||||
return [{ ...state, compactions, lifecycle }, events]
|
||||
}
|
||||
const result = yield* ToolStream.finish(ADAPTER, state.tools, event.index)
|
||||
const result = ToolStream.finish(ADAPTER, state.tools, event.index)
|
||||
if (ToolStream.isError(result)) return result
|
||||
const events: LLMEvent[] = []
|
||||
const resultEvents = result.events ?? []
|
||||
const signature = state.reasoningSignatures[event.index]
|
||||
@@ -1400,8 +1395,8 @@ 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, lifecycle, tools: result.tools, reasoningSignatures }, events]
|
||||
}
|
||||
|
||||
const onMessageDelta = (
|
||||
state: ParserState,
|
||||
@@ -1439,10 +1434,11 @@ const onMessageDelta = (
|
||||
]
|
||||
}
|
||||
|
||||
const onMessageStop = Effect.fn("AnthropicMessages.onMessageStop")(function* (state: ParserState) {
|
||||
const onMessageStop = (state: ParserState): StepResult | AIError => {
|
||||
if (Object.keys(state.compactions).length)
|
||||
return yield* ProviderShared.eventError(ADAPTER, "Response ended with an incomplete compaction block")
|
||||
const result = yield* ToolStream.finishAll(ADAPTER, state.tools)
|
||||
return ProviderShared.eventError(ADAPTER, "Response ended with an incomplete compaction block")
|
||||
const result = ToolStream.finishAll(ADAPTER, state.tools)
|
||||
if (ToolStream.isError(result)) return result
|
||||
const events: LLMEvent[] = []
|
||||
const lifecycle = result.events.length ? Lifecycle.stepStart(state.lifecycle, events) : state.lifecycle
|
||||
events.push(...result.events)
|
||||
@@ -1464,8 +1460,8 @@ const onMessageStop = Effect.fn("AnthropicMessages.onMessageStop")(function* (st
|
||||
usage: state.usage,
|
||||
providerMetadata: state.pendingFinish?.providerMetadata,
|
||||
})
|
||||
return [{ ...state, lifecycle: finished, tools: result.tools }, events] satisfies StepResult
|
||||
})
|
||||
return [{ ...state, lifecycle: finished, tools: result.tools }, events]
|
||||
}
|
||||
|
||||
// Prefix `error.type` so overloads, rate limits, and quota errors are visible
|
||||
// even when the provider message is generic or empty.
|
||||
@@ -1511,6 +1507,9 @@ const invalidStreamEvent = (event: AnthropicEvent) =>
|
||||
),
|
||||
)
|
||||
|
||||
const stepOutcome = (result: StepResult | AIError) =>
|
||||
result instanceof AIError ? Effect.fail(result) : Effect.succeed(result)
|
||||
|
||||
const step = (state: ParserState, event: AnthropicEvent) => {
|
||||
if (!SSE_EVENTS.has(event.type)) return Effect.succeed<StepResult>([state, NO_EVENTS])
|
||||
if (
|
||||
@@ -1556,15 +1555,15 @@ const step = (state: ParserState, event: AnthropicEvent) => {
|
||||
return Effect.succeed<StepResult>([state, NO_EVENTS])
|
||||
const decoded = decodeAnthropicStreamDelta(event.delta)
|
||||
if (Option.isNone(decoded)) return invalidStreamEvent(event)
|
||||
return onContentBlockDelta(state, { ...event, delta: decoded.value })
|
||||
return stepOutcome(onContentBlockDelta(state, { ...event, delta: decoded.value }))
|
||||
}
|
||||
if (event.type === "content_block_stop") return onContentBlockStop(state, event)
|
||||
if (event.type === "content_block_stop") return stepOutcome(onContentBlockStop(state, event))
|
||||
if (event.type === "message_delta") {
|
||||
const decoded = decodeAnthropicStreamDelta(event.delta)
|
||||
if (Option.isNone(decoded)) return invalidStreamEvent(event)
|
||||
return Effect.succeed(onMessageDelta(state, { ...event, delta: decoded.value }))
|
||||
}
|
||||
if (event.type === "message_stop") return onMessageStop(state)
|
||||
if (event.type === "message_stop") return stepOutcome(onMessageStop(state))
|
||||
if (event.type === "error") return onError(event)
|
||||
return Effect.succeed<StepResult>([state, NO_EVENTS])
|
||||
}
|
||||
|
||||
@@ -722,7 +722,8 @@ const step = (state: ParserState, event: BedrockEvent) =>
|
||||
|
||||
if (event.contentBlockStop) {
|
||||
const index = event.contentBlockStop.contentBlockIndex
|
||||
const result = yield* ToolStream.finish(ADAPTER, state.tools, index)
|
||||
const result = ToolStream.finish(ADAPTER, state.tools, index)
|
||||
if (ToolStream.isError(result)) return yield* result
|
||||
const events: LLMEvent[] = []
|
||||
const resultEvents = result.events ?? []
|
||||
const lifecycle = (() => {
|
||||
|
||||
@@ -252,7 +252,7 @@ const mapUsage = (usage: typeof NativeUsage.Type) =>
|
||||
})
|
||||
|
||||
// Lifecycle deltas open blocks on demand and ends are no-ops for closed blocks, so content-start needs no handling.
|
||||
const step = Effect.fn("CohereChat.step")(function* (state: State, event: Event) {
|
||||
const step = Effect.fnUntraced(function* (state: State, event: Event) {
|
||||
const events: LLMEvent[] = []
|
||||
switch (event.type) {
|
||||
case "message-start":
|
||||
@@ -292,11 +292,13 @@ const step = Effect.fn("CohereChat.step")(function* (state: State, event: Event)
|
||||
return [{ ...state, tools: result.tools }, result.events] as const
|
||||
}
|
||||
case "tool-call-end": {
|
||||
const result = yield* ToolStream.finish(ADAPTER, state.tools, event.index)
|
||||
const result = ToolStream.finish(ADAPTER, state.tools, event.index)
|
||||
if (ToolStream.isError(result)) return yield* result
|
||||
return [{ ...state, tools: result.tools }, result.events ?? []] as const
|
||||
}
|
||||
case "message-end": {
|
||||
const pending = yield* ToolStream.finishAll(ADAPTER, state.tools)
|
||||
const pending = ToolStream.finishAll(ADAPTER, state.tools)
|
||||
if (ToolStream.isError(pending)) return yield* pending
|
||||
events.push(...pending.events)
|
||||
const lifecycle = Lifecycle.finish(state.lifecycle, events, {
|
||||
reason: finishReason(event.delta.finish_reason),
|
||||
|
||||
@@ -364,7 +364,7 @@ const mapUsage = (usage: RawUsage | undefined, key: string) => {
|
||||
})
|
||||
}
|
||||
|
||||
const onStart = Effect.fn("GoogleInteractions.onStart")(function* (
|
||||
const onStart = Effect.fnUntraced(function* (
|
||||
state: ParserState,
|
||||
index: number,
|
||||
step: OutputStep,
|
||||
@@ -402,7 +402,7 @@ const onStart = Effect.fn("GoogleInteractions.onStart")(function* (
|
||||
return [{ ...state, lifecycle, tools, steps: { ...state.steps, [index]: step } }, events] satisfies StepResult
|
||||
})
|
||||
|
||||
const onDelta = Effect.fn("GoogleInteractions.onDelta")(function* (
|
||||
const onDelta = Effect.fnUntraced(function* (
|
||||
state: ParserState,
|
||||
index: number,
|
||||
delta: typeof Delta.Type,
|
||||
@@ -447,7 +447,7 @@ const onDelta = Effect.fn("GoogleInteractions.onDelta")(function* (
|
||||
return yield* ProviderShared.eventError(ADAPTER, `Unsupported Interactions delta: ${delta.type}`, encodeJson(delta))
|
||||
})
|
||||
|
||||
const onStop = Effect.fn("GoogleInteractions.onStop")(function* (state: ParserState, index: number) {
|
||||
const onStop = Effect.fnUntraced(function* (state: ParserState, index: number) {
|
||||
const step = state.steps[index]
|
||||
if (!step) return yield* ProviderShared.eventError(ADAPTER, "Interactions step.stop without step.start")
|
||||
const events: LLMEvent[] = []
|
||||
@@ -461,11 +461,12 @@ const onStop = Effect.fn("GoogleInteractions.onStop")(function* (state: ParserSt
|
||||
{ ...state, lifecycle: Lifecycle.textEnd(state.lifecycle, events, String(index)) },
|
||||
events,
|
||||
] satisfies StepResult
|
||||
const result = yield* ToolStream.finish(ADAPTER, state.tools, index)
|
||||
const result = ToolStream.finish(ADAPTER, state.tools, index)
|
||||
if (ToolStream.isError(result)) return yield* result
|
||||
return [{ ...state, tools: result.tools }, result.events ?? []] satisfies StepResult
|
||||
})
|
||||
|
||||
const step = Effect.fn("GoogleInteractions.step")(function* (state: ParserState, event: Event) {
|
||||
const step = Effect.fnUntraced(function* (state: ParserState, event: Event) {
|
||||
switch (event.event_type) {
|
||||
case "step.start":
|
||||
return yield* onStart(state, event.index, event.step)
|
||||
@@ -500,7 +501,8 @@ const step = Effect.fn("GoogleInteractions.step")(function* (state: ParserState,
|
||||
`Unexpected terminal Interactions status: ${interaction.status}`,
|
||||
encodeJson(event),
|
||||
)
|
||||
const pending = yield* ToolStream.finishAll(ADAPTER, state.tools)
|
||||
const pending = ToolStream.finishAll(ADAPTER, state.tools)
|
||||
if (ToolStream.isError(pending)) return yield* pending
|
||||
const events = [...pending.events]
|
||||
const lifecycle = Lifecycle.finish(state.lifecycle, events, {
|
||||
reason: {
|
||||
|
||||
@@ -95,7 +95,7 @@ const HOSTED_TOOLS = {
|
||||
image_generation_call: {
|
||||
name: "image_generation",
|
||||
input: () => ({}),
|
||||
result: Effect.fn("MetaResponses.imageResult")(function* (raw: ResponsesHostedTools.Item) {
|
||||
result: Effect.fnUntraced(function* (raw: ResponsesHostedTools.Item) {
|
||||
const item = yield* Schema.decodeUnknownEffect(ImageItem)(raw).pipe(
|
||||
Effect.mapError((cause) =>
|
||||
ProviderShared.eventError(
|
||||
@@ -136,7 +136,7 @@ const HOSTED_TOOLS = {
|
||||
},
|
||||
} satisfies ResponsesHostedTools.Definitions
|
||||
|
||||
const onEvent = Effect.fn("MetaResponses.onEvent")(function* (
|
||||
const onEvent = Effect.fnUntraced(function* (
|
||||
state: OpenResponses.ParserState,
|
||||
input: OpenResponses.Event,
|
||||
) {
|
||||
@@ -173,7 +173,7 @@ const onEvent = Effect.fn("MetaResponses.onEvent")(function* (
|
||||
] satisfies OpenResponses.StepResult
|
||||
})
|
||||
|
||||
const step = Effect.fn("MetaResponses.step")(function* (state: ParserState, input: OpenResponses.Event) {
|
||||
const step = Effect.fnUntraced(function* (state: ParserState, input: OpenResponses.Event) {
|
||||
const completedItems = new Set(state.completedItems)
|
||||
const event = OpenResponses.normalize(state, input)
|
||||
if (event.type === "response.output_item.done" && event.item && completedItems.has(event.item.id))
|
||||
|
||||
@@ -596,7 +596,7 @@ const toolText = (tool: MistralToolDelta) => {
|
||||
return value === null || value === undefined ? "" : ProviderShared.encodeJson(value)
|
||||
}
|
||||
|
||||
const appendTools = Effect.fn("MistralChat.appendTools")(function* (
|
||||
const appendTools = Effect.fnUntraced(function* (
|
||||
initial: ParserState,
|
||||
events: LLMEvent[],
|
||||
deltas: ReadonlyArray<MistralToolDelta>,
|
||||
@@ -662,7 +662,7 @@ const hasLateContent = (event: MistralEvent) => {
|
||||
)
|
||||
}
|
||||
|
||||
const step = Effect.fn("MistralChat.step")(function* (state: ParserState, event: MistralEvent) {
|
||||
const step = Effect.fnUntraced(function* (state: ParserState, event: MistralEvent) {
|
||||
if (event.error) {
|
||||
const body = ProviderShared.encodeJson(event)
|
||||
return yield* new AIError({
|
||||
@@ -712,8 +712,9 @@ const step = Effect.fn("MistralChat.step")(function* (state: ParserState, event:
|
||||
)
|
||||
const finished =
|
||||
!incomplete && Object.keys(withTools.tools).length > 0
|
||||
? yield* ToolStream.finishAll(ADAPTER, withTools.tools)
|
||||
? ToolStream.finishAll(ADAPTER, withTools.tools)
|
||||
: undefined
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
return [
|
||||
{
|
||||
...withTools,
|
||||
@@ -726,7 +727,7 @@ const step = Effect.fn("MistralChat.step")(function* (state: ParserState, event:
|
||||
] as const
|
||||
})
|
||||
|
||||
const finishEvents = Effect.fn("MistralChat.finishEvents")(function* (state: ParserState) {
|
||||
const finishEvents = Effect.fnUntraced(function* (state: ParserState) {
|
||||
if (!state.finishReason)
|
||||
return yield* new AIError({
|
||||
reason: new InvalidProviderOutputError({
|
||||
|
||||
@@ -1143,23 +1143,16 @@ const onReasoningSummaryPartDone = (state: ParserState, event: Event): StepResul
|
||||
]
|
||||
}
|
||||
|
||||
const onFunctionCallArgumentsDelta = Effect.fn("OpenResponses.onFunctionCallArgumentsDelta")(function* (
|
||||
state: ParserState,
|
||||
event: Event,
|
||||
) {
|
||||
if (event.item_id === undefined) return [state, NO_EVENTS] satisfies StepResult
|
||||
const onFunctionCallArgumentsDelta = (state: ParserState, event: Event): StepResult | AIError => {
|
||||
if (event.item_id === undefined) return [state, NO_EVENTS]
|
||||
const tool = state.tools[event.item_id]
|
||||
if (!tool) return [state, NO_EVENTS] satisfies StepResult
|
||||
if (!tool) return [state, NO_EVENTS]
|
||||
const final = event.type === "response.function_call_arguments.done" ? event.arguments : undefined
|
||||
if (event.type === "response.function_call_arguments.done" && final === undefined)
|
||||
return [state, NO_EVENTS] satisfies StepResult
|
||||
if (event.type === "response.function_call_arguments.done" && final === undefined) return [state, NO_EVENTS]
|
||||
if (final !== undefined && !final.startsWith(tool.input))
|
||||
return [
|
||||
{ ...state, tools: ToolStream.start(state.tools, event.item_id, { ...tool, input: final }) },
|
||||
NO_EVENTS,
|
||||
] satisfies StepResult
|
||||
return [{ ...state, tools: ToolStream.start(state.tools, event.item_id, { ...tool, input: final }) }, NO_EVENTS]
|
||||
const delta = final === undefined ? event.delta : final.slice(tool.input.length)
|
||||
if (!delta) return [state, NO_EVENTS] satisfies StepResult
|
||||
if (!delta) return [state, NO_EVENTS]
|
||||
const result = ToolStream.appendExisting(
|
||||
state.id,
|
||||
state.tools,
|
||||
@@ -1167,23 +1160,20 @@ const onFunctionCallArgumentsDelta = Effect.fn("OpenResponses.onFunctionCallArgu
|
||||
delta,
|
||||
`${state.name} tool argument delta is missing its tool call`,
|
||||
)
|
||||
if (ToolStream.isError(result)) return yield* result
|
||||
if (ToolStream.isError(result)) return result
|
||||
const events: LLMEvent[] = []
|
||||
const lifecycle = result.events.length ? Lifecycle.stepStart(state.lifecycle, events) : state.lifecycle
|
||||
events.push(...result.events)
|
||||
return [{ ...state, lifecycle, tools: result.tools }, events] satisfies StepResult
|
||||
})
|
||||
return [{ ...state, lifecycle, tools: result.tools }, events]
|
||||
}
|
||||
|
||||
const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
|
||||
state: ParserState,
|
||||
item: NormalizedEvent["item"],
|
||||
) {
|
||||
if (!item) return [state, NO_EVENTS] satisfies StepResult
|
||||
const onOutputItemDone = (state: ParserState, item: NormalizedEvent["item"]): StepResult | AIError => {
|
||||
if (!item) return [state, NO_EVENTS]
|
||||
|
||||
if (item.type === "compaction") {
|
||||
if (typeof item.encrypted_content !== "string")
|
||||
return yield* ProviderShared.eventError(state.id, "Compaction output is missing its encrypted content")
|
||||
if (state.completedCompactions.has(item.id)) return [state, NO_EVENTS] satisfies StepResult
|
||||
return ProviderShared.eventError(state.id, "Compaction output is missing its encrypted content")
|
||||
if (state.completedCompactions.has(item.id)) return [state, NO_EVENTS]
|
||||
const events: LLMEvent[] = []
|
||||
const lifecycle = Lifecycle.stepStart(state.lifecycle, events)
|
||||
events.push(
|
||||
@@ -1193,10 +1183,7 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
|
||||
encrypted: item.encrypted_content,
|
||||
}),
|
||||
)
|
||||
return [
|
||||
{ ...state, lifecycle, completedCompactions: new Set([...state.completedCompactions, item.id]) },
|
||||
events,
|
||||
] satisfies StepResult
|
||||
return [{ ...state, lifecycle, completedCompactions: new Set([...state.completedCompactions, item.id]) }, events]
|
||||
}
|
||||
|
||||
if (item.type === "message") {
|
||||
@@ -1221,11 +1208,11 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
|
||||
message: active ? undefined : state.message,
|
||||
},
|
||||
events,
|
||||
] satisfies StepResult
|
||||
]
|
||||
}
|
||||
|
||||
if (item.type === "function_call") {
|
||||
if (!item.call_id || !item.name) return [state, NO_EVENTS] satisfies StepResult
|
||||
if (!item.call_id || !item.name) return [state, NO_EVENTS]
|
||||
const metadata = providerMetadata(state, { itemId: item.id })
|
||||
const registered = state.tools[item.id] !== undefined
|
||||
const tools = registered
|
||||
@@ -1236,10 +1223,8 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
|
||||
namespace: item.namespace,
|
||||
providerMetadata: metadata,
|
||||
})
|
||||
const result =
|
||||
item.arguments === undefined
|
||||
? yield* ToolStream.finish(state.id, tools, item.id)
|
||||
: yield* ToolStream.finishWithInput(state.id, tools, item.id, item.arguments)
|
||||
const result = ToolStream.finish(state.id, tools, item.id, item.arguments)
|
||||
if (ToolStream.isError(result)) return result
|
||||
const events: LLMEvent[] = []
|
||||
const finished = result.events ?? []
|
||||
// A done-only call never streamed a start event, so open its lifecycle here.
|
||||
@@ -1267,7 +1252,7 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
|
||||
tools: result.tools,
|
||||
},
|
||||
events,
|
||||
] satisfies StepResult
|
||||
]
|
||||
}
|
||||
|
||||
if (item.type === "reasoning") {
|
||||
@@ -1299,18 +1284,18 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
|
||||
}
|
||||
const reasoningItems = { ...state.reasoningItems }
|
||||
delete reasoningItems[item.id]
|
||||
return [{ ...state, lifecycle, reasoningItems }, events] satisfies StepResult
|
||||
return [{ ...state, lifecycle, reasoningItems }, events]
|
||||
}
|
||||
const lifecycle = Lifecycle.stepStart(state.lifecycle, events)
|
||||
events.push(LLMEvent.reasoningStart({ id: item.id, providerMetadata: metadata }))
|
||||
events.push(LLMEvent.reasoningEnd({ id: item.id, providerMetadata: metadata, text: itemText }))
|
||||
return [{ ...state, lifecycle }, events] satisfies StepResult
|
||||
return [{ ...state, lifecycle }, events]
|
||||
}
|
||||
|
||||
return [state, NO_EVENTS] satisfies StepResult
|
||||
})
|
||||
return [state, NO_EVENTS]
|
||||
}
|
||||
|
||||
const onResponseFinish = Effect.fn("OpenResponses.onResponseFinish")(function* (state: ParserState, event: Event) {
|
||||
const onResponseFinish = (state: ParserState, event: Event): StepResult | AIError => {
|
||||
let current = state
|
||||
const events: LLMEvent[] = []
|
||||
if (event.type === "response.completed") {
|
||||
@@ -1318,19 +1303,21 @@ const onResponseFinish = Effect.fn("OpenResponses.onResponseFinish")(function* (
|
||||
for (const item of (event.response?.output ?? []).map((item, index) => resolveItem(state, item, index))) {
|
||||
// Terminal recovery cannot insert a checkpoint before already-emitted content.
|
||||
if (item.type === "compaction" && state.lifecycle.stepStarted && !state.completedCompactions.has(item.id))
|
||||
return yield* ProviderShared.eventError(
|
||||
return ProviderShared.eventError(
|
||||
state.id,
|
||||
"Cannot recover a compaction checkpoint after output has been emitted",
|
||||
)
|
||||
const recoverable =
|
||||
item.type === "compaction" || (item.type === "function_call" && current.tools[item.id] !== undefined)
|
||||
if (!recoverable) continue
|
||||
const [next, emitted] = yield* onOutputItemDone(current, item)
|
||||
current = next
|
||||
events.push(...emitted)
|
||||
const done = onOutputItemDone(current, item)
|
||||
if (done instanceof AIError) return done
|
||||
current = done[0]
|
||||
events.push(...done[1])
|
||||
}
|
||||
// Some compatible providers omit output_item.done even after completing the response.
|
||||
const pending = yield* ToolStream.finishAll(current.id, current.tools)
|
||||
const pending = ToolStream.finishAll(current.id, current.tools)
|
||||
if (ToolStream.isError(pending)) return pending
|
||||
current = {
|
||||
...current,
|
||||
tools: pending.tools,
|
||||
@@ -1354,8 +1341,8 @@ const onResponseFinish = Effect.fn("OpenResponses.onResponseFinish")(function* (
|
||||
})
|
||||
: undefined,
|
||||
})
|
||||
return [{ ...current, lifecycle }, events] satisfies StepResult
|
||||
})
|
||||
return [{ ...current, lifecycle }, events]
|
||||
}
|
||||
|
||||
/** Error code and message from wherever the frame put them; top-level fields win over nested ones. */
|
||||
export const errorDetail = (event: Event) => {
|
||||
@@ -1390,6 +1377,9 @@ export const providerFailure = (event: Event, fallback: string, body = ProviderS
|
||||
return new AIError({ reason })
|
||||
}
|
||||
|
||||
const stepOutcome = (result: StepResult | AIError) =>
|
||||
result instanceof AIError ? Effect.fail(result) : Effect.succeed(result)
|
||||
|
||||
// Callers must pass events through `normalize` first. The OpenAPI requires
|
||||
// string IDs but imposes no minLength; empty is not missing.
|
||||
export const step = (state: ParserState, event: NormalizedEvent) => {
|
||||
@@ -1449,10 +1439,11 @@ export const step = (state: ParserState, event: NormalizedEvent) => {
|
||||
}
|
||||
if (event.type === "response.function_call_arguments.delta" || event.type === "response.function_call_arguments.done")
|
||||
return event.item_id !== undefined
|
||||
? onFunctionCallArgumentsDelta(state, event)
|
||||
? stepOutcome(onFunctionCallArgumentsDelta(state, event))
|
||||
: ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
|
||||
if (event.type === "response.output_item.done") return onOutputItemDone(state, event.item)
|
||||
if (event.type === "response.completed" || event.type === "response.incomplete") return onResponseFinish(state, event)
|
||||
if (event.type === "response.output_item.done") return stepOutcome(onOutputItemDone(state, event.item))
|
||||
if (event.type === "response.completed" || event.type === "response.incomplete")
|
||||
return stepOutcome(onResponseFinish(state, event))
|
||||
if (event.type === "response.failed") return providerFailure(event, `${state.name} response failed`)
|
||||
if (event.type === "error") return providerFailure(event, `${state.name} stream error`)
|
||||
return Effect.succeed<StepResult>([state, NO_EVENTS])
|
||||
|
||||
@@ -855,7 +855,7 @@ export const fromRequest = Effect.fn("OpenAIChat.fromRequest")(function* (
|
||||
// Streaming parsers are small state machines: every event returns a new state
|
||||
// plus the common `LLMEvent`s produced by that event. Tool calls are accumulated
|
||||
// because OpenAI streams JSON arguments across multiple deltas.
|
||||
const mapFinishReason = Effect.fn("OpenAIChat.mapFinishReason")(function* (event: OpenAIChatEvent, reason: string) {
|
||||
const mapFinishReason = Effect.fnUntraced(function* (event: OpenAIChatEvent, reason: string) {
|
||||
switch (reason) {
|
||||
case "error":
|
||||
return yield* new AIError({
|
||||
@@ -1189,8 +1189,9 @@ const step = (state: ParserState, event: OpenAIChatEvent) =>
|
||||
!incompleteTools &&
|
||||
state.finishReason === undefined &&
|
||||
Object.keys(tools).length > 0
|
||||
? yield* ToolStream.finishAll(ADAPTER, tools)
|
||||
? ToolStream.finishAll(ADAPTER, tools)
|
||||
: undefined
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
return [
|
||||
{
|
||||
@@ -1214,7 +1215,7 @@ const step = (state: ParserState, event: OpenAIChatEvent) =>
|
||||
] as const
|
||||
})
|
||||
|
||||
const finishEvents = Effect.fn("OpenAIChat.finishEvents")(function* (state: ParserState) {
|
||||
const finishEvents = Effect.fnUntraced(function* (state: ParserState) {
|
||||
if (state.finishReason === undefined && state.requireFinishReason)
|
||||
return yield* new AIError({
|
||||
reason: new InvalidProviderOutputError({
|
||||
@@ -1224,10 +1225,12 @@ const finishEvents = Effect.fn("OpenAIChat.finishEvents")(function* (state: Pars
|
||||
}),
|
||||
})
|
||||
const events: LLMEvent[] = []
|
||||
const toolCallEvents =
|
||||
const finished =
|
||||
state.finishReason === undefined && Object.keys(state.tools).length > 0
|
||||
? (yield* ToolStream.finishAll(ADAPTER, state.tools)).events
|
||||
: state.toolCallEvents
|
||||
? ToolStream.finishAll(ADAPTER, state.tools)
|
||||
: undefined
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
const toolCallEvents = finished?.events ?? state.toolCallEvents
|
||||
const hasToolCalls = toolCallEvents.length > 0
|
||||
const reason = state.finishReason
|
||||
? {
|
||||
|
||||
@@ -237,7 +237,7 @@ const checkpointBody = {
|
||||
}),
|
||||
}
|
||||
|
||||
const hostedToolResult = Effect.fn("OpenAIResponses.hostedToolResult")(function* (item: ResponsesHostedTools.Item) {
|
||||
const hostedToolResult = Effect.fnUntraced(function* (item: ResponsesHostedTools.Item) {
|
||||
const isError = item.error !== undefined && item.error !== null
|
||||
if (item.type === "image_generation_call" && item.result) {
|
||||
yield* Effect.fromResult(Encoding.decodeBase64(item.result)).pipe(
|
||||
|
||||
@@ -11,7 +11,7 @@ interface State {
|
||||
readonly responseID?: string
|
||||
}
|
||||
|
||||
const onOutputItem = Effect.fn("ResponsesCheckpoint.onOutputItem")(function* (
|
||||
const onOutputItem = Effect.fnUntraced(function* (
|
||||
state: State,
|
||||
input: OpenResponses.Event,
|
||||
) {
|
||||
@@ -63,7 +63,7 @@ export const make = <Body>(body: RouteBody<Body>): TriggerCompactOperation =>
|
||||
checkpoints: {},
|
||||
}),
|
||||
terminal: OpenResponses.terminal,
|
||||
step: Effect.fn("ResponsesCheckpoint.step")(function* (state: State, event: OpenResponses.Event) {
|
||||
step: Effect.fnUntraced(function* (state: State, event: OpenResponses.Event) {
|
||||
if (event.response?.id && state.responseID && event.response.id !== state.responseID)
|
||||
return yield* ProviderShared.eventError(source.id, "Compaction response ID changed during execution")
|
||||
if (event.type === "response.created") return [{ ...state, responseID: event.response?.id }, []] as const
|
||||
|
||||
@@ -33,7 +33,7 @@ export const onDone: (
|
||||
state: OpenResponses.ParserState,
|
||||
item: Item,
|
||||
tools: Definitions,
|
||||
) => Effect.Effect<OpenResponses.StepResult, AIError> = Effect.fn("ResponsesHostedTools.onDone")(
|
||||
) => Effect.Effect<OpenResponses.StepResult, AIError> = Effect.fnUntraced(
|
||||
function* (state, item, tools) {
|
||||
const tool = tools[item.type]
|
||||
if (!tool) return [state, []] satisfies OpenResponses.StepResult
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
import { Effect, Option } from "effect"
|
||||
import { Option, Result, Schema } from "effect"
|
||||
import { AIError, LLMEvent, type ProviderMetadata, type ToolCall } from "../../schema/index.js"
|
||||
import { eventError, parseToolInput, type ToolAccumulator } from "../shared.js"
|
||||
import { Json, eventError, type ToolAccumulator } from "../shared.js"
|
||||
import { parse } from "./partial-json.js"
|
||||
|
||||
type StreamKey = string | number
|
||||
const parsePartialInput = Option.liftThrowable(parse)
|
||||
const decodeInput = Schema.decodeUnknownResult(Json)
|
||||
|
||||
/**
|
||||
* One pending streamed tool call. Providers emit the tool identity and JSON
|
||||
@@ -69,31 +70,26 @@ const inputDelta = (tool: PendingTool, text: string) =>
|
||||
input: Option.getOrElse(parsePartialInput(tool.input), () => ({})),
|
||||
})
|
||||
|
||||
const toolCall = (route: string, tool: PendingTool, inputOverride?: string) => {
|
||||
const toolCall = (route: string, tool: PendingTool, inputOverride?: string): ToolCall | AIError => {
|
||||
const raw = inputOverride ?? tool.input
|
||||
return parseToolInput(route, tool.name, raw).pipe(
|
||||
Effect.catch((error) =>
|
||||
tool.providerExecuted
|
||||
? Effect.fail(error)
|
||||
: Effect.succeed(
|
||||
Option.getOrElse(
|
||||
Option.map(parsePartialInput(raw), (input) => input ?? {}),
|
||||
() => ({}),
|
||||
),
|
||||
),
|
||||
),
|
||||
Effect.map(
|
||||
(input): ToolCall =>
|
||||
LLMEvent.toolCall({
|
||||
id: tool.id,
|
||||
name: tool.name,
|
||||
namespace: tool.namespace,
|
||||
input,
|
||||
providerExecuted: tool.providerExecuted ? true : undefined,
|
||||
providerMetadata: tool.providerMetadata,
|
||||
}),
|
||||
),
|
||||
)
|
||||
const body = raw || "{}"
|
||||
const parsed = decodeInput(body)
|
||||
if (Result.isFailure(parsed) && tool.providerExecuted)
|
||||
return eventError(route, `Invalid JSON input for ${route} tool call ${tool.name}`, body, parsed.failure)
|
||||
const input = Result.isSuccess(parsed)
|
||||
? parsed.success
|
||||
: Option.getOrElse(
|
||||
Option.map(parsePartialInput(raw), (value) => value ?? {}),
|
||||
() => ({}),
|
||||
)
|
||||
return LLMEvent.toolCall({
|
||||
id: tool.id,
|
||||
name: tool.name,
|
||||
namespace: tool.namespace,
|
||||
input,
|
||||
providerExecuted: tool.providerExecuted ? true : undefined,
|
||||
providerMetadata: tool.providerMetadata,
|
||||
})
|
||||
}
|
||||
|
||||
const finishEvents = (tool: PendingTool, event: ToolCall): ReadonlyArray<LLMEvent> => [
|
||||
@@ -123,8 +119,7 @@ const appendTool = <K extends StreamKey>(
|
||||
}
|
||||
}
|
||||
|
||||
export const isError = <K extends StreamKey>(result: AppendOutcome<K> | AIError): result is AIError =>
|
||||
result instanceof AIError
|
||||
export const isError = <T>(result: T | AIError): result is AIError => result instanceof AIError
|
||||
|
||||
/**
|
||||
* Register a tool call whose start event arrived before any argument deltas.
|
||||
@@ -198,52 +193,40 @@ export const appendExisting = <K extends StreamKey>(
|
||||
): AppendOutcome<K> | AIError => append(tools, key, text) ?? eventError(route, missingToolMessage)
|
||||
|
||||
/**
|
||||
* Finalize one pending tool call: parse the accumulated raw JSON, remove it
|
||||
* from state, and recover incomplete local arguments when needed.
|
||||
* Finalize one pending tool call: parse the accumulated raw JSON (or an
|
||||
* authoritative final `input` override from `response.output_item.done`),
|
||||
* remove it from state, and recover incomplete local arguments when needed.
|
||||
* Missing keys are a no-op because some providers emit stop events for
|
||||
* non-tool content blocks.
|
||||
*/
|
||||
export const finish = <K extends StreamKey>(route: string, tools: State<K>, key: K) =>
|
||||
Effect.gen(function* () {
|
||||
const tool = tools[key]
|
||||
if (!tool) return { tools }
|
||||
return {
|
||||
tools: withoutTool(tools, key),
|
||||
events: finishEvents(tool, yield* toolCall(route, tool)),
|
||||
}
|
||||
})
|
||||
|
||||
/**
|
||||
* Finalize one pending tool call with an authoritative final input string.
|
||||
* OpenAI Responses can send accumulated deltas and then repeat the completed
|
||||
* arguments on `response.output_item.done`; the final value wins.
|
||||
*/
|
||||
export const finishWithInput = <K extends StreamKey>(route: string, tools: State<K>, key: K, input: string) =>
|
||||
Effect.gen(function* () {
|
||||
const tool = tools[key]
|
||||
if (!tool) return { tools }
|
||||
return {
|
||||
tools: withoutTool(tools, key),
|
||||
events: finishEvents(tool, yield* toolCall(route, tool, input)),
|
||||
}
|
||||
})
|
||||
export const finish = <K extends StreamKey>(route: string, tools: State<K>, key: K, input?: string) => {
|
||||
const tool = tools[key]
|
||||
if (!tool) return { tools }
|
||||
const event = toolCall(route, tool, input)
|
||||
if (isError(event)) return event
|
||||
return {
|
||||
tools: withoutTool(tools, key),
|
||||
events: finishEvents(tool, event),
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Finalize every pending tool call at once. OpenAI Chat has this shape: it does
|
||||
* not emit per-tool stop events, so all accumulated calls finish independently
|
||||
* when the choice receives a terminal `finish_reason`.
|
||||
*/
|
||||
export const finishAll = <K extends StreamKey>(route: string, tools: State<K>) =>
|
||||
Effect.gen(function* () {
|
||||
const pending = Object.values<PendingTool | undefined>(tools).filter(
|
||||
(tool): tool is PendingTool => tool !== undefined,
|
||||
)
|
||||
return {
|
||||
tools: empty<K>(),
|
||||
events: yield* Effect.forEach(pending, (tool) =>
|
||||
toolCall(route, tool).pipe(Effect.map((event) => finishEvents(tool, event))),
|
||||
).pipe(Effect.map((events) => events.flat())),
|
||||
}
|
||||
})
|
||||
export const finishAll = <K extends StreamKey>(route: string, tools: State<K>) => {
|
||||
const events: LLMEvent[] = []
|
||||
for (const tool of Object.values<PendingTool | undefined>(tools)) {
|
||||
if (!tool) continue
|
||||
const event = toolCall(route, tool)
|
||||
if (isError(event)) return event
|
||||
events.push(...finishEvents(tool, event))
|
||||
}
|
||||
return {
|
||||
tools: empty<K>(),
|
||||
events,
|
||||
}
|
||||
}
|
||||
|
||||
export * as ToolStream from "./tool-stream.js"
|
||||
@@ -19,7 +19,8 @@ describe("ToolStream", () => {
|
||||
if (ToolStream.isError(first)) return yield* first
|
||||
const second = ToolStream.appendOrStart(ADAPTER, first.tools, 0, { text: ':"weather"}' }, "missing tool")
|
||||
if (ToolStream.isError(second)) return yield* second
|
||||
const finished = yield* ToolStream.finish(ADAPTER, second.tools, 0)
|
||||
const finished = ToolStream.finish(ADAPTER, second.tools, 0)
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
expect(first.events).toEqual([
|
||||
{ type: "tool-input-start", id: "call_1", name: "lookup" },
|
||||
@@ -91,7 +92,8 @@ describe("ToolStream", () => {
|
||||
"missing tool",
|
||||
)
|
||||
if (ToolStream.isError(second)) return yield* second
|
||||
const finished = yield* ToolStream.finish(ADAPTER, second.tools, 0)
|
||||
const finished = ToolStream.finish(ADAPTER, second.tools, 0)
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
expect(finished.events).toEqual([
|
||||
{ type: "tool-input-end", id: "call_1", name: "lookup" },
|
||||
@@ -114,7 +116,8 @@ describe("ToolStream", () => {
|
||||
name: "lookup",
|
||||
input: '{"query":"partial"}',
|
||||
})
|
||||
const finished = yield* ToolStream.finishWithInput(ADAPTER, tools, "item_1", '{"query":"final"}')
|
||||
const finished = ToolStream.finish(ADAPTER, tools, "item_1", '{"query":"final"}')
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
expect(finished).toEqual({
|
||||
tools: {},
|
||||
@@ -133,7 +136,8 @@ describe("ToolStream", () => {
|
||||
name: "lookup",
|
||||
input: '{"query":"partial',
|
||||
})
|
||||
const finished = yield* ToolStream.finish(ADAPTER, tools, "item_1")
|
||||
const finished = ToolStream.finish(ADAPTER, tools, "item_1")
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
expect(finished).toEqual({
|
||||
tools: {},
|
||||
@@ -152,7 +156,8 @@ describe("ToolStream", () => {
|
||||
name: "lookup",
|
||||
input: '{"path":"A\\H","text":"first\tsecond"}',
|
||||
})
|
||||
const finished = yield* ToolStream.finish(ADAPTER, tools, "item_1")
|
||||
const finished = ToolStream.finish(ADAPTER, tools, "item_1")
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
expect(finished.events).toEqual([
|
||||
{ type: "tool-input-end", id: "call_1", name: "lookup" },
|
||||
@@ -168,7 +173,8 @@ describe("ToolStream", () => {
|
||||
name: "lookup",
|
||||
input: "invalid",
|
||||
})
|
||||
const finished = yield* ToolStream.finish(ADAPTER, tools, "item_1")
|
||||
const finished = ToolStream.finish(ADAPTER, tools, "item_1")
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
expect(finished.events).toEqual([
|
||||
{ type: "tool-input-end", id: "call_1", name: "lookup" },
|
||||
@@ -189,7 +195,8 @@ describe("ToolStream", () => {
|
||||
name: "lookup",
|
||||
input: '{"query":"partial',
|
||||
})
|
||||
const finished = yield* ToolStream.finishAll(ADAPTER, tools)
|
||||
const finished = ToolStream.finishAll(ADAPTER, tools)
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
expect(finished).toEqual({
|
||||
tools: {},
|
||||
@@ -211,9 +218,9 @@ describe("ToolStream", () => {
|
||||
input: '{"query":"partial',
|
||||
providerExecuted: true,
|
||||
})
|
||||
const result = yield* Effect.exit(ToolStream.finish(ADAPTER, tools, "item_1"))
|
||||
const result = ToolStream.finish(ADAPTER, tools, "item_1")
|
||||
|
||||
expect(result._tag).toBe("Failure")
|
||||
expect(result).toBeInstanceOf(AIError)
|
||||
}),
|
||||
)
|
||||
|
||||
@@ -230,7 +237,8 @@ describe("ToolStream", () => {
|
||||
input: '{"query":"docs"}',
|
||||
providerExecuted: true,
|
||||
})
|
||||
const finished = yield* ToolStream.finishAll(ADAPTER, tools)
|
||||
const finished = ToolStream.finishAll(ADAPTER, tools)
|
||||
if (ToolStream.isError(finished)) return yield* finished
|
||||
|
||||
expect(finished).toEqual({
|
||||
tools: {},
|
||||
|
||||
Reference in new issue
Block a user