Compare commits

...
Author SHA1 Message Date
Aiden Cline 1ac220cf98 refactor(ai): track reasoning by output index 2026-08-24 22:03:12 -05:00
Aiden Cline 4d144edc80 refactor(ai): clarify reasoning id recovery 2026-08-24 18:22:45 -05:00
Aiden Cline 183971ef3a fix(ai): recover missing reasoning item ids
Keep reasoning streams coherent when compatible providers omit item IDs. Synthetic IDs remain internal, while provider-issued IDs are retained for safe replay when they arrive later.
2026-08-24 17:27:10 -05:00
2 changed files with 186 additions and 27 deletions

No files matched your search

+92 -27
View File
@@ -287,6 +287,7 @@ export const Event = Schema.StructWithRest(
delta: Schema.optional(Schema.String),
text: Schema.optional(Schema.String),
item_id: Schema.optional(Schema.String),
output_index: Schema.optional(Schema.Number),
summary_index: Schema.optional(Schema.Number),
item: Schema.optional(StreamItem),
response: Schema.optional(
@@ -348,6 +349,8 @@ export interface ParserState {
type ReasoningSummaryStatus = "active" | "can-conclude" | "concluded"
interface ReasoningStreamItem {
readonly providerItemID: string | undefined
readonly outputIndex: number | undefined
readonly encryptedContent: string | null | undefined
// Keyed by the wire protocol's numeric `summary_index`. JS object keys coerce to
// strings, but typing the map as `Record<number, ...>` documents intent
@@ -845,8 +848,32 @@ export const onReasoningDone = (state: ParserState, event: Event, itemID: string
return onReasoningDelta(state, { ...event, delta: event.text }, itemID)
}
const reasoningMetadata = (state: ParserState, item: StreamItem & { id: string }) =>
providerMetadata(state, { itemId: item.id, reasoningEncryptedContent: item.encrypted_content ?? null })
const reasoningMetadata = (state: ParserState, item: StreamItem) =>
providerMetadata(state, {
...(item.id ? { itemId: item.id } : {}),
reasoningEncryptedContent: item.encrypted_content ?? null,
})
const reasoningStreamMetadata = (state: ParserState, item: ReasoningStreamItem, includeEncrypted: boolean) => {
if (!item.providerItemID && !includeEncrypted) return undefined
return providerMetadata(state, {
...(item.providerItemID ? { itemId: item.providerItemID } : {}),
...(includeEncrypted ? { reasoningEncryptedContent: item.encryptedContent ?? null } : {}),
})
}
const findReasoningItem = (state: ParserState, event: Event, itemID: string | undefined) => {
if (itemID && state.reasoningItems[itemID]) return [itemID, state.reasoningItems[itemID]] as const
const items = Object.entries(state.reasoningItems)
if (itemID) {
const known = items.find((entry) => entry[1].providerItemID === itemID)
if (known) return known
}
if (event.output_index === undefined) return undefined
const indexed = items.find((entry) => entry[1].outputIndex === event.output_index)
if (indexed?.[1].providerItemID && itemID && indexed[1].providerItemID !== itemID) return undefined
return indexed
}
// Responses APIs stream reasoning items in a stable order:
// `output_item.added` (reasoning) →
@@ -873,15 +900,18 @@ const onOutputItemAdded = (state: ParserState, event: Event): StepResult => {
NO_EVENTS,
]
}
if (item && isReasoningItem(item)) {
if (item?.type === "reasoning") {
const id = item.id || `reasoning:${event.output_index}`
const events: LLMEvent[] = []
return [
{
...state,
lifecycle: Lifecycle.reasoningStart(state.lifecycle, events, `${item.id}:0`, reasoningMetadata(state, item)),
lifecycle: Lifecycle.reasoningStart(state.lifecycle, events, `${id}:0`, reasoningMetadata(state, item)),
reasoningItems: {
...state.reasoningItems,
[item.id]: {
[id]: {
providerItemID: item.id,
outputIndex: event.output_index,
encryptedContent: item.encrypted_content,
summaryParts: { 0: "active" },
deltaIndexes: new Set(),
@@ -928,7 +958,7 @@ const onReasoningSummaryPartAdded = (state: ParserState, event: Event): StepResu
lifecycle,
events,
`${event.item_id}:${entry[0]}`,
providerMetadata(state, { itemId: event.item_id }),
reasoningStreamMetadata(state, item, false),
),
state.lifecycle,
)
@@ -939,7 +969,7 @@ const onReasoningSummaryPartAdded = (state: ParserState, event: Event): StepResu
closed,
events,
`${event.item_id}:${event.summary_index}`,
providerMetadata(state, { itemId: event.item_id, reasoningEncryptedContent: item.encryptedContent ?? null }),
reasoningStreamMetadata(state, item, true),
),
reasoningItems: {
...state.reasoningItems,
@@ -974,7 +1004,7 @@ const onReasoningSummaryPartDone = (state: ParserState, event: Event): StepResul
state.lifecycle,
events,
`${event.item_id}:${event.summary_index}`,
providerMetadata(state, { itemId: event.item_id }),
reasoningStreamMetadata(state, item, false),
)
: state.lifecycle,
reasoningItems: {
@@ -1068,28 +1098,35 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
] satisfies StepResult
}
if (isReasoningItem(item)) {
if (item?.type === "reasoning") {
const id = findReasoningItem(state, event, item.id)?.[0] ?? item.id ?? `reasoning:${event.output_index}`
const events: LLMEvent[] = []
const metadata = reasoningMetadata(state, item)
const reasoningItem = state.reasoningItems[item.id]
const reasoningItem = state.reasoningItems[id]
const metadata = reasoningMetadata(state, {
...item,
id: item.id || reasoningItem?.providerItemID,
})
if (reasoningItem) {
const lifecycle = Object.entries(reasoningItem.summaryParts)
.filter((entry) => entry[1] === "active" || entry[1] === "can-conclude")
.reduce(
(lifecycle, entry) => Lifecycle.reasoningEnd(lifecycle, events, `${item.id}:${entry[0]}`, metadata),
(lifecycle, entry) => Lifecycle.reasoningEnd(lifecycle, events, `${id}:${entry[0]}`, metadata),
state.lifecycle,
)
const { [item.id]: _removed, ...reasoningItems } = state.reasoningItems
const { [id]: _removed, ...reasoningItems } = state.reasoningItems
return [{ ...state, lifecycle, reasoningItems }, events] satisfies StepResult
}
if (!state.lifecycle.reasoning.has(item.id)) {
if (!state.lifecycle.reasoning.has(id)) {
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 }))
events.push(LLMEvent.reasoningStart({ id, providerMetadata: metadata }))
events.push(LLMEvent.reasoningEnd({ id, providerMetadata: metadata }))
return [{ ...state, lifecycle }, events] satisfies StepResult
}
return [
{ ...state, lifecycle: Lifecycle.reasoningEnd(state.lifecycle, events, item.id, metadata) },
{
...state,
lifecycle: Lifecycle.reasoningEnd(state.lifecycle, events, id, metadata),
},
events,
] satisfies StepResult
}
@@ -1124,6 +1161,24 @@ const onResponseFinish = Effect.fn("OpenResponses.onResponseFinish")(function* (
return [{ ...state, lifecycle, hasFunctionCall, tools: pending.tools }, events] satisfies StepResult
})
const normalizeReasoningEvent = (state: ParserState, event: Event) => {
const entry = findReasoningItem(state, event, event.item_id)
if (!entry) return [state, event] as const
const id = entry[0]
const item = entry[1]
if (!event.item_id || item.providerItemID) return [state, { ...event, item_id: id }] as const
return [
{
...state,
reasoningItems: {
...state.reasoningItems,
[id]: { ...item, providerItemID: event.item_id },
},
},
{ ...event, item_id: id },
] as const
}
// Build the prettiest summary available from whatever the provider supplied.
// When both code and message are present, prefix the code so consumers see
// the failure mode (e.g. `rate_limit_exceeded: Slow down`) instead of just
@@ -1188,28 +1243,36 @@ export const step = (state: ParserState, event: Event) => {
)
}
if (event.type === "response.reasoning.delta" || event.type === "response.reasoning_summary_text.delta") {
if (!event.item_id) return ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
return Effect.succeed(onReasoningDelta(state, event, event.item_id))
const [next, normalized] = normalizeReasoningEvent(state, event)
if (!normalized.item_id) return ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
return Effect.succeed(onReasoningDelta(next, normalized, normalized.item_id))
}
if (
event.type === "response.reasoning.done" ||
event.type === "response.reasoning_summary_text.done" ||
event.type === "response.reasoning_text.done"
) {
if (!event.item_id) return ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
return Effect.succeed(onReasoningDone(state, event, event.item_id))
const [next, normalized] = normalizeReasoningEvent(state, event)
if (!normalized.item_id) return ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
return Effect.succeed(onReasoningDone(next, normalized, normalized.item_id))
}
if (event.type === "response.reasoning_summary_part.added")
return event.item_id
? Effect.succeed(onReasoningSummaryPartAdded(state, event))
if (event.type === "response.reasoning_summary_part.added") {
const [next, normalized] = normalizeReasoningEvent(state, event)
return normalized.item_id
? Effect.succeed(onReasoningSummaryPartAdded(next, normalized))
: ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
if (event.type === "response.reasoning_summary_part.done")
return event.item_id
? Effect.succeed(onReasoningSummaryPartDone(state, event))
}
if (event.type === "response.reasoning_summary_part.done") {
const [next, normalized] = normalizeReasoningEvent(state, event)
return normalized.item_id
? Effect.succeed(onReasoningSummaryPartDone(next, normalized))
: ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
}
if (event.type === "response.output_item.added") {
if (event.item?.type === "message" && !event.item.id)
return ProviderShared.eventError(state.id, `${event.type} message is missing id`)
if (event.item?.type === "reasoning" && !event.item.id && event.output_index === undefined)
return ProviderShared.eventError(state.id, `${event.type} reasoning is missing id and output_index`)
return Effect.succeed(onOutputItemAdded(state, event))
}
if (event.type === "response.function_call_arguments.delta")
@@ -1219,6 +1282,8 @@ export const step = (state: ParserState, event: Event) => {
if (event.type === "response.output_item.done") {
if (event.item?.type === "message" && !event.item.id)
return ProviderShared.eventError(state.id, `${event.type} message is missing id`)
if (event.item?.type === "reasoning" && !event.item.id && event.output_index === undefined)
return ProviderShared.eventError(state.id, `${event.type} reasoning is missing id and output_index`)
return onOutputItemDone(state, event)
}
if (event.type === "response.completed" || event.type === "response.incomplete") return onResponseFinish(state, event)
@@ -2133,6 +2133,100 @@ describe("OpenAI Responses route", () => {
}),
)
it.effect("tracks reasoning by output index when events omit item IDs", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
Effect.provide(
fixedResponse(
sseEvents(
{ type: "response.output_item.added", output_index: 0, item: { type: "reasoning" } },
{ type: "response.reasoning_summary_text.delta", output_index: 0, delta: "thinking" },
{
type: "response.output_item.done",
output_index: 0,
item: {
type: "reasoning",
encrypted_content: "encrypted-state",
summary: [{ type: "summary_text", text: "thinking" }],
},
},
{ type: "response.completed", response: { id: "resp_1" } },
),
),
),
)
const start = response.events.find(LLMEvent.is.reasoningStart)
const delta = response.events.find(LLMEvent.is.reasoningDelta)
const end = response.events.find(LLMEvent.is.reasoningEnd)
expect(start?.id).toBe("reasoning:0:0")
expect(delta?.id).toBe(start?.id)
expect(end?.id).toBe(start?.id)
expect(response.message.content).toEqual([
{
type: "reasoning",
text: "thinking",
providerMetadata: {
openai: {
reasoningEncryptedContent: "encrypted-state",
},
},
},
])
}),
)
it.effect("retains a provider reasoning ID that arrives after item start", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(
LLMRequest.update(request, { providerOptions: { store: false } }),
).pipe(
Effect.provide(
fixedResponse(
sseEvents(
{ type: "response.output_item.added", output_index: 0, item: { type: "reasoning" } },
{ type: "response.reasoning_summary_part.added", output_index: 0, item_id: "rs_real", summary_index: 0 },
{
type: "response.reasoning_summary_text.delta",
item_id: "rs_real",
summary_index: 0,
delta: "thinking",
},
{ type: "response.reasoning_summary_part.done", item_id: "rs_real", summary_index: 0 },
{
type: "response.output_item.done",
output_index: 0,
item: {
type: "reasoning",
encrypted_content: "encrypted-state",
summary: [{ type: "summary_text", text: "thinking" }],
},
},
{ type: "response.completed", response: { id: "resp_1" } },
),
),
),
)
const start = response.events.find(LLMEvent.is.reasoningStart)
expect(start?.id).toBe("reasoning:0:0")
expect(response.reasoning).toBe("thinking")
expect(response.message.content).toEqual([
{
type: "reasoning",
text: "thinking",
providerMetadata: {
openai: {
itemId: "rs_real",
reasoningEncryptedContent: "encrypted-state",
},
},
},
])
}),
)
it.effect("preserves encrypted reasoning metadata for continuation", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(