From 101888c011fa97bae10083bb4861f16a62541358 Mon Sep 17 00:00:00 2001 From: Jeff <36680501@qq.com> Date: Mon, 28 Sep 2026 00:08:15 +0800 Subject: [PATCH] fix: route OpenAI-compatible gateways to the Chat Completions adapter --- .env.example | 3 + apps/server/src/engine/tanstack-agent.ts | 9 +- tests/gateway-chat-completions.test.ts | 111 +++++++++++++++++++++++ tests/helpers/model.ts | 71 +++++++++++++++ 4 files changed, 193 insertions(+), 1 deletion(-) create mode 100644 tests/gateway-chat-completions.test.ts diff --git a/.env.example b/.env.example index 91e06682..32b0d87d 100644 --- a/.env.example +++ b/.env.example @@ -17,6 +17,9 @@ TASK_WORKER_ENABLED=true # Optional OpenAI-compatible Responses API endpoint: # OPENAI_BASE_URL= # For a gateway model ID such as vendor/model, set MODEL=openai/vendor/model. +# Providers whose /responses tool loop is incomplete (e.g. DeepSeek) can opt +# into the Chat Completions wire format instead: +# OPENAI_CHAT_COMPLETIONS=true # Live Google workspace: WORKSPACE_MODE=live, AGENT_BACKEND=model. # OPENMUSE_ACCESS_KEY= # random secret, at least 24 characters diff --git a/apps/server/src/engine/tanstack-agent.ts b/apps/server/src/engine/tanstack-agent.ts index 5701955a..2bc3644a 100644 --- a/apps/server/src/engine/tanstack-agent.ts +++ b/apps/server/src/engine/tanstack-agent.ts @@ -9,7 +9,7 @@ import { import { chat, maxIterations, type SchemaInput, toolDefinition } from "@tanstack/ai"; import { type AnthropicChatModel, anthropicText } from "@tanstack/ai-anthropic"; import { type GeminiTextModel, geminiText } from "@tanstack/ai-gemini"; -import { type OpenAIChatModel, openaiText } from "@tanstack/ai-openai"; +import { type OpenAIChatModel, openaiChatCompletions, openaiText } from "@tanstack/ai-openai"; import { map, mergeMap, type Observable } from "rxjs"; import { z } from "zod"; import { MODEL_MAX_RETRIES } from "../config.ts"; @@ -25,6 +25,13 @@ function adapter(spec: string) { const id = model.trim(); switch (provider.toLowerCase()) { case "openai": + // OpenAI-compatible endpoints that do not implement the Responses API + // tool loop (e.g. DeepSeek) can opt into the Chat Completions wire format. + if (process.env.OPENAI_CHAT_COMPLETIONS === "true") + return openaiChatCompletions(id as OpenAIChatModel, { + baseURL: process.env.OPENAI_BASE_URL, + maxRetries: MODEL_MAX_RETRIES, + }); return openaiText(id as OpenAIChatModel, { baseURL: process.env.OPENAI_BASE_URL, maxRetries: MODEL_MAX_RETRIES, diff --git a/tests/gateway-chat-completions.test.ts b/tests/gateway-chat-completions.test.ts new file mode 100644 index 00000000..83e28cc4 --- /dev/null +++ b/tests/gateway-chat-completions.test.ts @@ -0,0 +1,111 @@ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { EventType, type RunAgentInput } from "@ag-ui/core"; +import { type BuiltInAgent, defineTool } from "@copilotkit/runtime/v2"; +import { z } from "zod"; +import { tanstackAgent } from "../apps/server/src/engine/tanstack-agent.ts"; +import { chatCompletionsFixture, modelFixture } from "./helpers/model.ts"; + +const run = (agent: BuiltInAgent) => { + const input: RunAgentInput = { + threadId: "gateway-fixture", + runId: "gateway-fixture-run", + messages: [{ id: "m1", role: "user", content: "Use the tool once, then finish." }], + state: {}, + tools: [], + context: [], + forwardedProps: {}, + }; + return new Promise<{ error?: string; finished: boolean; text: string }>((resolve) => { + let error: string | undefined; + let finished = false; + let text = ""; + agent.run(input).subscribe({ + next: (event) => { + if ( + (event.type === EventType.TEXT_MESSAGE_CHUNK || + event.type === EventType.TEXT_MESSAGE_CONTENT) && + "delta" in event && + typeof event.delta === "string" + ) + text += event.delta; + if (event.type === EventType.RUN_ERROR && "message" in event) error = String(event.message); + if (event.type === EventType.RUN_FINISHED) finished = true; + }, + error: (cause) => { + if (error === undefined) error = String(cause); + resolve({ error, finished: false, text }); + }, + complete: () => resolve({ error, finished, text }), + }); + }); +}; + +test("OPENAI_CHAT_COMPLETIONS routes a gateway through /chat/completions and keeps the tool loop intact", async (t) => { + const previous = process.env.OPENAI_CHAT_COMPLETIONS; + process.env.OPENAI_CHAT_COMPLETIONS = "true"; + t.after(() => { + if (previous === undefined) delete process.env.OPENAI_CHAT_COMPLETIONS; + else process.env.OPENAI_CHAT_COMPLETIONS = previous; + }); + let toolRuns = 0; + const { requests } = await chatCompletionsFixture(t, (index) => + index === 0 ? { name: "note_step", arguments: { note: "step one done" } } : undefined, + ); + const agent = tanstackAgent({ + model: "openai/fixture", + maxSteps: 3, + tools: [ + defineTool({ + name: "note_step", + description: "Record a step note", + parameters: z.object({ note: z.string() }), + execute: async ({ note }) => { + toolRuns++; + return { recorded: note }; + }, + }), + ], + prompt: "Use the tool once, then finish.", + }); + const outcome = await run(agent); + assert.equal(outcome.error, undefined); + assert.equal(outcome.finished, true); + assert.ok(outcome.text.includes("Current state: empty."), "the final model reply arrived"); + assert.equal(toolRuns, 1, "the server tool executed exactly once"); + assert.ok(requests.length >= 2, "the tool result was sent back for a second model call"); + for (const request of requests) { + assert.equal(request.path, "/v1/chat/completions", "gateway traffic uses Chat Completions"); + } +}); + +test("without the flag the OpenAI adapter keeps using the Responses API", async (t) => { + const previous = process.env.OPENAI_CHAT_COMPLETIONS; + delete process.env.OPENAI_CHAT_COMPLETIONS; + t.after(() => { + if (previous === undefined) delete process.env.OPENAI_CHAT_COMPLETIONS; + else process.env.OPENAI_CHAT_COMPLETIONS = previous; + }); + const { requests } = await modelFixture(t, (index) => + index === 0 ? { name: "note_step", arguments: { note: "step one done" } } : undefined, + ); + const agent = tanstackAgent({ + model: "openai/fixture", + maxSteps: 3, + tools: [ + defineTool({ + name: "note_step", + description: "Record a step note", + parameters: z.object({ note: z.string() }), + execute: async () => ({ recorded: true }), + }), + ], + prompt: "Use the tool once, then finish.", + }); + const outcome = await run(agent); + assert.equal(outcome.error, undefined); + assert.equal(outcome.finished, true); + for (const request of requests) { + assert.equal(request.path, "/v1/responses", "the default wire format stays Responses"); + } +}); diff --git a/tests/helpers/model.ts b/tests/helpers/model.ts index 303bfce2..d79c4bfc 100644 --- a/tests/helpers/model.ts +++ b/tests/helpers/model.ts @@ -150,6 +150,77 @@ export async function modelFixture( assert.ok(address && typeof address !== "string"); const previousBase = process.env.OPENAI_BASE_URL; const previousKey = process.env.OPENAI_API_KEY; + const previousChatCompletions = process.env.OPENAI_CHAT_COMPLETIONS; + process.env.OPENAI_BASE_URL = `http://127.0.0.1:${address.port}/v1`; + process.env.OPENAI_API_KEY = "local-test-fixture"; + // This fixture serves the Responses protocol; a developer .env flag must + // not switch the adapter under test to the Chat Completions wire format. + delete process.env.OPENAI_CHAT_COMPLETIONS; + t.after(async () => { + if (previousBase === undefined) delete process.env.OPENAI_BASE_URL; + else process.env.OPENAI_BASE_URL = previousBase; + if (previousKey === undefined) delete process.env.OPENAI_API_KEY; + else process.env.OPENAI_API_KEY = previousKey; + if (previousChatCompletions === undefined) delete process.env.OPENAI_CHAT_COMPLETIONS; + else process.env.OPENAI_CHAT_COMPLETIONS = previousChatCompletions; + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + }); + return { requests }; +} + +// The same fixture protocol served in the OpenAI Chat Completions wire format +// (`/v1/chat/completions` SSE), for the optional gateway adapter path. +export async function chatCompletionsFixture( + t: TestContext, + reply: (index: number) => ModelCall | undefined | Promise, +) { + const requests: { path: string; body: string }[] = []; + const server = createServer(async (request, response) => { + let body = ""; + for await (const chunk of request) body += chunk; + const index = requests.length; + requests.push({ path: request.url ?? "", body }); + const call = await reply(index); + response.writeHead(200, { "Content-Type": "text/event-stream" }); + const chunk = (delta: object, finish: string | null) => + response.write( + `data: ${JSON.stringify({ + id: `chatcmpl-${index}`, + object: "chat.completion.chunk", + created: 1000, + model: "fixture", + choices: [{ index: 0, delta, finish_reason: finish }], + })}\n\n`, + ); + chunk({ role: "assistant" }, null); + if (call) { + chunk({ content: "I'll check first." }, null); + chunk( + { + tool_calls: [ + { + index: 0, + id: `call-${index}`, + type: "function", + function: { name: call.name, arguments: JSON.stringify(call.arguments) }, + }, + ], + }, + "tool_calls", + ); + } else { + chunk({ content: "Current state: empty." }, null); + } + chunk({}, "stop"); + response.end("data: [DONE]\n\n"); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const address = server.address(); + assert.ok(address && typeof address !== "string"); + const previousBase = process.env.OPENAI_BASE_URL; + const previousKey = process.env.OPENAI_API_KEY; process.env.OPENAI_BASE_URL = `http://127.0.0.1:${address.port}/v1`; process.env.OPENAI_API_KEY = "local-test-fixture"; t.after(async () => {