From cc0cc59700c7db27c12c4bd9b2f4f4737c848b30 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 31 Aug 2026 19:24:33 -0400 Subject: [PATCH] fix(ai): preserve done-only response messages (#46064) --- packages/ai/src/protocols/open-responses.ts | 22 ++++- .../provider/open-responses-finals.test.ts | 26 +++++ .../provider/open-responses-lifecycle.test.ts | 97 ++++++++++++++++++- 3 files changed, 138 insertions(+), 7 deletions(-) diff --git a/packages/ai/src/protocols/open-responses.ts b/packages/ai/src/protocols/open-responses.ts index 20ce29b6e5d..232eafb3e46 100644 --- a/packages/ai/src/protocols/open-responses.ts +++ b/packages/ai/src/protocols/open-responses.ts @@ -397,6 +397,9 @@ export interface ParserState { readonly lifecycle: Lifecycle.State readonly outputItems: Readonly> readonly message: { readonly id: string; readonly phase: MessagePhase | null | undefined } | undefined + // Item ids are response-scoped identities. Keep completed ids tombstoned so + // reconnect replay cannot reopen fragments already emitted downstream. + readonly completedMessages: ReadonlySet readonly reasoningItems: Readonly> } @@ -952,12 +955,16 @@ const onOutputItemAdded = (state: ParserState, event: Event): StepResult => { const item = event.item if (item?.type === "message" && item.id !== undefined) { const itemID = item.id + if (state.completedMessages.has(itemID)) return [state, NO_EVENTS] const phase = messagePhase(item.phase) + const completedMessages = new Set(state.completedMessages) + if (state.message !== undefined && state.message.id !== itemID) completedMessages.add(state.message.id) // A new message closes earlier messages, including ones that never streamed. const events: LLMEvent[] = [] const lifecycle = [...state.lifecycle.text] .filter((id) => id !== itemID) .reduce((lifecycle, id) => { + completedMessages.add(id) const openPhase = state.message?.id === id ? state.message.phase : undefined return Lifecycle.textEnd( lifecycle, @@ -970,6 +977,7 @@ const onOutputItemAdded = (state: ParserState, event: Event): StepResult => { { ...state, lifecycle, + completedMessages, message: { id: itemID, phase: phase === undefined && state.message?.id === itemID ? state.message.phase : phase, @@ -1086,7 +1094,12 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* ( if (!item) return [state, NO_EVENTS] satisfies StepResult if (item.type === "message" && item.id !== undefined) { - const message = state.message?.id === item.id ? state.message : undefined + if (state.completedMessages.has(item.id)) return [state, NO_EVENTS] satisfies StepResult + const completedMessages = new Set(state.completedMessages) + completedMessages.add(item.id) + if (state.message !== undefined && state.message.id !== item.id) + return [{ ...state, completedMessages }, NO_EVENTS] satisfies StepResult + const message = state.message const itemPhase = messagePhase(item.phase) const phase = itemPhase === undefined ? message?.phase : itemPhase const parts: ReadonlyArray = Array.isArray(item.content) ? item.content : [] @@ -1099,13 +1112,13 @@ const onOutputItemDone = Effect.fn("OpenResponses.onOutputItemDone")(function* ( const text = content.length > 0 ? content.join("") : undefined const metadata = providerMetadata(state, { itemId: item.id, ...(phase === undefined ? {} : { phase }) }) const events: LLMEvent[] = [] - const lifecycle = - message && text ? Lifecycle.textStart(state.lifecycle, events, item.id, metadata) : state.lifecycle + const lifecycle = text ? Lifecycle.textStart(state.lifecycle, events, item.id, metadata) : state.lifecycle return [ { ...state, lifecycle: Lifecycle.textEnd(lifecycle, events, item.id, metadata, text), - message: message ? undefined : state.message, + completedMessages, + message: undefined, }, events, ] satisfies StepResult @@ -1419,6 +1432,7 @@ export const initial = (request: LLMRequest, adapter: ProviderAdapter = BASE_ADA lifecycle: Lifecycle.initial(), outputItems: {}, message: undefined, + completedMessages: new Set(), reasoningItems: {}, }) diff --git a/packages/ai/test/provider/open-responses-finals.test.ts b/packages/ai/test/provider/open-responses-finals.test.ts index d46a335ed7a..b8ad309c63d 100644 --- a/packages/ai/test/provider/open-responses-finals.test.ts +++ b/packages/ai/test/provider/open-responses-finals.test.ts @@ -82,6 +82,32 @@ describe("Open Responses completed item text", () => { expect(response.events.filter(LLMEvent.is.textStart)).toEqual([]) }), ) + + it.effect("assembles a done-only message once across replayed item events", () => + Effect.gen(function* () { + const item = { + type: "message", + id: "msg_1", + content: [{ type: "output_text", text: "Recovered" }], + } + const response = yield* generate( + { type: "response.output_text.delta", item_id: "msg_1", delta: "Ignored after resume" }, + { type: "response.output_item.done", item }, + { type: "response.output_item.added", item }, + { type: "response.output_item.done", item }, + completed, + ) + expect(response.text).toBe("Recovered") + expect(response.message.content).toEqual([ + { + type: "text", + text: "Recovered", + providerMetadata: { "openai-compatible": { itemId: "msg_1" } }, + }, + ]) + expect(response.events.filter(LLMEvent.is.textEnd)).toHaveLength(1) + }), + ) }) describe("Open Responses completed item reasoning", () => { diff --git a/packages/ai/test/provider/open-responses-lifecycle.test.ts b/packages/ai/test/provider/open-responses-lifecycle.test.ts index dfd3c8e4a4a..9a35b4c9a83 100644 --- a/packages/ai/test/provider/open-responses-lifecycle.test.ts +++ b/packages/ai/test/provider/open-responses-lifecycle.test.ts @@ -216,7 +216,63 @@ describe("Open Responses basic-item lifecycles", () => { ]) }), ) - it.effect("allows a message to be registered again without inheriting its previous phase", () => + + it.effect("preserves non-empty done-only message content without replaying duplicates", () => + Effect.gen(function* () { + const text = { + type: "message", + id: "msg_text", + content: [{ type: "output_text", text: "Done-only text." }], + } + const refusal = { + type: "message", + id: "msg_refusal", + content: [{ type: "refusal", refusal: "Done-only refusal." }], + } + const events = yield* collect( + { type: "response.output_item.done", item: text }, + { type: "response.output_item.done", item: text }, + { + type: "response.output_item.done", + item: { type: "message", id: "msg_empty", content: [{ type: "output_text", text: "" }] }, + }, + { + type: "response.output_item.done", + item: { type: "message", id: "msg_empty", content: [{ type: "output_text", text: "Late" }] }, + }, + { type: "response.output_item.done", item: refusal }, + { type: "response.output_item.done", item: refusal }, + completed, + ) + + expect(events.filter((event) => event.type.startsWith("text-"))).toEqual([ + { + type: "text-start", + id: "msg_text", + providerMetadata: { "openai-compatible": { itemId: "msg_text" } }, + }, + { + type: "text-end", + id: "msg_text", + text: "Done-only text.", + providerMetadata: { "openai-compatible": { itemId: "msg_text" } }, + }, + { + type: "text-start", + id: "msg_refusal", + providerMetadata: { "openai-compatible": { itemId: "msg_refusal" } }, + }, + { + type: "text-end", + id: "msg_refusal", + text: "Done-only refusal.", + providerMetadata: { "openai-compatible": { itemId: "msg_refusal" } }, + }, + ]) + }), + ) + + it.effect("treats a repeated message lifecycle as replay", () => Effect.gen(function* () { const events = yield* collect( { type: "response.output_item.added", item: { type: "message", id: "msg_1", phase: "commentary" } }, @@ -233,9 +289,44 @@ describe("Open Responses basic-item lifecycles", () => { id: "msg_1", providerMetadata: { "openai-compatible": { itemId: "msg_1", phase: "commentary" } }, }, - { type: "text-end", id: "msg_1", providerMetadata: { "openai-compatible": { itemId: "msg_1" } } }, ]) - expect(events.filter(LLMEvent.is.textDelta).map((event) => event.text)).toEqual(["First", "Second"]) + expect(events.filter(LLMEvent.is.textDelta).map((event) => event.text)).toEqual(["First"]) + }), + ) + + it.effect("ignores a stale done-only message while another message is active", () => + Effect.gen(function* () { + const events = yield* collect( + { type: "response.output_item.added", item: { type: "message", id: "msg_1", phase: "commentary" } }, + { type: "response.output_text.delta", item_id: "msg_1", delta: "Draft" }, + { + type: "response.output_item.done", + item: { type: "message", id: "msg_2", content: [{ type: "output_text", text: "Recovered" }] }, + }, + { + type: "response.output_item.done", + item: { type: "message", id: "msg_1", content: [{ type: "output_text", text: "Final" }] }, + }, + { + type: "response.output_item.done", + item: { type: "message", id: "msg_2", content: [{ type: "output_text", text: "Late" }] }, + }, + completed, + ) + expect(events.filter((event) => event.type.startsWith("text-"))).toEqual([ + { + type: "text-start", + id: "msg_1", + providerMetadata: { "openai-compatible": { itemId: "msg_1", phase: "commentary" } }, + }, + { type: "text-delta", id: "msg_1", text: "Draft" }, + { + type: "text-end", + id: "msg_1", + text: "Final", + providerMetadata: { "openai-compatible": { itemId: "msg_1", phase: "commentary" } }, + }, + ]) }), ) ;[undefined, "fc_1"].forEach((id) => {