Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 8 additions & 1 deletion apps/server/src/engine/tanstack-agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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,
Expand Down
111 changes: 111 additions & 0 deletions tests/gateway-chat-completions.test.ts
Original file line number Diff line number Diff line change
@@ -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");
}
});
71 changes: 71 additions & 0 deletions tests/helpers/model.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>((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<ModelCall | undefined>,
) {
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 () => {
Expand Down