mirror of
https://github.com/anomalyco/opencode.git
synced 2026-07-29 06:31:55 -04:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 92ee702eeb |
@@ -55,6 +55,12 @@ const OpenResponsesOutputText = Schema.Struct({
|
||||
text: Schema.String,
|
||||
})
|
||||
|
||||
const ExtendedOutputText = Schema.Struct({
|
||||
type: Schema.tag("output_text"),
|
||||
text: Schema.String,
|
||||
annotations: Schema.Array(Schema.Unknown),
|
||||
})
|
||||
|
||||
const OpenResponsesReasoningSummaryText = Schema.Struct({
|
||||
type: Schema.tag("summary_text"),
|
||||
text: Schema.String,
|
||||
@@ -86,10 +92,19 @@ const OpenResponsesFunctionCallOutput = Schema.Union([
|
||||
Schema.Array(OpenResponsesFunctionCallOutputContent),
|
||||
])
|
||||
|
||||
const OpenResponsesInputItem = Schema.Union([
|
||||
export const InputItem = Schema.Union([
|
||||
Schema.Struct({ role: Schema.tag("system"), content: Schema.String }),
|
||||
Schema.Struct({ role: Schema.tag("user"), content: Schema.Array(OpenResponsesInputContent) }),
|
||||
Schema.Struct({ role: Schema.tag("assistant"), content: Schema.Array(OpenResponsesOutputText) }),
|
||||
Schema.Struct({ role: Schema.tag("assistant"), content: Schema.String, phase: Schema.optionalKey(Schema.Unknown) }),
|
||||
Schema.Struct({
|
||||
type: Schema.tag("message"),
|
||||
id: Schema.String,
|
||||
status: Schema.Literals(["in_progress", "completed", "incomplete"]),
|
||||
role: Schema.tag("assistant"),
|
||||
content: Schema.Array(ExtendedOutputText),
|
||||
phase: Schema.optionalKey(Schema.Unknown),
|
||||
}),
|
||||
OpenResponsesReasoningItem,
|
||||
OpenResponsesItemReference,
|
||||
Schema.Struct({
|
||||
@@ -104,7 +119,7 @@ const OpenResponsesInputItem = Schema.Union([
|
||||
output: OpenResponsesFunctionCallOutput,
|
||||
}),
|
||||
])
|
||||
type OpenResponsesInputItem = Schema.Schema.Type<typeof OpenResponsesInputItem>
|
||||
type OpenResponsesInputItem = Schema.Schema.Type<typeof InputItem>
|
||||
|
||||
// Mutable counterpart of the schema reasoning item so `lowerMessages` can fold
|
||||
// multiple streamed summary parts into the same item before flushing.
|
||||
@@ -135,7 +150,7 @@ export const ToolChoice = Schema.Union([
|
||||
// transports in sync without a destructure-and-strip dance.
|
||||
export const coreFields = {
|
||||
model: Schema.String,
|
||||
input: Schema.Array(OpenResponsesInputItem),
|
||||
input: Schema.Array(InputItem),
|
||||
instructions: Schema.optional(Schema.String),
|
||||
tools: optionalArray(Tool),
|
||||
tool_choice: Schema.optional(ToolChoice),
|
||||
@@ -238,6 +253,25 @@ export interface Extension {
|
||||
readonly media: ProviderShared.ValidatedMedia
|
||||
readonly request: LLMRequest
|
||||
}) => MediaInput | undefined
|
||||
readonly messageMetadata?: (
|
||||
item: StreamItem,
|
||||
id: string,
|
||||
previous: ProviderMetadata | undefined,
|
||||
providerMetadataKey: string,
|
||||
) => ProviderMetadata | undefined
|
||||
readonly messageContentMetadata?: (
|
||||
providerMetadata: ProviderMetadata | undefined,
|
||||
item: StreamItem,
|
||||
index: number,
|
||||
) => ProviderMetadata | undefined
|
||||
}
|
||||
|
||||
export type ExtendedExtension = Extension & {
|
||||
readonly lowerAssistantText?: (
|
||||
parts: ReadonlyArray<TextPart>,
|
||||
request: LLMRequest,
|
||||
message: LLMRequest["messages"][number],
|
||||
) => ReadonlyArray<OpenResponsesInputItem>
|
||||
}
|
||||
|
||||
const BASE: Extension = { id: ADAPTER, name: NAME }
|
||||
@@ -249,10 +283,20 @@ export interface ParserState {
|
||||
readonly tools: ToolStream.State<string>
|
||||
readonly hasFunctionCall: boolean
|
||||
readonly lifecycle: Lifecycle.State
|
||||
readonly messageItems: Readonly<Record<string, MessageStreamItem>>
|
||||
readonly messageContentIDs: ReadonlySet<string>
|
||||
readonly messageContentMetadata?: Extension["messageContentMetadata"]
|
||||
readonly messageMetadata?: Extension["messageMetadata"]
|
||||
readonly nextMessageContentID: number
|
||||
readonly reasoningItems: Readonly<Record<string, ReasoningStreamItem>>
|
||||
readonly store: boolean | undefined
|
||||
}
|
||||
|
||||
interface MessageStreamItem {
|
||||
readonly providerMetadata?: ProviderMetadata
|
||||
readonly content: Readonly<Record<number, { readonly id: string; readonly text: string }>>
|
||||
}
|
||||
|
||||
type ReasoningSummaryStatus = "active" | "can-conclude" | "concluded"
|
||||
|
||||
interface ReasoningStreamItem {
|
||||
@@ -377,7 +421,10 @@ const lowerToolResultOutput = Effect.fn("OpenResponses.lowerToolResultOutput")(f
|
||||
return yield* Effect.forEach(content, (item) => lowerToolResultContentItem(item, request, extension))
|
||||
})
|
||||
|
||||
const lowerMessages = Effect.fn("OpenResponses.lowerMessages")(function* (request: LLMRequest, extension: Extension) {
|
||||
const lowerMessages = Effect.fn("OpenResponses.lowerMessages")(function* (
|
||||
request: LLMRequest,
|
||||
extension: ExtendedExtension,
|
||||
) {
|
||||
const system: OpenResponsesInputItem[] =
|
||||
request.system.length === 0 ? [] : [{ role: "system", content: ProviderShared.joinText(request.system) }]
|
||||
const input: OpenResponsesInputItem[] = [...system]
|
||||
@@ -412,7 +459,14 @@ const lowerMessages = Effect.fn("OpenResponses.lowerMessages")(function* (reques
|
||||
const hostedToolReferences = new Set<string>()
|
||||
const flushText = () => {
|
||||
if (content.length === 0) return
|
||||
input.push({ role: "assistant", content: content.map((part) => ({ type: "output_text", text: part.text })) })
|
||||
input.push(
|
||||
...(extension.lowerAssistantText?.(content, request, message) ?? [
|
||||
{
|
||||
role: "assistant" as const,
|
||||
content: content.map((part) => ({ type: "output_text" as const, text: part.text })),
|
||||
},
|
||||
]),
|
||||
)
|
||||
content.splice(0, content.length)
|
||||
}
|
||||
for (const part of message.content) {
|
||||
@@ -515,7 +569,7 @@ const lowerOptions = (request: LLMRequest) => {
|
||||
|
||||
export const fromRequest = Effect.fn("OpenResponses.fromRequest")(function* (
|
||||
request: LLMRequest,
|
||||
extension: Extension = BASE,
|
||||
extension: ExtendedExtension = BASE,
|
||||
) {
|
||||
const generation = request.generation
|
||||
const toolSchemaCompatibility = request.model.compatibility?.toolSchema
|
||||
@@ -595,8 +649,83 @@ const NO_EVENTS: StepResult["1"] = []
|
||||
const TERMINAL_TYPES = new Set(["response.completed", "response.incomplete", "response.failed"])
|
||||
export const terminal = (event: Event) => TERMINAL_TYPES.has(event.type)
|
||||
|
||||
const ensureMessageContent = (state: ParserState, event: Event) => {
|
||||
const itemID = event.item_id ?? "text-0"
|
||||
const index = typeof event.content_index === "number" ? event.content_index : 0
|
||||
const item = state.messageItems[itemID] ?? { content: {} }
|
||||
const existing = item.content[index]
|
||||
if (existing) return { state, itemID, index, item, content: existing }
|
||||
const findID = (next: number): readonly [string, number] => {
|
||||
const id = `responses-text-${next}`
|
||||
return state.messageContentIDs.has(id) ? findID(next + 1) : [id, next + 1]
|
||||
}
|
||||
const [id, nextMessageContentID] =
|
||||
index === 0 && !state.messageContentIDs.has(itemID)
|
||||
? ([itemID, state.nextMessageContentID] as const)
|
||||
: findID(state.nextMessageContentID)
|
||||
const content = { id, text: "" }
|
||||
const nextItem = { ...item, content: { ...item.content, [index]: content } }
|
||||
return {
|
||||
state: {
|
||||
...state,
|
||||
messageItems: { ...state.messageItems, [itemID]: nextItem },
|
||||
messageContentIDs: new Set([...state.messageContentIDs, id]),
|
||||
nextMessageContentID,
|
||||
},
|
||||
itemID,
|
||||
index,
|
||||
item: nextItem,
|
||||
content,
|
||||
}
|
||||
}
|
||||
|
||||
const updateMessageContent = (
|
||||
state: ParserState,
|
||||
itemID: string,
|
||||
index: number,
|
||||
content: { readonly id: string; readonly text: string },
|
||||
): ParserState => ({
|
||||
...state,
|
||||
messageItems: {
|
||||
...state.messageItems,
|
||||
[itemID]: {
|
||||
...state.messageItems[itemID],
|
||||
content: { ...state.messageItems[itemID]?.content, [index]: content },
|
||||
},
|
||||
},
|
||||
})
|
||||
|
||||
const closeOtherMessageContent = (state: ParserState, events: LLMEvent[], item: MessageStreamItem, index: number) =>
|
||||
Object.entries(item.content).reduce(
|
||||
(lifecycle, entry) =>
|
||||
Number(entry[0]) === index ? lifecycle : Lifecycle.textEnd(lifecycle, events, entry[1].id, item.providerMetadata),
|
||||
state.lifecycle,
|
||||
)
|
||||
|
||||
const appendOutputText = (state: ParserState, event: Event, text: string): StepResult => {
|
||||
const ensured = ensureMessageContent(state, event)
|
||||
const events: LLMEvent[] = []
|
||||
const lifecycle = Lifecycle.textStart(
|
||||
closeOtherMessageContent(ensured.state, events, ensured.item, ensured.index),
|
||||
events,
|
||||
ensured.content.id,
|
||||
ensured.item.providerMetadata,
|
||||
)
|
||||
return [
|
||||
{
|
||||
...updateMessageContent(ensured.state, ensured.itemID, ensured.index, {
|
||||
...ensured.content,
|
||||
text: ensured.content.text + text,
|
||||
}),
|
||||
lifecycle: Lifecycle.textDelta(lifecycle, events, ensured.content.id, text),
|
||||
},
|
||||
events,
|
||||
]
|
||||
}
|
||||
|
||||
const onOutputTextDelta = (state: ParserState, event: Event): StepResult => {
|
||||
if (!event.delta) return [state, NO_EVENTS]
|
||||
if (state.messageMetadata) return appendOutputText(state, event, event.delta)
|
||||
const events: LLMEvent[] = []
|
||||
return [
|
||||
{ ...state, lifecycle: Lifecycle.textDelta(state.lifecycle, events, event.item_id ?? "text-0", event.delta) },
|
||||
@@ -605,6 +734,30 @@ const onOutputTextDelta = (state: ParserState, event: Event): StepResult => {
|
||||
}
|
||||
|
||||
const onOutputTextDone = (state: ParserState, event: Event): StepResult => {
|
||||
if (state.messageMetadata) {
|
||||
const text = typeof event.text === "string" ? event.text : undefined
|
||||
const ensured = ensureMessageContent(state, event)
|
||||
if (text === undefined) return [ensured.state, NO_EVENTS]
|
||||
if (text === ensured.content.text) {
|
||||
if (ensured.state.lifecycle.text.has(ensured.content.id)) return [ensured.state, NO_EVENTS]
|
||||
const events: LLMEvent[] = []
|
||||
return [
|
||||
{
|
||||
...ensured.state,
|
||||
lifecycle: Lifecycle.textStart(
|
||||
closeOtherMessageContent(ensured.state, events, ensured.item, ensured.index),
|
||||
events,
|
||||
ensured.content.id,
|
||||
ensured.item.providerMetadata,
|
||||
),
|
||||
},
|
||||
events,
|
||||
]
|
||||
}
|
||||
if (text.startsWith(ensured.content.text))
|
||||
return appendOutputText(ensured.state, event, text.slice(ensured.content.text.length))
|
||||
return [ensured.state, NO_EVENTS]
|
||||
}
|
||||
const events: LLMEvent[] = []
|
||||
return [{ ...state, lifecycle: Lifecycle.textEnd(state.lifecycle, events, event.item_id ?? "text-0") }, events]
|
||||
}
|
||||
@@ -643,6 +796,28 @@ const reasoningMetadata = (state: ParserState, item: StreamItem & { id: string }
|
||||
// best-effort, not guaranteed.
|
||||
const onOutputItemAdded = (state: ParserState, event: Event): StepResult => {
|
||||
const item = event.item
|
||||
if (item?.type === "message" && item.id && state.messageMetadata) {
|
||||
const existing = state.messageItems[item.id]
|
||||
return [
|
||||
{
|
||||
...state,
|
||||
messageItems: {
|
||||
...state.messageItems,
|
||||
[item.id]: {
|
||||
...existing,
|
||||
providerMetadata: state.messageMetadata(
|
||||
item,
|
||||
item.id,
|
||||
existing?.providerMetadata,
|
||||
state.providerMetadataKey,
|
||||
),
|
||||
content: existing?.content ?? {},
|
||||
},
|
||||
},
|
||||
},
|
||||
NO_EVENTS,
|
||||
]
|
||||
}
|
||||
if (item && isReasoningItem(item)) {
|
||||
const events: LLMEvent[] = []
|
||||
return [
|
||||
@@ -799,6 +974,24 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
|
||||
const item = event.item
|
||||
if (!item) return [state, NO_EVENTS] satisfies StepResult
|
||||
|
||||
if (item.type === "message" && item.id && state.messageMetadata) {
|
||||
const events: LLMEvent[] = []
|
||||
const messageItem = state.messageItems[item.id]
|
||||
const { [item.id]: _finished, ...messageItems } = state.messageItems
|
||||
const metadata = state.messageMetadata(item, item.id, messageItem?.providerMetadata, state.providerMetadataKey)
|
||||
const lifecycle = Object.entries(messageItem?.content ?? {}).reduce(
|
||||
(lifecycle, entry) =>
|
||||
Lifecycle.textEnd(
|
||||
lifecycle,
|
||||
events,
|
||||
entry[1].id,
|
||||
state.messageContentMetadata?.(metadata, item, Number(entry[0])) ?? metadata,
|
||||
),
|
||||
state.lifecycle,
|
||||
)
|
||||
return [{ ...state, lifecycle, messageItems }, events] satisfies StepResult
|
||||
}
|
||||
|
||||
if (item.type === "message" && item.id) return onOutputTextDone(state, { ...event, item_id: item.id })
|
||||
|
||||
if (item.type === "function_call") {
|
||||
@@ -933,6 +1126,11 @@ export const initial = (request: LLMRequest, extension: Extension = BASE): Parse
|
||||
hasFunctionCall: false,
|
||||
tools: ToolStream.empty<string>(),
|
||||
lifecycle: Lifecycle.initial(),
|
||||
messageItems: {},
|
||||
messageContentIDs: new Set<string>(),
|
||||
messageContentMetadata: extension.messageContentMetadata,
|
||||
messageMetadata: extension.messageMetadata,
|
||||
nextMessageContentID: 0,
|
||||
reasoningItems: {},
|
||||
store: OpenResponsesOptions.resolve(request).store,
|
||||
})
|
||||
|
||||
@@ -4,11 +4,19 @@ import { Auth } from "../route/auth"
|
||||
import { Endpoint } from "../route/endpoint"
|
||||
import { Protocol } from "../route/protocol"
|
||||
import { HttpTransport, WebSocketTransport } from "../route/transport"
|
||||
import { LLMEvent, LLMRequest, type JsonSchema, type ToolDefinition } from "../schema"
|
||||
import {
|
||||
LLMEvent,
|
||||
LLMRequest,
|
||||
type JsonSchema,
|
||||
type ProviderMetadata,
|
||||
type TextPart,
|
||||
type ToolDefinition,
|
||||
} from "../schema"
|
||||
import { OpenResponses } from "./open-responses"
|
||||
import { optionalArray, ProviderShared } from "./shared"
|
||||
import { Lifecycle } from "./utils/lifecycle"
|
||||
import { OpenAIImage } from "./utils/openai-image"
|
||||
import { OpenResponsesOptions } from "./utils/open-responses-options"
|
||||
import { ToolSchemaProjection } from "./utils/tool-schema"
|
||||
|
||||
const ADAPTER = "openai-responses"
|
||||
@@ -35,8 +43,36 @@ const OpenAIResponsesToolChoice = Schema.Union([
|
||||
Schema.Struct({ type: Schema.tag("image_generation") }),
|
||||
])
|
||||
|
||||
const OpenAIResponsesOutputText = Schema.Struct({
|
||||
type: Schema.tag("output_text"),
|
||||
text: Schema.String,
|
||||
annotations: Schema.Array(Schema.Unknown),
|
||||
})
|
||||
|
||||
const OpenAIResponsesMessagePhase = Schema.Literals(["commentary", "final_answer"])
|
||||
type OpenAIResponsesMessagePhase = Schema.Schema.Type<typeof OpenAIResponsesMessagePhase>
|
||||
|
||||
const OpenAIResponsesInputItem = Schema.Union([
|
||||
Schema.Struct({
|
||||
role: Schema.tag("assistant"),
|
||||
content: Schema.String,
|
||||
phase: Schema.optionalKey(Schema.NullOr(OpenAIResponsesMessagePhase)),
|
||||
}),
|
||||
Schema.Struct({
|
||||
type: Schema.tag("message"),
|
||||
id: Schema.String,
|
||||
status: Schema.Literals(["in_progress", "completed", "incomplete"]),
|
||||
role: Schema.tag("assistant"),
|
||||
content: Schema.Array(OpenAIResponsesOutputText),
|
||||
phase: Schema.optionalKey(Schema.NullOr(OpenAIResponsesMessagePhase)),
|
||||
}),
|
||||
OpenResponses.InputItem,
|
||||
])
|
||||
type OpenAIResponsesInputItem = Schema.Schema.Type<typeof OpenAIResponsesInputItem>
|
||||
|
||||
const OpenAIResponsesCoreFields = {
|
||||
...OpenResponses.coreFields,
|
||||
input: Schema.Array(OpenAIResponsesInputItem),
|
||||
tools: optionalArray(OpenAIResponsesTools),
|
||||
tool_choice: Schema.optional(OpenAIResponsesToolChoice),
|
||||
}
|
||||
@@ -57,9 +93,161 @@ const OpenAIResponsesWebSocketMessage = Schema.StructWithRest(
|
||||
type OpenAIResponsesWebSocketMessage = Schema.Schema.Type<typeof OpenAIResponsesWebSocketMessage>
|
||||
const encodeWebSocketMessage = Schema.encodeSync(Schema.fromJsonString(OpenAIResponsesWebSocketMessage))
|
||||
|
||||
const messagePhase = (part: TextPart, providerMetadataKey: string): OpenAIResponsesMessagePhase | null | undefined => {
|
||||
const metadata = ProviderShared.isRecord(part.providerMetadata)
|
||||
? part.providerMetadata[providerMetadataKey]
|
||||
: undefined
|
||||
if (!ProviderShared.isRecord(metadata)) return undefined
|
||||
return metadata.phase === "commentary" || metadata.phase === "final_answer" || metadata.phase === null
|
||||
? metadata.phase
|
||||
: undefined
|
||||
}
|
||||
|
||||
const messageItemID = (part: TextPart, providerMetadataKey: string) => {
|
||||
const metadata = ProviderShared.isRecord(part.providerMetadata)
|
||||
? part.providerMetadata[providerMetadataKey]
|
||||
: undefined
|
||||
return ProviderShared.isRecord(metadata) && typeof metadata.itemId === "string" && metadata.itemId.length > 0
|
||||
? metadata.itemId
|
||||
: undefined
|
||||
}
|
||||
|
||||
const messageStatus = (part: TextPart, providerMetadataKey: string) => {
|
||||
const metadata = ProviderShared.isRecord(part.providerMetadata)
|
||||
? part.providerMetadata[providerMetadataKey]
|
||||
: undefined
|
||||
if (!ProviderShared.isRecord(metadata)) return undefined
|
||||
return metadata.status === "in_progress" || metadata.status === "completed" || metadata.status === "incomplete"
|
||||
? metadata.status
|
||||
: undefined
|
||||
}
|
||||
|
||||
const messageAnnotations = (part: TextPart, providerMetadataKey: string) => {
|
||||
const metadata = ProviderShared.isRecord(part.providerMetadata)
|
||||
? part.providerMetadata[providerMetadataKey]
|
||||
: undefined
|
||||
return ProviderShared.isRecord(metadata) && Array.isArray(metadata.annotations) ? metadata.annotations : []
|
||||
}
|
||||
|
||||
const lowerAssistantText = (
|
||||
parts: ReadonlyArray<TextPart>,
|
||||
request: LLMRequest,
|
||||
message: LLMRequest["messages"][number],
|
||||
) => {
|
||||
const providerMetadataKey = request.model.route.providerMetadataKey ?? "openai"
|
||||
const dropItemIdentity =
|
||||
OpenResponsesOptions.resolve(request).store === false &&
|
||||
message.content.some((part) => {
|
||||
if (part.type !== "reasoning") return false
|
||||
const metadata = part.providerMetadata?.[providerMetadataKey]
|
||||
return (
|
||||
ProviderShared.isRecord(metadata) &&
|
||||
typeof metadata.itemId === "string" &&
|
||||
typeof metadata.reasoningEncryptedContent !== "string"
|
||||
)
|
||||
})
|
||||
const input: OpenAIResponsesInputItem[] = []
|
||||
const content: TextPart[] = []
|
||||
let phase: OpenAIResponsesMessagePhase | null | undefined
|
||||
let itemID: string | undefined
|
||||
let status: "in_progress" | "completed" | "incomplete" | undefined
|
||||
const flush = () => {
|
||||
if (content.length === 0) return
|
||||
if (phase === undefined && content.every((part) => part.text.length === 0)) {
|
||||
content.splice(0, content.length)
|
||||
itemID = undefined
|
||||
status = undefined
|
||||
return
|
||||
}
|
||||
if (phase === undefined && itemID === undefined) {
|
||||
input.push({
|
||||
role: "assistant",
|
||||
content: content.map((part) => ({ type: "output_text", text: part.text })),
|
||||
})
|
||||
content.splice(0, content.length)
|
||||
status = undefined
|
||||
return
|
||||
}
|
||||
input.push(
|
||||
itemID
|
||||
? {
|
||||
type: "message",
|
||||
id: itemID,
|
||||
status: status ?? "completed",
|
||||
role: "assistant",
|
||||
content: content.map((part) => ({
|
||||
type: "output_text",
|
||||
text: part.text,
|
||||
annotations: messageAnnotations(part, providerMetadataKey),
|
||||
})),
|
||||
...(phase !== undefined ? { phase } : {}),
|
||||
}
|
||||
: {
|
||||
role: "assistant",
|
||||
content: ProviderShared.joinText(content),
|
||||
...(phase !== undefined ? { phase } : {}),
|
||||
},
|
||||
)
|
||||
content.splice(0, content.length)
|
||||
phase = undefined
|
||||
itemID = undefined
|
||||
status = undefined
|
||||
}
|
||||
for (const part of parts) {
|
||||
const nextPhase = messagePhase(part, providerMetadataKey)
|
||||
const nextItemID = dropItemIdentity ? undefined : messageItemID(part, providerMetadataKey)
|
||||
if (content.length > 0 && (phase !== nextPhase || itemID !== nextItemID)) flush()
|
||||
phase = nextPhase
|
||||
itemID = nextItemID
|
||||
status = messageStatus(part, providerMetadataKey) ?? status
|
||||
content.push(part)
|
||||
}
|
||||
flush()
|
||||
return input
|
||||
}
|
||||
|
||||
const messageMetadata = (
|
||||
item: OpenResponses.StreamItem,
|
||||
id: string,
|
||||
previous: ProviderMetadata | undefined,
|
||||
providerMetadataKey: string,
|
||||
) => {
|
||||
const prior = previous?.[providerMetadataKey]
|
||||
const metadata = ProviderShared.isRecord(prior) ? prior : {}
|
||||
const phase = item.phase !== undefined ? item.phase : metadata.phase
|
||||
const status = item.status !== undefined ? item.status : metadata.status
|
||||
if (previous === undefined && phase === undefined && status === undefined) return undefined
|
||||
return {
|
||||
[providerMetadataKey]: {
|
||||
...metadata,
|
||||
itemId: id,
|
||||
...(phase === "commentary" || phase === "final_answer" || phase === null ? { phase } : {}),
|
||||
...(status === "in_progress" || status === "completed" || status === "incomplete" ? { status } : {}),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
const messageContentMetadata = (
|
||||
providerMetadata: ProviderMetadata | undefined,
|
||||
item: OpenResponses.StreamItem,
|
||||
index: number,
|
||||
) => {
|
||||
const content = Array.isArray(item.content) ? item.content[index] : undefined
|
||||
if (!ProviderShared.isRecord(content) || content.type !== "output_text" || !Array.isArray(content.annotations))
|
||||
return providerMetadata
|
||||
if (!providerMetadata) return undefined
|
||||
const key = Object.keys(providerMetadata)[0]
|
||||
if (!key) return providerMetadata
|
||||
const metadata = providerMetadata[key]
|
||||
return { [key]: { ...(ProviderShared.isRecord(metadata) ? metadata : {}), annotations: content.annotations } }
|
||||
}
|
||||
|
||||
const extension = {
|
||||
id: ADAPTER,
|
||||
name: NAME,
|
||||
lowerAssistantText,
|
||||
messageMetadata,
|
||||
messageContentMetadata,
|
||||
lowerMedia: ({ part, media, request }) => {
|
||||
if (request.model.provider !== "xai" || media.mime !== "application/pdf") return undefined
|
||||
return {
|
||||
@@ -69,7 +257,7 @@ const extension = {
|
||||
mime_type: media.mime,
|
||||
}
|
||||
},
|
||||
} satisfies OpenResponses.Extension
|
||||
} satisfies OpenResponses.ExtendedExtension
|
||||
|
||||
const nativeImageToolInput = (tool: ToolDefinition) => {
|
||||
const native = tool.native?.openai
|
||||
|
||||
@@ -17,7 +17,7 @@ export const stepStart = (state: State, events: LLMEvent[]): State => {
|
||||
export const textStart = (state: State, events: LLMEvent[], id: string, providerMetadata?: ProviderMetadata): State => {
|
||||
if (state.text.has(id)) return state
|
||||
const stepped = stepStart(state, events)
|
||||
events.push(LLMEvent.textStart({ id, providerMetadata }))
|
||||
events.push(LLMEvent.textStart({ id, ...(providerMetadata ? { providerMetadata } : {}) }))
|
||||
return { ...stepped, text: new Set([...stepped.text, id]) }
|
||||
}
|
||||
|
||||
@@ -68,7 +68,7 @@ export const reasoningEnd = (
|
||||
export const textEnd = (state: State, events: LLMEvent[], id: string, providerMetadata?: ProviderMetadata): State => {
|
||||
if (!state.text.has(id)) return state
|
||||
const stepped = stepStart(state, events)
|
||||
events.push(LLMEvent.textEnd({ id, providerMetadata }))
|
||||
events.push(LLMEvent.textEnd({ id, ...(providerMetadata ? { providerMetadata } : {}) }))
|
||||
const text = new Set(stepped.text)
|
||||
text.delete(id)
|
||||
return { ...stepped, text }
|
||||
|
||||
+52
File diff suppressed because one or more lines are too long
@@ -0,0 +1,89 @@
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Effect } from "effect"
|
||||
import { LLM, Message } from "../../src"
|
||||
import { configure } from "../../src/providers/openai"
|
||||
import { OpenAIResponses } from "../../src/protocols/openai-responses"
|
||||
import { LLMClient } from "../../src/route"
|
||||
import { weatherTool } from "../recorded-scenarios"
|
||||
import { recordedTests } from "../recorded-test"
|
||||
|
||||
const model = configure({
|
||||
apiKey: process.env.OPENAI_API_KEY ?? "fixture",
|
||||
}).responses("gpt-5.6-sol")
|
||||
|
||||
const recorded = recordedTests({
|
||||
prefix: "openai-responses-phase",
|
||||
provider: "openai",
|
||||
protocol: "openai-responses",
|
||||
requires: ["OPENAI_API_KEY"],
|
||||
})
|
||||
|
||||
describe("OpenAI Responses phase recorded", () => {
|
||||
recorded.effect.with("round-trips commentary into a final answer", { tags: ["phase", "tool"] }, () =>
|
||||
Effect.gen(function* () {
|
||||
const user = Message.user("What is the weather in Paris?")
|
||||
const first = yield* LLMClient.generate(
|
||||
LLM.request({
|
||||
model,
|
||||
system:
|
||||
"Before calling get_weather, briefly tell the user you are checking. Then call get_weather exactly once. Do not provide the final answer until its result is available.",
|
||||
messages: [user],
|
||||
tools: [weatherTool],
|
||||
generation: { maxTokens: 100 },
|
||||
}),
|
||||
)
|
||||
const call = first.toolCalls[0]
|
||||
if (!call) throw new Error("OpenAI Responses did not return the expected weather tool call")
|
||||
|
||||
expect(call).toMatchObject({ name: "get_weather", input: { city: "Paris" } })
|
||||
const commentary = first.message.content.find(
|
||||
(part) => part.type === "text" && part.providerMetadata?.openai?.phase === "commentary",
|
||||
)
|
||||
if (!commentary || commentary.type !== "text") throw new Error("OpenAI Responses did not return commentary text")
|
||||
const itemID = commentary.providerMetadata?.openai?.itemId
|
||||
if (typeof itemID !== "string") throw new Error("OpenAI Responses commentary did not include an item ID")
|
||||
expect(commentary).toEqual({
|
||||
type: "text",
|
||||
text: "I’ll check the current weather in Paris.",
|
||||
providerMetadata: {
|
||||
openai: { itemId: itemID, phase: "commentary", status: "completed", annotations: [] },
|
||||
},
|
||||
})
|
||||
|
||||
const continuation = LLM.request({
|
||||
model,
|
||||
system:
|
||||
"Before calling get_weather, briefly tell the user you are checking. Then call get_weather exactly once. After its result, answer exactly: Paris is sunny.",
|
||||
messages: [
|
||||
user,
|
||||
first.message,
|
||||
Message.tool({
|
||||
id: call.id,
|
||||
name: call.name,
|
||||
result: { temperature: 22, condition: "sunny" },
|
||||
}),
|
||||
],
|
||||
tools: [weatherTool],
|
||||
generation: { maxTokens: 100 },
|
||||
})
|
||||
const prepared = yield* LLMClient.prepare<OpenAIResponses.OpenAIResponsesBody>(continuation)
|
||||
expect(prepared.body.input).toContainEqual({
|
||||
type: "message",
|
||||
id: itemID,
|
||||
status: "completed",
|
||||
role: "assistant",
|
||||
content: [{ type: "output_text", text: commentary.text, annotations: [] }],
|
||||
phase: "commentary",
|
||||
})
|
||||
|
||||
const second = yield* LLMClient.generate(continuation)
|
||||
|
||||
expect(second.text.trim()).toBe("Paris is sunny.")
|
||||
expect(
|
||||
second.message.content.some(
|
||||
(part) => part.type === "text" && part.providerMetadata?.openai?.phase === "final_answer",
|
||||
),
|
||||
).toBeTrue()
|
||||
}),
|
||||
)
|
||||
})
|
||||
@@ -920,6 +920,7 @@ describe("OpenAI Responses route", () => {
|
||||
sseEvents(
|
||||
{ type: "response.output_text.delta", item_id: "msg_1", delta: "First" },
|
||||
{ type: "response.output_text.done", item_id: "msg_1" },
|
||||
{ type: "response.output_item.done", item: { type: "message", id: "msg_1" } },
|
||||
{ type: "response.output_text.delta", item_id: "msg_2", delta: "Second" },
|
||||
{ type: "response.output_item.done", item: { type: "message", id: "msg_2" } },
|
||||
{ type: "response.completed", response: { id: "resp_1" } },
|
||||
@@ -939,6 +940,94 @@ describe("OpenAI Responses route", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("preserves and replays OpenAI response message phases", () =>
|
||||
Effect.gen(function* () {
|
||||
const response = yield* LLMClient.generate(request).pipe(
|
||||
Effect.provide(
|
||||
fixedResponse(
|
||||
sseEvents(
|
||||
{
|
||||
type: "response.output_item.added",
|
||||
item: { type: "message", id: "msg_commentary", status: "in_progress", phase: "commentary" },
|
||||
},
|
||||
{ type: "response.output_text.delta", item_id: "msg_commentary", delta: "Checking." },
|
||||
{
|
||||
type: "response.output_item.done",
|
||||
item: {
|
||||
type: "message",
|
||||
id: "msg_commentary",
|
||||
status: "completed",
|
||||
phase: "commentary",
|
||||
content: [{ type: "output_text", text: "Checking.", annotations: [] }],
|
||||
},
|
||||
},
|
||||
{
|
||||
type: "response.output_item.added",
|
||||
item: { type: "message", id: "msg_final", status: "in_progress", phase: "final_answer" },
|
||||
},
|
||||
{ type: "response.output_text.done", item_id: "msg_final", content_index: 0, text: "Finished." },
|
||||
{
|
||||
type: "response.output_item.done",
|
||||
item: {
|
||||
type: "message",
|
||||
id: "msg_final",
|
||||
status: "incomplete",
|
||||
phase: "final_answer",
|
||||
content: [{ type: "output_text", text: "Finished.", annotations: [{ type: "test" }] }],
|
||||
},
|
||||
},
|
||||
{ type: "response.completed", response: { id: "resp_1" } },
|
||||
),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
expect(response.message.content).toEqual([
|
||||
{
|
||||
type: "text",
|
||||
text: "Checking.",
|
||||
providerMetadata: {
|
||||
openai: { itemId: "msg_commentary", status: "completed", phase: "commentary", annotations: [] },
|
||||
},
|
||||
},
|
||||
{
|
||||
type: "text",
|
||||
text: "Finished.",
|
||||
providerMetadata: {
|
||||
openai: {
|
||||
itemId: "msg_final",
|
||||
status: "incomplete",
|
||||
phase: "final_answer",
|
||||
annotations: [{ type: "test" }],
|
||||
},
|
||||
},
|
||||
},
|
||||
])
|
||||
|
||||
const prepared = yield* LLMClient.prepare<OpenAIResponses.OpenAIResponsesBody>(
|
||||
LLM.request({ model, messages: [response.message] }),
|
||||
)
|
||||
expect(prepared.body.input).toEqual([
|
||||
{
|
||||
type: "message",
|
||||
id: "msg_commentary",
|
||||
status: "completed",
|
||||
role: "assistant",
|
||||
phase: "commentary",
|
||||
content: [{ type: "output_text", text: "Checking.", annotations: [] }],
|
||||
},
|
||||
{
|
||||
type: "message",
|
||||
id: "msg_final",
|
||||
status: "incomplete",
|
||||
role: "assistant",
|
||||
phase: "final_answer",
|
||||
content: [{ type: "output_text", text: "Finished.", annotations: [{ type: "test" }] }],
|
||||
},
|
||||
])
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("parses reasoning summary stream fixtures", () =>
|
||||
Effect.gen(function* () {
|
||||
const body = sseEvents(
|
||||
@@ -1347,7 +1436,13 @@ describe("OpenAI Responses route", () => {
|
||||
},
|
||||
},
|
||||
},
|
||||
{ type: "text", text: "The parser changed." },
|
||||
{
|
||||
type: "text",
|
||||
text: "The parser changed.",
|
||||
providerMetadata: {
|
||||
openai: { itemId: "msg_1", phase: "final_answer", status: "completed" },
|
||||
},
|
||||
},
|
||||
]),
|
||||
Message.user("Summarize it."),
|
||||
],
|
||||
@@ -1358,7 +1453,7 @@ describe("OpenAI Responses route", () => {
|
||||
expect(prepared.body).toMatchObject({
|
||||
input: [
|
||||
{ role: "user", content: [{ type: "input_text", text: "What changed?" }] },
|
||||
{ role: "assistant", content: [{ type: "output_text", text: "The parser changed." }] },
|
||||
{ role: "assistant", content: "The parser changed.", phase: "final_answer" },
|
||||
{ role: "user", content: [{ type: "input_text", text: "Summarize it." }] },
|
||||
],
|
||||
store: false,
|
||||
|
||||
Reference in New Issue
Block a user