Compare commits

..
1 Commits
13 changed files with 178 additions and 188 deletions

No files matched your search

+35 -36
View File
@@ -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 = (() => {
+5 -3
View File
@@ -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: {
+3 -3
View File
@@ -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))
+5 -4
View File
@@ -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({
+40 -49
View File
@@ -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])
+9 -6
View File
@@ -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
+49 -66
View File
@@ -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"
+18 -10
View File
@@ -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: {},