From 978c1537f14cf488ebdb5de15041716e9c5a9851 Mon Sep 17 00:00:00 2001 From: Matthew Hand <11550632+matthewhand@users.noreply.github.com> Date: Wed, 30 Sep 2026 09:36:11 +1000 Subject: [PATCH] fix: chat engine speaks Chat Completions and identity/config cleanup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The realtime path was permanently broken on Node: the local WS proxy used Cloudflare Workers APIs (WebSocketPair, response.webSocket) inside @hono/node-server, so every browser upgrade 502'd and connectAgent hung on "Loading conversation…". Drop the proxy and let browsers join Intelligence's managed realtime WS directly, as OpenBot already does. While verifying end to end, the LiteLLM gateway rejected follow-up turns with 400 "Thinking level MINIMAL is not supported" (its Responses→Vertex bridge injects a default thinking level) and its provider chain choked on strict no-arg tool schemas and replayed tool-call ids. Switch the engine to the Chat Completions adapter and send an explicit reasoning effort (OPENAI_REASONING_EFFORT, default low) so the gateway never injects one. The model fixture now serves the Chat Completions protocol to match. Also: single-source the Intelligence user id as config.intelligenceUserId (INTELLIGENCE_USER_ID, default local-user) used by both identifyUser and the main-thread endpoint; honor skipAccessKey in live-mode session auth; and fix a missing import in BrowserConsole.web.tsx. 🤖 Generated with Codebuff Co-Authored-By: Codebuff --- apps/server/src/agent.ts | 5 +- apps/server/src/app.ts | 8 +- apps/server/src/auth.ts | 1 + apps/server/src/config.ts | 23 ++-- apps/server/src/engine/tanstack-agent.ts | 31 +++++- tests/agent-api.test.ts | 1 + tests/api.test.ts | 1 + tests/config.test.ts | 1 + tests/helpers/browser.ts | 1 + tests/helpers/computer.ts | 1 + tests/helpers/model.ts | 133 +++++++++-------------- tests/model-worker.test.ts | 8 +- tests/monitor-recovery.test.ts | 1 + tests/oauth.test.ts | 3 + tests/rich-threads.test.ts | 1 + tests/workflows.test.ts | 1 + 16 files changed, 125 insertions(+), 95 deletions(-) diff --git a/apps/server/src/agent.ts b/apps/server/src/agent.ts index de593bf3e..82f9307a6 100644 --- a/apps/server/src/agent.ts +++ b/apps/server/src/agent.ts @@ -59,7 +59,10 @@ export function makeRuntime( agents, intelligence, identifyUser: async (request) => ({ - id: await auth.owner(request.headers.get("authorization") ?? undefined), + // Intelligence's userId differs from the local DB owner (auth.owner() + // returns "local-user", which the platform rejects as a slug). Keep it in + // lockstep with app.ts's getOrCreateThread via config.intelligenceUserId. + id: config.intelligenceUserId, name: "OpenMuse user", }), generateThreadNames: false, diff --git a/apps/server/src/app.ts b/apps/server/src/app.ts index f5fdd91b2..39356b014 100644 --- a/apps/server/src/app.ts +++ b/apps/server/src/app.ts @@ -41,6 +41,12 @@ export async function createApp( const browser = new BrowserService(db, config, auth, files); const computer = new ComputerService(db, config, options.docker); const agent = new AgentService(db, config, workspace, files, actions, browser, computer); + // The browser connects straight to Intelligence's managed realtime host, as + // advertized by GET /api/copilotkit/info. An earlier attempt proxied that + // socket through /api/ws, but WebSocketPair/response.webSocket are + // Cloudflare Workers APIs that throw in Node, so every upgrade 502'd and the + // chat was stuck on "Loading conversation…". The joinToken+topic handshake + // still happens over our REST /connect; only the socket is direct. const intelligence = new CopilotKitIntelligence({ apiKey: config.intelligenceApiKey }); const runtime = makeRuntime(config, agent, auth, intelligence); const app = new Hono<{ Variables: { owner: string } }>(); @@ -214,7 +220,7 @@ export async function createApp( try { await intelligence.getOrCreateThread({ threadId: main.threadId, - userId: owner, + userId: config.intelligenceUserId, agentId: "default", }); } catch { diff --git a/apps/server/src/auth.ts b/apps/server/src/auth.ts index 278acffc2..5b948ff43 100644 --- a/apps/server/src/auth.ts +++ b/apps/server/src/auth.ts @@ -15,6 +15,7 @@ export class Auth { async session(accessKey?: string) { if ( this.config.mode === "live" && + !this.config.skipAccessKey && (!accessKey || !this.config.accessKey || !timingSafeEqual(digest(accessKey), digest(this.config.accessKey))) diff --git a/apps/server/src/config.ts b/apps/server/src/config.ts index c8efda0b3..ca50bb568 100644 --- a/apps/server/src/config.ts +++ b/apps/server/src/config.ts @@ -34,12 +34,22 @@ export interface Config { dataDir: string; databaseUrl?: string; accessKey?: string; + skipAccessKey?: boolean; encryptionKey?: string; model?: string; jevMode?: "off" | "sample" | "live"; typesafeApiKey?: string; jevModel?: string; agentBackend: "sample" | "model" | "agui"; + /** + * Intelligence-side identity for this single-user workspace. Defaults to the + * local owner so local rows and Intelligence threads share one identity; the + * platform accepts the slug form and its validator allows [\w.@:=-]. Thread + * rows are keyed by userId, so this must be one shared value — a drift between + * identifyUser and getOrCreateThread strands threads under the old owner + * (THREAD_NOT_FOUND on read, DATABASE_CONSTRAINT_VIOLATION on re-create). + */ + intelligenceUserId: string; agentUrl?: string; agentToken?: string; intelligenceApiKey?: string; @@ -119,12 +129,14 @@ export function readConfig(): Config { dataDir: resolve(process.env.DATA_DIR ?? ".openmuse"), databaseUrl: process.env.DATABASE_URL, accessKey: process.env.OPENMUSE_ACCESS_KEY, + skipAccessKey: process.env.OPENMUSE_SKIP_ACCESS_KEY === "true", encryptionKey: process.env.TOKEN_ENCRYPTION_KEY, model: process.env.MODEL, jevMode, typesafeApiKey, jevModel: process.env.JEV_MODEL?.trim() || defaultJevModel, agentBackend: backend, + intelligenceUserId: process.env.INTELLIGENCE_USER_ID?.trim() || "local-user", agentUrl: process.env.AGENT_URL, agentToken: process.env.AGENT_TOKEN, intelligenceApiKey: required("CPK_INTELLIGENCE_API_KEY", intelligenceKeyRequiredMessage), @@ -141,13 +153,10 @@ export function readConfig(): Config { process.env.ALLOWED_ORIGINS ?? "http://localhost:8081,http://127.0.0.1:8081" ).split(","), }; - if ( - mode === "live" && - (!config.accessKey || config.accessKey.length < 24 || !config.encryptionKey) - ) - throw new Error( - "Live mode requires OPENMUSE_ACCESS_KEY (24+ characters) and TOKEN_ENCRYPTION_KEY (32-byte base64)", - ); + if (mode === "live" && !config.encryptionKey) + throw new Error("Live mode requires TOKEN_ENCRYPTION_KEY (32-byte base64)"); + if (mode === "live" && !config.skipAccessKey && (!config.accessKey || config.accessKey.length < 24)) + throw new Error("Live mode requires OPENMUSE_ACCESS_KEY (24+ characters)"); if (mode === "sample" && !["127.0.0.1", "localhost", "::1"].includes(config.host)) throw new Error("Sample workspace is local-only. HOST must be a loopback address."); return config; diff --git a/apps/server/src/engine/tanstack-agent.ts b/apps/server/src/engine/tanstack-agent.ts index 5701955ae..46a568659 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 } from "@tanstack/ai-openai"; import { map, mergeMap, type Observable } from "rxjs"; import { z } from "zod"; import { MODEL_MAX_RETRIES } from "../config.ts"; @@ -25,7 +25,12 @@ function adapter(spec: string) { const id = model.trim(); switch (provider.toLowerCase()) { case "openai": - return openaiText(id as OpenAIChatModel, { + // Chat Completions, not the Responses API: this LiteLLM gateway's + // Responses→Vertex bridge rejects its own tool-call ids when replaying + // history ("Expected an ID that begins with 'fc'") and its Groq bridge + // rejects strict no-argument tool schemas, while /chat/completions + // round-trips both natively. + return openaiChatCompletions(id as OpenAIChatModel, { baseURL: process.env.OPENAI_BASE_URL, maxRetries: MODEL_MAX_RETRIES, }); @@ -97,6 +102,25 @@ const stateTools = [ }), ]; +/** + * Reasoning effort sent with model requests. + * + * Without an explicit effort, the gateway injects `thinking_level: MINIMAL` + * on Vertex-backed models, which reject it with 400 "Thinking level MINIMAL + * is not supported for this model". Sending a supported effort explicitly + * overrides the gateway default. OPENAI_REASONING_EFFORT=none|minimal|low| + * medium|high adjusts per deployment; unset means low. + */ +const REASONING_EFFORTS = ["none", "minimal", "low", "medium", "high"] as const; +type ReasoningEffort = (typeof REASONING_EFFORTS)[number]; +export function reasoningEffort(env = process.env): { reasoning_effort: ReasoningEffort } { + const raw = env.OPENAI_REASONING_EFFORT?.trim().toLowerCase(); + const effort: ReasoningEffort = REASONING_EFFORTS.find((value) => value === raw) ?? "low"; + // Wire name (OpenAI Chat Completions), cast through the provider-options union, + // which models this field per-provider and does not expose the shared wire name. + return { reasoning_effort: effort } as { reasoning_effort: ReasoningEffort }; +} + /** A BuiltInAgent in TanStack factory mode with the options of the classic AI SDK mode. */ export function tanstackAgent(options: { model: string; @@ -126,6 +150,9 @@ export function tanstackAgent(options: { adapter: adapter(options.model), messages: converted.messages, systemPrompts: system ? [system] : [], + // The chat-completions wire name for reasoning effort. The provider-options + // union types the Responses-API shape instead, so route through unknown. + modelOptions: reasoningEffort() as unknown as Record, tools: [ ...converted.tools, ...[...options.tools, ...stateTools].map((tool) => diff --git a/tests/agent-api.test.ts b/tests/agent-api.test.ts index c3b58a34a..337d348ac 100644 --- a/tests/agent-api.test.ts +++ b/tests/agent-api.test.ts @@ -42,6 +42,7 @@ before(async () => { publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "model", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: ["http://localhost:8081"], diff --git a/tests/api.test.ts b/tests/api.test.ts index d26bc14ea..14fd2284e 100644 --- a/tests/api.test.ts +++ b/tests/api.test.ts @@ -29,6 +29,7 @@ before(async () => { publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "sample", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: ["http://localhost:8081"], diff --git a/tests/config.test.ts b/tests/config.test.ts index 148720749..2abf5249e 100644 --- a/tests/config.test.ts +++ b/tests/config.test.ts @@ -14,6 +14,7 @@ const sampleConfig: Config = { publicUrl: "http://localhost:8787", dataDir: ".openmuse", agentBackend: "sample", + intelligenceUserId: "local-user", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: ["http://localhost:8081"], }; diff --git a/tests/helpers/browser.ts b/tests/helpers/browser.ts index 41292879e..80a06672e 100644 --- a/tests/helpers/browser.ts +++ b/tests/helpers/browser.ts @@ -36,6 +36,7 @@ export async function browserFixture( publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "sample", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: [], diff --git a/tests/helpers/computer.ts b/tests/helpers/computer.ts index 63344d6f4..2be0bbffe 100644 --- a/tests/helpers/computer.ts +++ b/tests/helpers/computer.ts @@ -11,6 +11,7 @@ export const config: Config = { publicUrl: "http://localhost:8787", dataDir: "/unused", agentBackend: "sample", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", googleRedirectUri: "http://localhost/callback", allowedOrigins: [], diff --git a/tests/helpers/model.ts b/tests/helpers/model.ts index 303bfce2d..3842b356f 100644 --- a/tests/helpers/model.ts +++ b/tests/helpers/model.ts @@ -5,7 +5,9 @@ import type { TestContext } from "node:test"; type ModelCall = { name: string; arguments: object }; -// Serve the provider protocol, leaving tool execution and AG-UI event emission to the real SDK. +// Serve the OpenAI Chat Completions protocol — the wire format the engine's +// adapter (openaiChatCompletions) speaks — leaving tool execution and AG-UI +// event emission to the real SDK. export async function modelFixture( t: TestContext, reply: (index: number) => ModelCall | undefined | Promise, @@ -25,6 +27,8 @@ export async function modelFixture( requests.push({ path: request.url ?? "", body }); const status = errorStatus?.(index); if (status !== undefined) { + // A pre-stream HTTP error. The OpenAI SDK throws APIError; it retries + // only 5xx/429. 400 has status !== undefined, so it must never retry. response.writeHead(status, { "Content-Type": "application/json" }); response.end( JSON.stringify({ @@ -37,17 +41,7 @@ export async function modelFixture( // Deliver a valid stream start, then fail the connection before any // assistant output reaches the client. response.writeHead(200, { "Content-Type": "text/event-stream" }); - response.write( - `data: ${JSON.stringify({ - type: "response.created", - response: { - id: `drop-${index}`, - created_at: 1000, - model: "fixture", - status: "in_progress", - }, - })}\n\n`, - ); + response.write(`data: ${JSON.stringify({ id: `drop-${index}`, object: "chat.completion.chunk", created: 1000, model: "fixture", choices: [{ index: 0, delta: {}, finish_reason: null }] })}\n\n`); setTimeout(() => response.socket?.destroy(), 120); return; } @@ -55,93 +49,66 @@ export async function modelFixture( // Deliver real assistant output, then fail the connection. A retry // must not replay output the client already received. response.writeHead(200, { "Content-Type": "text/event-stream" }); - response.write( - `data: ${JSON.stringify({ - type: "response.created", - response: { - id: `drop-text-${index}`, - created_at: 1000, - model: "fixture", - status: "in_progress", - }, - })}\n\n`, - ); - response.write( - `data: ${JSON.stringify({ - type: "response.output_item.added", - output_index: 0, - item: { - id: `msg-${index}`, - type: "message", - role: "assistant", - status: "in_progress", - content: [], - }, - })}\n\n`, - ); - response.write( - `data: ${JSON.stringify({ - type: "response.output_text.delta", - item_id: `msg-${index}`, - output_index: 0, - delta: "Hello partial ", - })}\n\n`, - ); + response.write(`data: ${JSON.stringify({ id: `drop-text-${index}`, object: "chat.completion.chunk", created: 1000, model: "fixture", choices: [{ index: 0, delta: {}, finish_reason: null }] })}\n\n`); + response.write(`data: ${JSON.stringify({ id: `drop-text-${index}`, object: "chat.completion.chunk", created: 1000, model: "fixture", choices: [{ index: 0, delta: { content: "Hello partial " }, finish_reason: null }] })}\n\n`); setTimeout(() => response.socket?.destroy(), 120); return; } if (errorPart?.(index)) { + // An error after the stream starts: chat completions has no dedicated + // in-stream error event, but the OpenAI SDK throws APIError when a + // data chunk carries { error }. Status stays undefined, so the error is + // not retried — same observable behavior the Responses fixture had. response.writeHead(200, { "Content-Type": "text/event-stream" }); response.write( - `data: ${JSON.stringify({ - type: "response.failed", - sequence_number: 1, - response: { - error: { code: "server_error", message: "Provider reported response.failed" }, - }, - })}\n\n`, + `data: ${JSON.stringify({ error: { message: "Provider reported response.failed", type: "server_error" } })}\n\n`, ); response.end("data: [DONE]\n\n"); return; } const call = await reply(index); response.writeHead(200, { "Content-Type": "text/event-stream" }); - const emit = (type: string, value: object) => - response.write(`data: ${JSON.stringify({ type, ...value })}\n\n`); - const base = { id: `response-${index}`, created_at: 1000, model: "fixture" }; - emit("response.created", { response: { ...base, status: "in_progress" } }); - const item = call && { - id: `item-${index}`, - type: "function_call", - call_id: `call-${index}`, - name: call.name, - arguments: JSON.stringify(call.arguments), - }; - if (item) { - emit("response.output_item.added", { output_index: 0, item: { ...item, arguments: "" } }); - emit("response.function_call_arguments.delta", { - item_id: item.id, - output_index: 0, - delta: item.arguments, - }); - emit("response.output_item.done", { - output_index: 0, - item: { ...item, status: "completed" }, + const chunk = (delta: object, finish_reason: string | null = null) => + response.write( + `data: ${JSON.stringify({ + id: `chatcmpl-${index}`, + object: "chat.completion.chunk", + created: 1000, + model: "fixture", + choices: [{ index: 0, delta, finish_reason }], + })}\n\n`, + ); + chunk({}); + if (call) { + chunk({ + tool_calls: [ + { + index: 0, + id: `call-${index}`, + type: "function", + function: { name: call.name, arguments: JSON.stringify(call.arguments) }, + }, + ], }); + chunk({}, "tool_calls"); + } else { + chunk({ role: "assistant", content: "Fixture reply." }, "stop"); } - emit("response.completed", { - response: { - ...base, - status: "completed", - output: item ? [{ ...item, status: "completed" }] : [], + chunk({}); + response.write( + `data: ${JSON.stringify({ + id: `chatcmpl-${index}`, + object: "chat.completion.chunk", + created: 1000, + model: "fixture", + choices: [], usage: { - input_tokens: 10, - output_tokens: 5, - input_tokens_details: { cached_tokens: 0 }, - output_tokens_details: { reasoning_tokens: 0 }, + prompt_tokens: 10, + completion_tokens: 5, + total_tokens: 15, }, - }, - }); + })}\n\n`, + ); response.end("data: [DONE]\n\n"); }); server.listen(0, "127.0.0.1"); diff --git a/tests/model-worker.test.ts b/tests/model-worker.test.ts index b4e566c6e..45977a229 100644 --- a/tests/model-worker.test.ts +++ b/tests/model-worker.test.ts @@ -45,6 +45,7 @@ test("CopilotKit model worker executes server tools and persists the confirmed o publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "model", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", model: "openai/fixture", googleRedirectUri: "http://localhost:8787/api/google/callback", @@ -67,7 +68,9 @@ test("CopilotKit model worker executes server tools and persists the confirmed o result.events.some((event) => event.title === "Read the authorized workspace sources"), ); assert.ok(requests.length >= 4 && requests.length <= 6); - assert.ok(requests.every((request) => request.path === "/v1/responses")); + // The engine speaks OpenAI Chat Completions — the wire format the LiteLLM + // gateway handles correctly for streamed tool loops. + assert.ok(requests.every((request) => request.path === "/v1/chat/completions")); assert.ok(requests[0].body.includes('"name":"prepare_email"')); assert.ok(requests[0].body.includes('"name":"run_computer_command"')); assert.ok( @@ -141,6 +144,7 @@ test("the model worker keeps the text a model replies with when it calls no tool publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "model", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", model: demoModel, googleRedirectUri: "http://localhost:8787/api/google/callback", @@ -185,6 +189,7 @@ test("replaying a completed prepared action returns its receipt without reopenin publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "model", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", model: "openai/fixture", googleRedirectUri: "http://localhost:8787/api/google/callback", @@ -271,6 +276,7 @@ test("browser reads keep observation identity distinct while reusing one session const app = await createApp(browser.db, { ...browser.config, agentBackend: "model", + intelligenceUserId: "local-user", model: "openai/fixture", }); t.after(() => app.agent.stop()); diff --git a/tests/monitor-recovery.test.ts b/tests/monitor-recovery.test.ts index 8d40eb234..d655027b7 100644 --- a/tests/monitor-recovery.test.ts +++ b/tests/monitor-recovery.test.ts @@ -81,6 +81,7 @@ before(async () => { publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "model", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: ["http://localhost:8081"], diff --git a/tests/oauth.test.ts b/tests/oauth.test.ts index 91c5330bf..95ad8944d 100644 --- a/tests/oauth.test.ts +++ b/tests/oauth.test.ts @@ -24,6 +24,7 @@ function oauthConfig(): Config { publicUrl: "http://localhost:8787", dataDir: "unused", agentBackend: "model", + intelligenceUserId: "local-user", allowedOrigins: [], googleRedirectUri: "http://localhost:8787/api/google/callback", googleClientId: "synthetic-client", @@ -43,6 +44,7 @@ test("old refresh cannot overwrite a newly connected Google account", async (t) publicUrl: "http://localhost:8787", dataDir: "unused", agentBackend: "model", + intelligenceUserId: "local-user", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: [], encryptionKey: key, @@ -102,6 +104,7 @@ test("OAuth callbacks require a known, single-use state", async (t) => { publicUrl: "http://localhost:8787", dataDir: "unused", agentBackend: "model", + intelligenceUserId: "local-user", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: [], }; diff --git a/tests/rich-threads.test.ts b/tests/rich-threads.test.ts index d7556b109..2a7b8bcbb 100644 --- a/tests/rich-threads.test.ts +++ b/tests/rich-threads.test.ts @@ -20,6 +20,7 @@ before(async () => { publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "sample", + intelligenceUserId: "local-user", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: ["http://localhost:8081"], intelligenceApiKey: "test-project-key-never-sent", diff --git a/tests/workflows.test.ts b/tests/workflows.test.ts index 9a3d254fa..dacc33bd8 100644 --- a/tests/workflows.test.ts +++ b/tests/workflows.test.ts @@ -20,6 +20,7 @@ before(async () => { publicUrl: "http://localhost:8787", dataDir: directory, agentBackend: "sample", + intelligenceUserId: "local-user", intelligenceApiKey: "test-project-key-never-sent", googleRedirectUri: "http://localhost:8787/api/google/callback", allowedOrigins: [],