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: [],