fix(ai): reconcile responses snapshots

This commit is contained in:
Aiden Cline 2026-07-25 14:37:30 -05:00
commit 02e8f398b8
4 changed files with 328 additions and 37 deletions

View file

@ -193,6 +193,14 @@ const OpenResponsesUsage = Schema.Struct({
})
type OpenResponsesUsage = Schema.Schema.Type<typeof OpenResponsesUsage>
const StreamContent = Schema.StructWithRest(
Schema.Struct({
type: Schema.String,
text: Schema.optional(Schema.String),
}),
[Schema.Record(Schema.String, Schema.Unknown)],
)
export const StreamItem = Schema.StructWithRest(
Schema.Struct({
type: Schema.String,
@ -201,6 +209,8 @@ export const StreamItem = Schema.StructWithRest(
name: Schema.optional(Schema.String),
arguments: Schema.optional(Schema.String),
encrypted_content: optionalNull(Schema.String),
content: optionalArray(StreamContent),
summary: optionalArray(StreamContent),
}),
[Schema.Record(Schema.String, Schema.Unknown)],
)
@ -221,7 +231,10 @@ export const Event = Schema.StructWithRest(
type: Schema.String,
delta: Schema.optional(Schema.String),
text: Schema.optional(Schema.String),
arguments: Schema.optional(Schema.String),
item_id: Schema.optional(Schema.String),
output_index: Schema.optional(Schema.Number),
content_index: Schema.optional(Schema.Number),
summary_index: Schema.optional(Schema.Number),
item: Schema.optional(StreamItem),
response: Schema.optional(
@ -232,6 +245,7 @@ export const Event = Schema.StructWithRest(
incomplete_details: optionalNull(Schema.Struct({ reason: Schema.optional(Schema.String) })),
usage: optionalNull(OpenResponsesUsage),
error: optionalNull(OpenResponsesErrorPayload),
output: optionalArray(StreamItem),
}),
[Schema.Record(Schema.String, Schema.Unknown)],
),
@ -268,6 +282,8 @@ export interface ParserState {
readonly messageItems: ReadonlySet<string>
readonly messagePhase: (value: unknown) => MessagePhase | null | undefined
readonly messagePhases: Readonly<Record<string, MessagePhase | null>>
readonly outputText: Readonly<Record<string, string>>
readonly reasoningText: Readonly<Record<string, string>>
readonly reasoningItems: Readonly<Record<string, ReasoningStreamItem>>
readonly store: boolean | undefined
}
@ -647,19 +663,43 @@ const onOutputTextDelta = (state: ParserState, event: Event, id: string): StepRe
const metadata = phase === undefined ? undefined : providerMetadata(state, { phase })
const lifecycle = Lifecycle.textStart(state.lifecycle, events, id, metadata)
return [
{ ...state, lifecycle: Lifecycle.textDelta(lifecycle, events, id, event.delta) },
{
...state,
lifecycle: Lifecycle.textDelta(lifecycle, events, id, event.delta),
outputText: { ...state.outputText, [id]: `${state.outputText[id] ?? ""}${event.delta}` },
},
events,
]
}
const onOutputTextDone = (state: ParserState, event: Event, id: string): StepResult => {
if (state.messageItems.has(id)) {
if (state.lifecycle.text.has(id) || event.text === undefined) return [state, NO_EVENTS]
return onOutputTextDelta(state, { ...event, delta: event.text }, id)
}
const authoritativeSuffix = Effect.fn("OpenResponses.authoritativeSuffix")(function* (
state: ParserState,
kind: string,
id: string,
current: string,
value: string,
) {
if (value.startsWith(current)) return value.slice(current.length)
return yield* ProviderShared.eventError(
state.id,
`${kind} ${id} completed with content that conflicts with its streamed deltas`,
value,
)
})
const onOutputTextDone = Effect.fn("OpenResponses.onOutputTextDone")(function* (
state: ParserState,
event: Event,
id: string,
) {
const suffix =
event.text === undefined ? "" : yield* authoritativeSuffix(state, "output text", id, state.outputText[id] ?? "", event.text)
const [reconciled, deltaEvents] = onOutputTextDelta(state, { ...event, delta: suffix }, id)
if (reconciled.messageItems.has(id)) return [reconciled, deltaEvents] satisfies StepResult
const events: LLMEvent[] = []
return [{ ...state, lifecycle: Lifecycle.textEnd(state.lifecycle, events, id) }, events]
}
const lifecycle = Lifecycle.textEnd(reconciled.lifecycle, events, id)
return [{ ...reconciled, lifecycle }, [...deltaEvents, ...events]] satisfies StepResult
})
export const onReasoningDelta = (state: ParserState, event: Event, itemID: string): StepResult => {
if (!event.delta) return [state, NO_EVENTS]
@ -670,12 +710,39 @@ export const onReasoningDelta = (state: ParserState, event: Event, itemID: strin
{
...state,
lifecycle: Lifecycle.reasoningDelta(state.lifecycle, events, id, event.delta),
reasoningText: { ...state.reasoningText, [id]: `${state.reasoningText[id] ?? ""}${event.delta}` },
},
events,
]
}
export const onReasoningDone = (state: ParserState, _event: Event): StepResult => [state, NO_EVENTS]
export const onReasoningDone = Effect.fn("OpenResponses.onReasoningDone")(function* (
state: ParserState,
event: Event,
) {
if (!event.item_id || event.text === undefined) return [state, NO_EVENTS] satisfies StepResult
const id =
event.summary_index !== undefined || state.reasoningItems[event.item_id]
? `${event.item_id}:${event.summary_index ?? 0}`
: event.item_id
const suffix = yield* authoritativeSuffix(state, "reasoning", id, state.reasoningText[id] ?? "", event.text)
const [reconciled, events] = onReasoningDelta(state, { ...event, delta: suffix }, event.item_id)
const item = reconciled.reasoningItems[event.item_id]
if (!suffix || !item || event.summary_index === undefined) return [reconciled, events] satisfies StepResult
return [
{
...reconciled,
reasoningItems: {
...reconciled.reasoningItems,
[event.item_id]: {
...item,
summaryParts: { ...item.summaryParts, [event.summary_index]: "active" },
},
},
},
events,
] satisfies StepResult
})
const reasoningMetadata = (state: ParserState, item: StreamItem & { id: string }) =>
providerMetadata(state, { itemId: item.id, reasoningEncryptedContent: item.encrypted_content ?? null })
@ -865,23 +932,32 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
if (item.type === "message" && item.id) {
const itemPhase = state.messagePhase(item.phase)
const phase = itemPhase === undefined ? state.messagePhases[item.id] : itemPhase
const messageItems = new Set([...state.messageItems, item.id])
const messagePhases = phase === undefined ? state.messagePhases : { ...state.messagePhases, [item.id]: phase }
const text = item.content
?.filter((part) => part.type === "output_text" && part.text !== undefined)
.map((part) => part.text)
.join("")
const [reconciled, reconciledEvents] =
text === undefined
? [{ ...state, messageItems, messagePhases }, NO_EVENTS]
: yield* onOutputTextDone({ ...state, messageItems, messagePhases }, { ...event, text }, item.id)
const events: LLMEvent[] = []
const messageItems = new Set(state.messageItems)
messageItems.delete(item.id)
const { [item.id]: _phase, ...messagePhases } = state.messagePhases
const { [item.id]: _phase, ...remainingPhases } = reconciled.messagePhases
return [
{
...state,
...reconciled,
lifecycle: Lifecycle.textEnd(
state.lifecycle,
reconciled.lifecycle,
events,
item.id,
phase === undefined ? undefined : providerMetadata(state, { phase }),
),
messageItems,
messagePhases,
messagePhases: remainingPhases,
},
events,
[...reconciledEvents, ...events],
] satisfies StepResult
}
@ -914,25 +990,75 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
if (isReasoningItem(item)) {
const events: LLMEvent[] = []
const metadata = reasoningMetadata(state, item)
const reasoningItem = state.reasoningItems[item.id]
const summaries = item.summary?.filter((part) => part.type === "summary_text" && part.text !== undefined) ?? []
if (summaries.length === 1 && !state.reasoningItems[item.id] && state.reasoningText[item.id] !== undefined) {
const [reconciled, reconciledEvents] = yield* onReasoningDone(state, {
...event,
item_id: item.id,
text: summaries[0]?.text,
})
events.push(...reconciledEvents)
return [
{ ...reconciled, lifecycle: Lifecycle.reasoningEnd(reconciled.lifecycle, events, item.id, metadata) },
events,
] satisfies StepResult
}
const needsSummaryReconciliation = summaries.some(
(part, index) => part.text !== state.reasoningText[`${item.id}:${index}`],
)
const seeded =
needsSummaryReconciliation && !state.reasoningItems[item.id]
? onOutputItemAdded(state, { ...event, item })
: ([state, NO_EVENTS] satisfies StepResult)
events.push(...seeded[1])
const reconciled = yield* (needsSummaryReconciliation ? summaries : [])
.map((part, index) => ({ part, index }))
.reduce<Effect.Effect<ParserState, LLMError>>(
(effect, entry) =>
effect.pipe(
Effect.flatMap(
Effect.fnUntraced(function* (current) {
const added = current.reasoningItems[item.id]?.summaryParts[entry.index]
? ([current, NO_EVENTS] satisfies StepResult)
: onReasoningSummaryPartAdded(current, {
...event,
item_id: item.id,
summary_index: entry.index,
})
events.push(...added[1])
const [next, nextEvents] = yield* onReasoningDone(added[0], {
...event,
item_id: item.id,
summary_index: entry.index,
text: entry.part.text,
})
events.push(...nextEvents)
return next
}),
),
),
Effect.succeed(seeded[0]),
)
const reasoningItem = reconciled.reasoningItems[item.id]
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),
state.lifecycle,
reconciled.lifecycle,
)
const { [item.id]: _removed, ...reasoningItems } = state.reasoningItems
return [{ ...state, lifecycle, reasoningItems }, events] satisfies StepResult
const { [item.id]: _removed, ...reasoningItems } = reconciled.reasoningItems
return [{ ...reconciled, lifecycle, reasoningItems }, events] satisfies StepResult
}
if (!state.lifecycle.reasoning.has(item.id)) {
const lifecycle = Lifecycle.stepStart(state.lifecycle, events)
if (summaries.length) return [reconciled, events] satisfies StepResult
if (!reconciled.lifecycle.reasoning.has(item.id)) {
const lifecycle = Lifecycle.stepStart(reconciled.lifecycle, events)
events.push(LLMEvent.reasoningStart({ id: item.id, providerMetadata: metadata }))
events.push(LLMEvent.reasoningEnd({ id: item.id, providerMetadata: metadata }))
return [{ ...state, lifecycle }, events] satisfies StepResult
return [{ ...reconciled, lifecycle }, events] satisfies StepResult
}
return [
{ ...state, lifecycle: Lifecycle.reasoningEnd(state.lifecycle, events, item.id, metadata) },
{ ...reconciled, lifecycle: Lifecycle.reasoningEnd(reconciled.lifecycle, events, item.id, metadata) },
events,
] satisfies StepResult
}
@ -940,11 +1066,24 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* (
return [state, NO_EVENTS] satisfies StepResult
})
const onResponseFinish = (state: ParserState, event: Event): StepResult => {
const onResponseFinish = Effect.fn("OpenResponses.onResponseFinish")(function* (state: ParserState, event: Event) {
const events: LLMEvent[] = []
const lifecycle = Lifecycle.finish(state.lifecycle, events, {
const reconciled = yield* (event.response?.output ?? []).reduce<Effect.Effect<ParserState, LLMError>>(
(effect, item) =>
effect.pipe(
Effect.flatMap(
Effect.fnUntraced(function* (current) {
const [next, nextEvents] = yield* onOutputItemDone(current, { ...event, item })
events.push(...nextEvents)
return next
}),
),
),
Effect.succeed(state),
)
const lifecycle = Lifecycle.finish(reconciled.lifecycle, events, {
reason: {
normalized: mapFinishReason(event, state.hasFunctionCall),
normalized: mapFinishReason(event, reconciled.hasFunctionCall),
raw: event.response?.incomplete_details?.reason,
},
usage: mapUsage(event.response?.usage, state.providerMetadataKey),
@ -956,8 +1095,8 @@ const onResponseFinish = (state: ParserState, event: Event): StepResult => {
})
: undefined,
})
return [{ ...state, lifecycle }, events]
}
return [{ ...reconciled, lifecycle }, events] satisfies StepResult
})
// Build a single human-readable message from whatever the provider supplied.
// When both code and message are present, prefix the code so consumers see
@ -985,11 +1124,9 @@ const providerError = (state: ParserState, event: Event, fallback: string) => {
export const step = (state: ParserState, event: Event) => {
if (event.type === "response.output_text.delta" || event.type === "response.output_text.done") {
if (!event.item_id) return ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
return Effect.succeed(
event.type === "response.output_text.delta"
? onOutputTextDelta(state, event, event.item_id)
: onOutputTextDone(state, event, event.item_id),
)
return event.type === "response.output_text.delta"
? Effect.succeed(onOutputTextDelta(state, event, event.item_id))
: onOutputTextDone(state, event, event.item_id)
}
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`)
@ -997,7 +1134,7 @@ export const step = (state: ParserState, event: Event) => {
}
if (event.type === "response.reasoning.done" || event.type === "response.reasoning_summary_text.done") {
if (!event.item_id) return ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
return Effect.succeed(onReasoningDone(state, event))
return onReasoningDone(state, event)
}
if (event.type === "response.reasoning_summary_part.added")
return event.item_id
@ -1013,13 +1150,29 @@ export const step = (state: ParserState, event: Event) => {
return Effect.succeed(onOutputItemAdded(state, event))
}
if (event.type === "response.function_call_arguments.delta") return onFunctionCallArgumentsDelta(state, event)
if (event.type === "response.function_call_arguments.done") {
if (!event.item_id) return ProviderShared.eventError(state.id, `${event.type} is missing item_id`)
if (event.arguments === undefined) return ProviderShared.eventError(state.id, `${event.type} is missing arguments`)
const tool = state.tools[event.item_id]
if (!tool)
return ProviderShared.eventError(state.id, `${state.name} completed tool arguments are missing their tool call`)
return Effect.succeed(
[
{
...state,
tools: { ...state.tools, [event.item_id]: { ...tool, input: event.arguments } },
},
NO_EVENTS,
] satisfies StepResult,
)
}
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`)
return onOutputItemDone(state, event)
}
if (event.type === "response.completed" || event.type === "response.incomplete")
return Effect.succeed(onResponseFinish(state, event))
return onResponseFinish(state, event)
if (event.type === "response.failed") return providerError(state, event, `${state.name} response failed`)
if (event.type === "error") return providerError(state, event, `${state.name} stream error`)
return Effect.succeed<StepResult>([state, NO_EVENTS])
@ -1042,6 +1195,8 @@ export const initial = (request: LLMRequest, extension: Extension = BASE): Parse
messageItems: new Set<string>(),
messagePhase: (value) => messagePhase(value, extension),
messagePhases: {},
outputText: {},
reasoningText: {},
reasoningItems: {},
store: OpenResponsesOptions.resolve(request).store,
})

View file

@ -224,7 +224,7 @@ const step = (state: OpenResponses.ParserState, event: OpenResponses.Event) => {
: ProviderShared.eventError(ADAPTER, `${event.type} is missing item_id`)
if (event.type === "response.reasoning_text.done" || event.type === "response.reasoning_summary.done")
return event.item_id
? Effect.succeed(OpenResponses.onReasoningDone(state, event))
? OpenResponses.onReasoningDone(state, event)
: ProviderShared.eventError(ADAPTER, `${event.type} is missing item_id`)
if (event.type === "response.output_item.done" && event.item && isHostedToolItem(event.item))
return onHostedToolDone(state, event.item)

View file

@ -882,6 +882,87 @@ describe("OpenAI Responses route", () => {
}),
)
it.effect("reconciles authoritative output text done values", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
Effect.provide(
fixedResponse(
sseEvents(
{ type: "response.output_item.added", item: { type: "message", id: "msg_1" } },
{ type: "response.output_text.delta", item_id: "msg_1", delta: "Hello" },
{ type: "response.output_text.done", item_id: "msg_1", text: "Hello!" },
{
type: "response.output_item.done",
item: {
type: "message",
id: "msg_1",
content: [{ type: "output_text", text: "Hello!" }],
},
},
{ type: "response.completed", response: { id: "resp_1" } },
),
),
),
)
expect(response.text).toBe("Hello!")
expect(response.events.filter(LLMEvent.is.textDelta)).toEqual([
{ type: "text-delta", id: "msg_1", text: "Hello" },
{ type: "text-delta", id: "msg_1", text: "!" },
])
}),
)
it.effect("recovers output from the authoritative terminal response", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
Effect.provide(
fixedResponse(
sseEvents(
{
type: "response.completed",
response: {
id: "resp_1",
output: [
{
type: "message",
id: "msg_1",
content: [{ type: "output_text", text: "Terminal only." }],
},
],
},
},
),
),
),
)
expect(response.text).toBe("Terminal only.")
expect(response.events.filter(LLMEvent.is.textDelta)).toEqual([
{ type: "text-delta", id: "msg_1", text: "Terminal only." },
])
}),
)
it.effect("rejects authoritative text that conflicts with streamed deltas", () =>
Effect.gen(function* () {
const error = yield* LLMClient.generate(request).pipe(
Effect.provide(
fixedResponse(
sseEvents(
{ type: "response.output_text.delta", item_id: "msg_1", delta: "Hello" },
{ type: "response.output_text.done", item_id: "msg_1", text: "Goodbye" },
),
),
),
Effect.flip,
)
expect(error.reason._tag).toBe("InvalidProviderOutput")
expect(error.message).toContain("completed with content that conflicts with its streamed deltas")
}),
)
it.effect("preserves and replays assistant message phases", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
@ -1086,6 +1167,24 @@ describe("OpenAI Responses route", () => {
}),
)
it.effect("reconciles authoritative reasoning done values", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
Effect.provide(
fixedResponse(
sseEvents(
{ type: "response.reasoning_summary_text.delta", item_id: "rs_1", delta: "thinking" },
{ type: "response.reasoning_summary_text.done", item_id: "rs_1", text: "thinking..." },
{ type: "response.completed", response: { id: "resp_1" } },
),
),
),
)
expect(response.reasoning).toBe("thinking...")
}),
)
it.effect("preserves encrypted reasoning metadata for continuation", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
@ -1568,6 +1667,43 @@ describe("OpenAI Responses route", () => {
}),
)
it.effect("uses authoritative function argument done values", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(
LLMRequest.update(request, {
tools: [ToolDefinition.make({ name: "lookup", description: "Lookup data", inputSchema: { type: "object" } })],
}),
).pipe(
Effect.provide(
fixedResponse(
sseEvents(
{
type: "response.output_item.added",
item: { type: "function_call", id: "item_1", call_id: "call_1", name: "lookup", arguments: "" },
},
{ type: "response.function_call_arguments.delta", item_id: "item_1", delta: '{"query":"stale"}' },
{
type: "response.function_call_arguments.done",
item_id: "item_1",
arguments: '{"query":"authoritative"}',
},
{
type: "response.output_item.done",
item: { type: "function_call", id: "item_1", call_id: "call_1", name: "lookup" },
},
{ type: "response.completed", response: {} },
),
),
),
)
expect(response.events.find(LLMEvent.is.toolCall)).toMatchObject({
id: "call_1",
input: { query: "authoritative" },
})
}),
)
it.effect("emits malformed final function arguments as an unexecuted tool error", () =>
Effect.gen(function* () {
const body = sseEvents(