diff --git a/packages/ai/src/providers/openai-responses-server-schema.ts b/packages/ai/src/providers/openai-responses-server-schema.ts index 330dd9dc0..c9b13d573 100644 --- a/packages/ai/src/providers/openai-responses-server-schema.ts +++ b/packages/ai/src/providers/openai-responses-server-schema.ts @@ -92,8 +92,11 @@ const systemMessageItemSchema = type({ const assistantMessageItemSchema = type({ "type?": "'message'", + "id?": "string", role: "'assistant'", "content?": type("string").or(outputContentBlockSchema.array()), + "status?": "'in_progress' | 'completed' | 'incomplete'", + "phase?": "'commentary' | 'final_answer' | null", }); const reasoningItemSchema = type({ diff --git a/packages/ai/src/providers/openai-responses-server.ts b/packages/ai/src/providers/openai-responses-server.ts index e81b9caf6..6a4d947d8 100644 --- a/packages/ai/src/providers/openai-responses-server.ts +++ b/packages/ai/src/providers/openai-responses-server.ts @@ -33,6 +33,7 @@ import { type OpenAIResponsesTool, openaiResponsesRequestSchema, } from "./openai-responses-server-schema"; +import { encodeTextSignatureV1, parseTextSignature } from "./openai-shared"; export type { ParsedRequest }; @@ -54,6 +55,20 @@ function asString(v: unknown): string | undefined { return typeof v === "string" ? v : undefined; } +type AssistantItemPhase = "commentary" | "final_answer"; +type MessageSignature = { id: string; phase?: AssistantItemPhase }; + +function parseAssistantItemPhase(value: unknown): AssistantItemPhase | undefined { + return value === "commentary" || value === "final_answer" ? value : undefined; +} + +function messageTextSignature(id: unknown, phase: unknown): string | undefined { + const parsedPhase = parseAssistantItemPhase(phase); + if (typeof id === "string" && id.length > 0) return encodeTextSignatureV1(id, parsedPhase); + if (!parsedPhase) return undefined; + return encodeTextSignatureV1(makeMsgId(), parsedPhase); +} + // ─── id helpers ───────────────────────────────────────────────────────────── function uuidNoDashes(): string { @@ -147,20 +162,27 @@ type OutputBlockUnion = | { type: "text"; text: string } | { type: "refusal"; refusal: string }; -function outputTextOf(blocks: OpenAIResponsesOutputContent[] | string | undefined): TextContent[] { - if (typeof blocks === "string") return blocks.length > 0 ? [{ type: "text", text: blocks }] : []; +function outputTextOf( + blocks: OpenAIResponsesOutputContent[] | string | undefined, + message?: { id?: unknown; phase?: unknown }, +): TextContent[] { + const textSignature = messageTextSignature(message?.id, message?.phase); + const textContent = (text: string): TextContent => + textSignature ? { type: "text", text, textSignature } : { type: "text", text }; + if (typeof blocks === "string") return blocks.length > 0 ? [textContent(blocks)] : []; if (!blocks) return []; - const out: TextContent[] = []; + const parts: string[] = []; for (const raw of blocks) { const block = raw as OutputBlockUnion; if (block.type === "output_text" || block.type === "text") { - out.push({ type: "text", text: block.text }); + parts.push(block.text); } else if (block.type === "refusal") { // Preserve the refusal reason so history replay still carries it. - out.push({ type: "text", text: `[refusal: ${block.refusal}]` }); + parts.push(`[refusal: ${block.refusal}]`); } } - return out; + const text = parts.join(""); + return text.length > 0 ? [textContent(text)] : []; } // The schema accepts a much wider tool_choice union than the SDK type so the @@ -288,6 +310,8 @@ export function parseRequest(body: unknown, headers?: Headers): ParsedRequest { const msg = item as { role?: string; content?: OpenAIResponsesInputContent[] | OpenAIResponsesOutputContent[] | string; + id?: unknown; + phase?: unknown; }; switch (msg.role) { case "system": { @@ -303,7 +327,10 @@ export function parseRequest(body: unknown, headers?: Headers): ParsedRequest { break; } case "assistant": { - const parts = outputTextOf(msg.content as OpenAIResponsesOutputContent[] | string | undefined); + const parts = outputTextOf(msg.content as OpenAIResponsesOutputContent[] | string | undefined, { + id: msg.id, + phase: msg.phase, + }); messages.push({ role: "assistant", content: parts, @@ -503,6 +530,7 @@ type MessageOutputItem = { role: "assistant"; status: "completed"; content: Array<{ type: "output_text"; text: string; annotations: never[] }>; + phase?: AssistantItemPhase; }; type FunctionCallOutputItem = { @@ -594,23 +622,32 @@ function wireCallId(id: string): string { function buildOutputItems(message: AssistantMessage): OutputItem[] { const out: OutputItem[] = []; let pendingMessage: MessageOutputItem | null = null; + let pendingMessageSignature: { id: string; phase?: AssistantItemPhase } | undefined; const flushMessage = () => { if (pendingMessage) { out.push(pendingMessage); pendingMessage = null; + pendingMessageSignature = undefined; } }; for (const part of message.content) { if (part.type === "text") { + const signature = parseTextSignature(part.textSignature); + const sameSignature = + !pendingMessage || + (pendingMessageSignature?.id === signature?.id && pendingMessageSignature?.phase === signature?.phase); + if (!sameSignature) flushMessage(); if (!pendingMessage) { pendingMessage = { type: "message", - id: makeMsgId(), + id: signature?.id ?? makeMsgId(), role: "assistant", status: "completed", content: [], + ...(signature?.phase ? { phase: signature.phase } : {}), }; + pendingMessageSignature = signature; } pendingMessage.content.push({ type: "output_text", text: part.text, annotations: [] }); } else if (part.type === "thinking") { @@ -619,7 +656,8 @@ function buildOutputItems(message: AssistantMessage): OutputItem[] { } else if (part.type === "toolCall") { flushMessage(); if (part.customWireName) { - const rawInput = typeof part.arguments?.input === "string" ? (part.arguments.input as string) : ""; + const input = part.arguments?.input; + const rawInput = typeof input === "string" ? input : ""; out.push({ type: "custom_tool_call", id: part.thoughtSignature ?? makeCustomCallId(), @@ -701,6 +739,7 @@ interface OpenMessage { contentIndex: number; currentPartText: string; content: Array<{ type: "output_text"; text: string; annotations: never[] }>; + signature?: MessageSignature; } interface OpenReasoning { kind: "reasoning"; @@ -768,15 +807,16 @@ export function encodeStream( usage: null, }); - const openMessage = (): OpenMessage => { + const openMessage = (signature?: MessageSignature): OpenMessage => { const itemOutputIndex = allocateOutputIndex(); - const itemId = makeMsgId(); + const itemId = signature?.id ?? makeMsgId(); const item = { type: "message" as const, id: itemId, - status: "in_progress", + status: "in_progress" as const, role: "assistant" as const, content: [] as Array<{ type: "output_text"; text: string; annotations: never[] }>, + ...(signature?.phase ? { phase: signature.phase } : {}), }; emit("response.output_item.added", { output_index: itemOutputIndex, item }); const next: OpenMessage = { @@ -786,6 +826,7 @@ export function encodeStream( contentIndex: 0, currentPartText: "", content: [], + ...(signature ? { signature } : {}), }; state.open = next; return next; @@ -907,20 +948,15 @@ export function encodeStream( if (!state.open) return; if (state.open.kind === "message") { const item = { - type: "message", + type: "message" as const, id: state.open.itemId, - status: "completed", - role: "assistant", + status: "completed" as const, + role: "assistant" as const, content: state.open.content, + ...(state.open.signature?.phase ? { phase: state.open.signature.phase } : {}), }; emit("response.output_item.done", { output_index: state.open.outputIndex, item }); - finishedItems.push({ - type: "message", - id: state.open.itemId, - role: "assistant", - status: "completed", - content: state.open.content, - }); + finishedItems.push(item); state.open = null; } else if (state.open.kind === "reasoning") { const summary = [{ type: "summary_text" as const, text: state.open.reasoningText ?? "" }]; @@ -973,20 +1009,33 @@ export function encodeStream( } case "text_start": { let cur: OpenMessage; + const textBlock = ev.partial.content[ev.contentIndex]; + const signature = + textBlock?.type === "text" ? parseTextSignature(textBlock.textSignature) : undefined; if (state.open && state.open.kind === "message") { - // continue same message item, new content part - cur = state.open; - cur.currentPartText = ""; + const sameSignature = + (!signature && !state.open.signature) || + (signature !== undefined && + state.open.signature?.id === signature.id && + state.open.signature.phase === signature.phase); + if (sameSignature) { + // Continue same message item, new content part. + cur = state.open; + cur.currentPartText = ""; + } else { + closeOpen(); + cur = openMessage(signature); + } } else { if (state.open && state.open.kind !== "function_call") closeOpen(); - cur = openMessage(); + cur = openMessage(signature); } - const part = { type: "output_text", text: "", annotations: [] as never[] }; + const contentPart = { type: "output_text", text: "", annotations: [] as never[] }; emit("response.content_part.added", { item_id: cur.itemId, output_index: cur.outputIndex, content_index: cur.contentIndex, - part, + part: contentPart, }); break; } diff --git a/packages/ai/src/providers/openai-shared.ts b/packages/ai/src/providers/openai-shared.ts index 6fffcc485..a024d79b5 100644 --- a/packages/ai/src/providers/openai-shared.ts +++ b/packages/ai/src/providers/openai-shared.ts @@ -1927,7 +1927,11 @@ export async function processResponsesStream( registerOpenItem(event.output_index, item.id, { item, block }); stream.push({ type: "thinking_start", contentIndex: contentIndexOf(block), partial: output }); } else if (item.type === "message") { - const block: TextContent = { type: "text", text: "" }; + const block: TextContent = { + type: "text", + text: "", + textSignature: encodeTextSignatureV1(item.id, item.phase ?? undefined), + }; output.content.push(block); registerOpenItem(event.output_index, item.id, { item, block }); stream.push({ type: "text_start", contentIndex: contentIndexOf(block), partial: output }); diff --git a/packages/ai/test/auth-gateway-openai-responses.test.ts b/packages/ai/test/auth-gateway-openai-responses.test.ts index 4c50ccbcc..cbb68ec66 100644 --- a/packages/ai/test/auth-gateway-openai-responses.test.ts +++ b/packages/ai/test/auth-gateway-openai-responses.test.ts @@ -69,7 +69,12 @@ describe("openai-responses parseRequest", () => { { type: "message", role: "assistant", - content: [{ type: "output_text", text: "Let me think." }], + id: "msg_commentary", + phase: "commentary", + content: [ + { type: "output_text", text: "Let me " }, + { type: "output_text", text: "think." }, + ], }, reasoningItem, { @@ -121,6 +126,9 @@ describe("openai-responses parseRequest", () => { expect(a.model).toBe("gpt-5.3-codex-spark"); expect(a.content).toHaveLength(3); expect(a.content[0]).toMatchObject({ type: "text", text: "Let me think." }); + const commentary = a.content[0]; + if (commentary?.type !== "text") throw new Error("expected commentary text"); + expect(commentary.textSignature).toBe(JSON.stringify({ v: 1, id: "msg_commentary", phase: "commentary" })); expect(a.content[1]).toMatchObject({ type: "thinking", thinking: "The user wants arithmetic.", @@ -306,6 +314,51 @@ describe("openai-responses encodeResponse", () => { }); }); + it("encodes assistant message phase from text signatures", () => { + const message: AssistantMessage = { + role: "assistant", + api: "openai-responses", + provider: "openai", + model: "gpt-5", + content: [ + { + type: "text", + text: "Intermediate update", + textSignature: JSON.stringify({ v: 1, id: "msg_commentary", phase: "commentary" }), + }, + { + type: "text", + text: "Final answer", + textSignature: JSON.stringify({ v: 1, id: "msg_final", phase: "final_answer" }), + }, + ], + usage: zeroUsage(), + stopReason: "stop", + timestamp: 1_700_000_000_000, + }; + + const body = encodeResponse(message, "gpt-5-requested"); + const output = body.output as Array>; + + expect(output).toHaveLength(2); + expect(output[0]).toMatchObject({ + type: "message", + id: "msg_commentary", + role: "assistant", + status: "completed", + phase: "commentary", + content: [{ type: "output_text", text: "Intermediate update", annotations: [] }], + }); + expect(output[1]).toMatchObject({ + type: "message", + id: "msg_final", + role: "assistant", + status: "completed", + phase: "final_answer", + content: [{ type: "output_text", text: "Final answer", annotations: [] }], + }); + }); + it("marks length-limited responses incomplete", () => { const message: AssistantMessage = { role: "assistant", @@ -495,6 +548,44 @@ describe("openai-responses encodeStream", () => { expect(output[2]!.id).not.toBe(output[2]!.call_id); }); + it("streams assistant message phase from text signatures", async () => { + const stream = new AssistantMessageEventStream(); + const textSignature = JSON.stringify({ v: 1, id: "msg_commentary", phase: "commentary" }); + const message: AssistantMessage = { + role: "assistant", + api: "openai-responses", + provider: "openai", + model: "gpt-5", + content: [{ type: "text", text: "Working", textSignature }], + usage: { ...zeroUsage(), input: 1, output: 1 }, + stopReason: "stop", + timestamp: 1_700_000_000_000, + }; + + queueMicrotask(() => { + stream.push({ type: "start", partial: { ...message, content: [] } }); + stream.push({ type: "text_start", contentIndex: 0, partial: message }); + stream.push({ type: "text_delta", contentIndex: 0, delta: "Working", partial: message }); + stream.push({ type: "text_end", contentIndex: 0, content: "Working", partial: message }); + stream.push({ type: "done", reason: "stop", message }); + }); + + const raw = await collectStream(encodeStream(stream, "gpt-5-requested")); + const frames = parseSse(raw); + const messageItems = frames + .filter(f => f.event === "response.output_item.added" || f.event === "response.output_item.done") + .map(f => (f.data as Record).item as Record) + .filter(item => item.type === "message"); + const completed = frames.find(f => f.event === "response.completed")?.data as Record | undefined; + const response = completed?.response as Record | undefined; + const output = response?.output as Array> | undefined; + + expect(messageItems).toHaveLength(2); + expect(messageItems[0]).toMatchObject({ id: "msg_commentary", phase: "commentary" }); + expect(messageItems[1]).toMatchObject({ id: "msg_commentary", phase: "commentary" }); + expect(output?.[0]).toMatchObject({ id: "msg_commentary", phase: "commentary" }); + }); + it("routes late tool-call deltas by contentIndex after later parallel starts", async () => { const stream = new AssistantMessageEventStream(); const base: AssistantMessage = { diff --git a/packages/ai/test/openai-responses-delta-input.test.ts b/packages/ai/test/openai-responses-delta-input.test.ts index d4135bf81..46cce96f6 100644 --- a/packages/ai/test/openai-responses-delta-input.test.ts +++ b/packages/ai/test/openai-responses-delta-input.test.ts @@ -85,4 +85,29 @@ describe("buildResponsesDeltaInput streaming-symbol scrub", () => { }; expect(buildResponsesDeltaInput(previous, [items[1]], current)).toBeNull(); }); + + it("treats assistant message phase as part of chained-prefix equality", () => { + const user = baselineItems()[0]!; + const previousAssistant: ResponseInputItem = { + type: "message", + role: "assistant", + content: "intermediate update", + phase: "commentary", + }; + const appended: ResponseInputItem = { + type: "message", + role: "user", + content: [{ type: "input_text", text: "follow-up" }], + }; + const previous = { input: [user] }; + + expect( + buildResponsesDeltaInput(previous, [previousAssistant], { input: [user, previousAssistant, appended] }), + ).toEqual([appended]); + + const wrongPhaseAssistant: ResponseInputItem = { ...previousAssistant, phase: "final_answer" }; + expect( + buildResponsesDeltaInput(previous, [previousAssistant], { input: [user, wrongPhaseAssistant, appended] }), + ).toBeNull(); + }); }); diff --git a/packages/ai/test/openai-responses-history-payload.test.ts b/packages/ai/test/openai-responses-history-payload.test.ts index e3af39a14..8f6c004fc 100644 --- a/packages/ai/test/openai-responses-history-payload.test.ts +++ b/packages/ai/test/openai-responses-history-payload.test.ts @@ -222,6 +222,7 @@ const incrementalItems1 = [ content: [{ type: "output_text", text: "First response" }], status: "completed", id: "msg_1", + phase: "commentary", }, ]; @@ -232,6 +233,7 @@ const incrementalItems2 = [ content: [{ type: "output_text", text: "Second response" }], status: "completed", id: "msg_2", + phase: "final_answer", }, ];