From 2863a9c4e7f58964e6576c69283ce300c3f71a00 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sat, 27 Jun 2026 13:05:56 +0200 Subject: [PATCH] feat(ai): enabled textVerbosity option for openai api stream requests - Added textVerbosity option to OpenAIResponsesOptions and SimpleStreamOptions. - Updated stream mapping logic to propagate verbosity setting to the API request body. - Implemented endpoint validation to ensure verbosity settings are only applied to official OpenAI endpoints. - Added comprehensive test coverage for request payload inspection and stream event handling. --- packages/ai/src/providers/openai-responses.ts | 14 +++++ packages/ai/src/stream.ts | 3 + packages/ai/src/types.ts | 2 + packages/ai/test/openai-codex-stream.test.ts | 53 +++++++++++++++- .../openai-responses-cache-affinity.test.ts | 62 ++++++++++++++++++- 5 files changed, 129 insertions(+), 5 deletions(-) diff --git a/packages/ai/src/providers/openai-responses.ts b/packages/ai/src/providers/openai-responses.ts index d38998eda..55f5e447e 100644 --- a/packages/ai/src/providers/openai-responses.ts +++ b/packages/ai/src/providers/openai-responses.ts @@ -98,6 +98,7 @@ export interface OpenAIResponsesOptions extends StreamOptions { reasoning?: "minimal" | "low" | "medium" | "high" | "xhigh"; reasoningSummary?: "auto" | "detailed" | "concise" | null; serviceTier?: ServiceTier; + textVerbosity?: "low" | "medium" | "high"; toolChoice?: ToolChoice; openrouterVariant?: string; maxTokensExplicit?: boolean; @@ -772,6 +773,16 @@ const streamOpenAIResponsesOnce = ( export const streamOpenAIResponses: StreamFunction<"openai-responses"> = (model, context, options) => withEmptyCompletionRetry(model, context, options, streamOpenAIResponsesOnce); +function isOfficialOpenAIResponsesEndpoint(model: Model<"openai-responses">): boolean { + if (model.provider !== "openai") return false; + if (!model.baseUrl) return true; + try { + return new URL(model.baseUrl).hostname === "api.openai.com"; + } catch { + return false; + } +} + export function buildParams( model: Model<"openai-responses">, context: Context, @@ -859,6 +870,9 @@ export function buildParams( }); applyCommonResponsesSamplingParams(params, { ...options, maxTokens: outputToken?.value }, model); + if (options?.textVerbosity && isOfficialOpenAIResponsesEndpoint(model)) { + params.text = { ...params.text, verbosity: options.textVerbosity }; + } // TODO: openai responses has no top-level `stop`/`stop_sequences`; surface via reasoning.stop? // `StreamOptions.stopSequences` is intentionally dropped for this provider. // TODO: openai responses has no top-level `frequency_penalty` field as of the current SDK; diff --git a/packages/ai/src/stream.ts b/packages/ai/src/stream.ts index 8a628c190..d16b8568e 100644 --- a/packages/ai/src/stream.ts +++ b/packages/ai/src/stream.ts @@ -1493,6 +1493,7 @@ function mapOptionsForApi( openrouterVariant: options?.openrouterVariant, maxTokensExplicit: rawOptions?.maxTokens !== undefined, disableReasoning: options?.disableReasoning, + textVerbosity: options?.textVerbosity, }); } return castApi<"openai-completions">({ @@ -1527,6 +1528,7 @@ function mapOptionsForApi( openrouterVariant: options?.openrouterVariant, maxTokensExplicit: rawOptions?.maxTokens !== undefined, disableReasoning: options?.disableReasoning, + textVerbosity: options?.textVerbosity, }); case "azure-openai-responses": @@ -1546,6 +1548,7 @@ function mapOptionsForApi( serviceTier: options?.serviceTier, preferWebsockets: options?.preferWebsockets, reasoningSummary: options?.hideThinkingSummary ? null : "detailed", + textVerbosity: options?.textVerbosity, }); case "google-generative-ai": { diff --git a/packages/ai/src/types.ts b/packages/ai/src/types.ts index a292cf26b..f19e736c5 100644 --- a/packages/ai/src/types.ts +++ b/packages/ai/src/types.ts @@ -380,6 +380,8 @@ export interface SimpleStreamOptions extends Omit { * Useful when the UI hides thinking blocks anyway and the summary is wasted bandwidth. */ hideThinkingSummary?: boolean; + /** OpenAI Responses/Codex `text.verbosity` response detail level. */ + textVerbosity?: "low" | "medium" | "high"; /** Custom token budgets for thinking levels (token-based providers only) */ thinkingBudgets?: ThinkingBudgets; /** Cursor exec handlers for local tool execution */ diff --git a/packages/ai/test/openai-codex-stream.test.ts b/packages/ai/test/openai-codex-stream.test.ts index 936c00260..931604f05 100644 --- a/packages/ai/test/openai-codex-stream.test.ts +++ b/packages/ai/test/openai-codex-stream.test.ts @@ -1,4 +1,5 @@ import { afterEach, describe, expect, it, vi } from "bun:test"; +import { streamSimple } from "@oh-my-pi/pi-ai"; import { getOpenAICodexTransportDetails, getOpenAICodexWebSocketDebugStats, @@ -306,6 +307,34 @@ describe("openai-codex streaming", () => { expect(capturedBody?.prompt_cache_key).toBe("replacement-cache-key"); }); + it("forwards SimpleStreamOptions textVerbosity into the Codex request body", async () => { + const tempDir = TempDir.createSync("@pi-codex-stream-"); + setAgentDir(tempDir.path()); + const token = createCodexTestToken(); + const context = createCodexTestContext(); + const model = { ...createCodexTestModel("https://chatgpt.com/backend-api"), preferWebsockets: false }; + let capturedText: unknown; + const fetchMock: FetchImpl = async (_input, init) => { + if (typeof init?.body === "string") { + const parsed: { text?: unknown } = JSON.parse(init.body); + capturedText = parsed.text; + } + return new Response(createCompletedCodexSse("Hello"), { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + }; + + const result = await streamSimple(model, context, { + apiKey: token, + fetch: fetchMock, + textVerbosity: "low", + }).result(); + + expect(result.stopReason).toBe("stop"); + expect(capturedText).toEqual({ verbosity: "low" }); + }); + it("maps end_turn=false on the terminal event to a pause_turn stop", async () => { const tempDir = TempDir.createSync("@pi-codex-stream-"); setAgentDir(tempDir.path()); @@ -935,10 +964,18 @@ describe("openai-codex streaming", () => { ).toBase64(); const token = `aaa.${payload}.bbb`; + const textSignature = JSON.stringify({ v: 1, id: "msg_1", phase: "commentary" }); const sse = `${[ `data: ${JSON.stringify({ type: "response.output_item.added", - item: { type: "message", id: "msg_1", role: "assistant", status: "in_progress", content: [] }, + item: { + type: "message", + id: "msg_1", + role: "assistant", + status: "in_progress", + phase: "commentary", + content: [], + }, })}`, `data: ${JSON.stringify({ type: "response.content_part.added", part: { type: "output_text", text: "" } })}`, `data: ${JSON.stringify({ type: "response.output_text.delta", delta: "Hello" })}`, @@ -949,6 +986,7 @@ describe("openai-codex streaming", () => { id: "msg_1", role: "assistant", status: "completed", + phase: "commentary", content: [{ type: "output_text", text: "Hello" }], }, })}`, @@ -1018,18 +1056,28 @@ describe("openai-codex streaming", () => { const streamResult = streamOpenAICodexResponses(model, context, { apiKey: token, fetch: fetchMock as FetchImpl }); let sawTextDelta = false; + let sawTextStart = false; let sawDone = false; for await (const event of streamResult) { + if (event.type === "text_start") { + sawTextStart = true; + const block = event.partial.content[event.contentIndex]; + if (block?.type !== "text") throw new Error("expected text block"); + expect(block.textSignature).toBe(textSignature); + } if (event.type === "text_delta") { sawTextDelta = true; } if (event.type === "done") { sawDone = true; - expect(event.message.content.find(c => c.type === "text")?.text).toBe("Hello"); + const block = event.message.content.find(c => c.type === "text"); + expect(block?.text).toBe("Hello"); + expect(block?.textSignature).toBe(textSignature); } } + expect(sawTextStart).toBe(true); expect(sawTextDelta).toBe(true); expect(sawDone).toBe(true); }); @@ -2469,7 +2517,6 @@ describe("openai-codex streaming", () => { }); }); - it("uses websocket v2 beta header when v2 mode is enabled", async () => { const tempDir = TempDir.createSync("@pi-codex-stream-"); setAgentDir(tempDir.path()); diff --git a/packages/ai/test/openai-responses-cache-affinity.test.ts b/packages/ai/test/openai-responses-cache-affinity.test.ts index 988d0495a..1410d5824 100644 --- a/packages/ai/test/openai-responses-cache-affinity.test.ts +++ b/packages/ai/test/openai-responses-cache-affinity.test.ts @@ -1,7 +1,7 @@ import { afterEach, describe, expect, it, vi } from "bun:test"; import { type OpenAIResponsesOptions, streamOpenAIResponses } from "@oh-my-pi/pi-ai/providers/openai-responses"; -import { stream as streamModel } from "@oh-my-pi/pi-ai/stream"; -import type { Context, FetchImpl, Model, ProviderSessionState } from "@oh-my-pi/pi-ai/types"; +import { stream as streamModel, streamSimple } from "@oh-my-pi/pi-ai/stream"; +import type { Context, FetchImpl, Model, ProviderSessionState, SimpleStreamOptions } from "@oh-my-pi/pi-ai/types"; import { buildOpenAIResponsesCompat } from "@oh-my-pi/pi-catalog/compat/openai"; import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; @@ -189,6 +189,58 @@ async function captureDispatchedOpenAIResponseHeaders( return captured; } +async function captureSimpleOpenAIResponseBody( + options: SimpleStreamOptions, + requestModel: Model<"openai-responses"> = model, +): Promise | null> { + let body: Record | null = null; + const fetchMock: FetchImpl = vi.fn(async (_input: string | URL | Request, init?: RequestInit) => { + body = typeof init?.body === "string" ? (JSON.parse(init.body) as Record) : null; + return createSseResponse([ + { + type: "response.output_item.added", + item: { type: "message", id: "msg_1", role: "assistant", status: "in_progress", content: [] }, + }, + { type: "response.content_part.added", part: { type: "output_text", text: "" } }, + { type: "response.output_text.delta", delta: "Hello" }, + { + type: "response.output_item.done", + item: { + type: "message", + id: "msg_1", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "Hello" }], + }, + }, + { + type: "response.completed", + response: { + status: "completed", + usage: { + input_tokens: 5, + output_tokens: 3, + total_tokens: 8, + input_tokens_details: { cached_tokens: 0 }, + }, + }, + }, + ]); + }); + + const context: Context = { + systemPrompt: ["stable system", "stable durable context"], + messages: [{ role: "user", content: "hi", timestamp: Date.now() }], + }; + const stream = streamSimple(requestModel, context, { apiKey: "test-key", ...options, fetch: fetchMock }); + + for await (const event of stream) { + if (event.type === "done" || event.type === "error") break; + } + + return body; +} + afterEach(() => { vi.restoreAllMocks(); }); @@ -201,6 +253,12 @@ describe("openai-responses cache affinity", () => { expect(captured.clientRequestId).toBe("session-123"); expect(captured.body?.prompt_cache_key).toBe("session-123"); }); + + it("forwards textVerbosity through streamSimple to official OpenAI Responses text config", async () => { + const body = await captureSimpleOpenAIResponseBody({ textVerbosity: "low" }); + + expect(body?.text).toEqual({ verbosity: "low" }); + }); it("keeps prompt cache key separate from OpenAI routing headers when both are provided", async () => { const captured = await captureOpenAIResponseHeaders({ sessionId: "side-channel-456",