|
|
@@ -46,12 +46,8 @@ type OpenAIResponsesInputContent = Schema.Schema.Type<typeof OpenAIResponsesInpu
|
|
|
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 OpenAIResponsesReasoningSummaryText = Schema.Struct({
|
|
|
type: Schema.tag("summary_text"),
|
|
|
text: Schema.String,
|
|
|
@@ -82,19 +78,7 @@ const OpenAIResponsesFunctionCallOutput = Schema.Union([
|
|
|
const OpenAIResponsesInputItem = Schema.Union([
|
|
|
Schema.Struct({ role: Schema.tag("system"), content: Schema.String }),
|
|
|
Schema.Struct({ role: Schema.tag("user"), content: Schema.Array(OpenAIResponsesInputContent) }),
|
|
|
- Schema.Struct({
|
|
|
- role: Schema.tag("assistant"),
|
|
|
- content: Schema.String,
|
|
|
- phase: optionalNull(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: optionalNull(OpenAIResponsesMessagePhase),
|
|
|
- }),
|
|
|
+ Schema.Struct({ role: Schema.tag("assistant"), content: Schema.Array(OpenAIResponsesOutputText) }),
|
|
|
OpenAIResponsesReasoningItem,
|
|
|
OpenAIResponsesItemReference,
|
|
|
Schema.Struct({
|
|
|
@@ -210,9 +194,7 @@ const OpenAIResponsesStreamItem = Schema.Struct({
|
|
|
server_label: Schema.optional(Schema.String),
|
|
|
output: Schema.optional(Schema.Unknown),
|
|
|
error: Schema.optional(Schema.Unknown),
|
|
|
- content: Schema.optional(Schema.Array(Schema.Unknown)),
|
|
|
encrypted_content: optionalNull(Schema.String),
|
|
|
- phase: optionalNull(OpenAIResponsesMessagePhase),
|
|
|
})
|
|
|
type OpenAIResponsesStreamItem = Schema.Schema.Type<typeof OpenAIResponsesStreamItem>
|
|
|
|
|
|
@@ -230,9 +212,7 @@ const OpenAIResponsesErrorPayload = Schema.Struct({
|
|
|
const OpenAIResponsesEvent = Schema.Struct({
|
|
|
type: Schema.String,
|
|
|
delta: Schema.optional(Schema.String),
|
|
|
- text: Schema.optional(Schema.String),
|
|
|
item_id: Schema.optional(Schema.String),
|
|
|
- content_index: Schema.optional(Schema.Number),
|
|
|
summary_index: Schema.optional(Schema.Number),
|
|
|
item: Schema.optional(OpenAIResponsesStreamItem),
|
|
|
response: Schema.optional(
|
|
|
@@ -257,18 +237,10 @@ 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 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 {
|
|
|
@@ -326,26 +298,6 @@ const lowerReasoning = (part: ReasoningPart): OpenAIResponsesReasoningInput | un
|
|
|
}
|
|
|
}
|
|
|
|
|
|
-const messagePhase = (part: TextPart): OpenAIResponsesMessagePhase | null | undefined => {
|
|
|
- const phase = part.providerMetadata?.openai?.phase
|
|
|
- return phase === "commentary" || phase === "final_answer" || phase === null ? phase : undefined
|
|
|
-}
|
|
|
-
|
|
|
-const messageItemID = (part: TextPart) => {
|
|
|
- const itemID = part.providerMetadata?.openai?.itemId
|
|
|
- return typeof itemID === "string" && itemID.length > 0 ? itemID : undefined
|
|
|
-}
|
|
|
-
|
|
|
-const messageStatus = (part: TextPart) => {
|
|
|
- const status = part.providerMetadata?.openai?.status
|
|
|
- return status === "in_progress" || status === "completed" || status === "incomplete" ? status : undefined
|
|
|
-}
|
|
|
-
|
|
|
-const messageAnnotations = (part: TextPart) => {
|
|
|
- const annotations = part.providerMetadata?.openai?.annotations
|
|
|
- return Array.isArray(annotations) ? annotations : []
|
|
|
-}
|
|
|
-
|
|
|
const hostedToolItemID = (part: ToolResultPart) => {
|
|
|
const openai = part.providerMetadata?.openai
|
|
|
return ProviderShared.isRecord(openai) && typeof openai.itemId === "string" && openai.itemId.length > 0
|
|
|
@@ -416,49 +368,17 @@ const lowerMessages = Effect.fn("OpenAIResponses.lowerMessages")(function* (requ
|
|
|
}
|
|
|
|
|
|
if (message.role === "assistant") {
|
|
|
- const inputStart = input.length
|
|
|
const content: TextPart[] = []
|
|
|
- let phase: OpenAIResponsesMessagePhase | null | undefined
|
|
|
- let itemID: string | undefined
|
|
|
- let status: "in_progress" | "completed" | "incomplete" | undefined
|
|
|
const reasoningItems: Record<string, OpenAIResponsesReasoningReplay> = {}
|
|
|
const reasoningReferences = new Set<string>()
|
|
|
const hostedToolReferences = new Set<string>()
|
|
|
const flushText = () => {
|
|
|
if (content.length === 0) 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),
|
|
|
- })),
|
|
|
- ...(phase !== undefined ? { phase } : {}),
|
|
|
- }
|
|
|
- : {
|
|
|
- role: "assistant",
|
|
|
- content: ProviderShared.joinText(content),
|
|
|
- ...(phase !== undefined ? { phase } : {}),
|
|
|
- },
|
|
|
- )
|
|
|
+ input.push({ role: "assistant", content: content.map((part) => ({ type: "output_text", text: part.text })) })
|
|
|
content.splice(0, content.length)
|
|
|
- phase = undefined
|
|
|
- itemID = undefined
|
|
|
- status = undefined
|
|
|
}
|
|
|
for (const part of message.content) {
|
|
|
if (part.type === "text") {
|
|
|
- const nextPhase = messagePhase(part)
|
|
|
- const nextItemID = messageItemID(part)
|
|
|
- if (content.length > 0 && (phase !== nextPhase || itemID !== nextItemID)) flushText()
|
|
|
- phase = nextPhase
|
|
|
- itemID = nextItemID
|
|
|
- status = messageStatus(part) ?? status
|
|
|
content.push(part)
|
|
|
continue
|
|
|
}
|
|
|
@@ -509,20 +429,6 @@ const lowerMessages = Effect.fn("OpenAIResponses.lowerMessages")(function* (requ
|
|
|
])
|
|
|
}
|
|
|
flushText()
|
|
|
- if (store === false && Object.values(reasoningItems).some((item) => typeof item.encrypted_content !== "string"))
|
|
|
- input.splice(
|
|
|
- inputStart,
|
|
|
- input.length - inputStart,
|
|
|
- ...input.slice(inputStart).map((item) =>
|
|
|
- "type" in item && item.type === "message"
|
|
|
- ? {
|
|
|
- role: "assistant" as const,
|
|
|
- content: ProviderShared.joinText(item.content),
|
|
|
- ...(item.phase !== undefined ? { phase: item.phase } : {}),
|
|
|
- }
|
|
|
- : item,
|
|
|
- ),
|
|
|
- )
|
|
|
continue
|
|
|
}
|
|
|
|
|
|
@@ -706,134 +612,15 @@ const NO_EVENTS: StepResult["1"] = []
|
|
|
// the protocol's `terminal` predicate stay in sync.
|
|
|
const TERMINAL_TYPES = new Set(["response.completed", "response.incomplete", "response.failed"])
|
|
|
|
|
|
-const messageMetadata = (item: OpenAIResponsesStreamItem, id: string, previous?: ProviderMetadata) => {
|
|
|
- const openai = previous?.openai
|
|
|
- const phase = item.phase !== undefined ? item.phase : openai?.phase
|
|
|
- const status =
|
|
|
- item.status === "in_progress" || item.status === "completed" || item.status === "incomplete"
|
|
|
- ? item.status
|
|
|
- : openai?.status
|
|
|
- return openaiMetadata({
|
|
|
- itemId: id,
|
|
|
- ...(phase === "commentary" || phase === "final_answer" || phase === null ? { phase } : {}),
|
|
|
- ...(status === "in_progress" || status === "completed" || status === "incomplete" ? { status } : {}),
|
|
|
- })
|
|
|
-}
|
|
|
-
|
|
|
-const messageContentMetadata = (
|
|
|
- providerMetadata: ProviderMetadata,
|
|
|
- item: OpenAIResponsesStreamItem,
|
|
|
- index: number,
|
|
|
-): ProviderMetadata => {
|
|
|
- const content = item.content?.[index]
|
|
|
- if (!ProviderShared.isRecord(content) || content.type !== "output_text" || !Array.isArray(content.annotations))
|
|
|
- return providerMetadata
|
|
|
- return openaiMetadata({ ...providerMetadata.openai, annotations: content.annotations })
|
|
|
-}
|
|
|
-
|
|
|
-const ensureMessageContent = (state: ParserState, event: OpenAIResponsesEvent) => {
|
|
|
- const itemID = event.item_id ?? "text-0"
|
|
|
- const index = 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 = `openai-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: OpenAIResponsesEvent, text: string): StepResult => {
|
|
|
- const ensured = ensureMessageContent(state, event)
|
|
|
+const onOutputTextDelta = (state: ParserState, event: OpenAIResponsesEvent): StepResult => {
|
|
|
+ if (!event.delta) return [state, NO_EVENTS]
|
|
|
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),
|
|
|
- },
|
|
|
+ { ...state, lifecycle: Lifecycle.textDelta(state.lifecycle, events, event.item_id ?? "text-0", event.delta) },
|
|
|
events,
|
|
|
]
|
|
|
}
|
|
|
|
|
|
-const onOutputTextDelta = (state: ParserState, event: OpenAIResponsesEvent): StepResult => {
|
|
|
- if (!event.delta) return [state, NO_EVENTS]
|
|
|
- return appendOutputText(state, event, event.delta)
|
|
|
-}
|
|
|
-
|
|
|
-const onOutputTextDone = (state: ParserState, event: OpenAIResponsesEvent): StepResult => {
|
|
|
- if (event.text === undefined) return [state, NO_EVENTS]
|
|
|
- const ensured = ensureMessageContent(state, event)
|
|
|
- if (event.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 (event.text.startsWith(ensured.content.text))
|
|
|
- return appendOutputText(ensured.state, event, event.text.slice(ensured.content.text.length))
|
|
|
- return [ensured.state, NO_EVENTS]
|
|
|
-}
|
|
|
-
|
|
|
const onReasoningDelta = (state: ParserState, event: OpenAIResponsesEvent): StepResult => {
|
|
|
if (!event.delta) return [state, NO_EVENTS]
|
|
|
const events: LLMEvent[] = []
|
|
|
@@ -868,23 +655,6 @@ const reasoningMetadata = (item: OpenAIResponsesStreamItem & { id: string }) =>
|
|
|
// best-effort, not guaranteed.
|
|
|
const onOutputItemAdded = (state: ParserState, event: OpenAIResponsesEvent): StepResult => {
|
|
|
const item = event.item
|
|
|
- if (item?.type === "message" && item.id) {
|
|
|
- const existing = state.messageItems[item.id]
|
|
|
- return [
|
|
|
- {
|
|
|
- ...state,
|
|
|
- messageItems: {
|
|
|
- ...state.messageItems,
|
|
|
- [item.id]: {
|
|
|
- ...existing,
|
|
|
- providerMetadata: messageMetadata(item, item.id, existing?.providerMetadata),
|
|
|
- content: existing?.content ?? {},
|
|
|
- },
|
|
|
- },
|
|
|
- },
|
|
|
- NO_EVENTS,
|
|
|
- ]
|
|
|
- }
|
|
|
if (item && isReasoningItem(item)) {
|
|
|
const events: LLMEvent[] = []
|
|
|
return [
|
|
|
@@ -1042,32 +812,6 @@ const onOutputItemDone = Effect.fn("OpenAIResponses.onOutputItemDone")(function*
|
|
|
const item = event.item
|
|
|
if (!item) return [state, NO_EVENTS] satisfies StepResult
|
|
|
|
|
|
- if (item.type === "message" && item.id) {
|
|
|
- const events: LLMEvent[] = []
|
|
|
- const itemID = item.id
|
|
|
- const messageItem = state.messageItems[itemID]
|
|
|
- const { [itemID]: _finished, ...messageItems } = state.messageItems
|
|
|
- const providerMetadata = messageMetadata(item, itemID, messageItem?.providerMetadata)
|
|
|
- const lifecycle = Object.entries(messageItem?.content ?? {}).reduce(
|
|
|
- (lifecycle, entry) =>
|
|
|
- Lifecycle.textEnd(
|
|
|
- lifecycle,
|
|
|
- events,
|
|
|
- entry[1].id,
|
|
|
- messageContentMetadata(providerMetadata, item, Number(entry[0])),
|
|
|
- ),
|
|
|
- state.lifecycle,
|
|
|
- )
|
|
|
- return [
|
|
|
- {
|
|
|
- ...state,
|
|
|
- lifecycle,
|
|
|
- messageItems,
|
|
|
- },
|
|
|
- events,
|
|
|
- ] satisfies StepResult
|
|
|
- }
|
|
|
-
|
|
|
if (item.type === "function_call") {
|
|
|
if (!item.id || !item.call_id || !item.name) return [state, NO_EVENTS] satisfies StepResult
|
|
|
const tools = state.tools[item.id]
|
|
|
@@ -1195,7 +939,6 @@ const step = (state: ParserState, event: OpenAIResponsesEvent) => {
|
|
|
if (event.type === "response.reasoning_summary_part.done")
|
|
|
return Effect.succeed(onReasoningSummaryPartDone(state, event))
|
|
|
if (event.type === "response.output_item.added") return Effect.succeed(onOutputItemAdded(state, event))
|
|
|
- if (event.type === "response.output_text.done") return Effect.succeed(onOutputTextDone(state, event))
|
|
|
if (event.type === "response.function_call_arguments.delta") return onFunctionCallArgumentsDelta(state, event)
|
|
|
if (event.type === "response.output_item.done") return onOutputItemDone(state, event)
|
|
|
if (event.type === "response.completed" || event.type === "response.incomplete")
|
|
|
@@ -1225,9 +968,6 @@ export const protocol = Protocol.make({
|
|
|
hasFunctionCall: false,
|
|
|
tools: ToolStream.empty<string>(),
|
|
|
lifecycle: Lifecycle.initial(),
|
|
|
- messageItems: {},
|
|
|
- messageContentIDs: new Set<string>(),
|
|
|
- nextMessageContentID: 0,
|
|
|
reasoningItems: {},
|
|
|
store: OpenAIOptions.store(request),
|
|
|
}),
|