diff --git a/apps/server/src/agent.ts b/apps/server/src/agent.ts index de593bf3e..15d6bcc80 100644 --- a/apps/server/src/agent.ts +++ b/apps/server/src/agent.ts @@ -8,6 +8,7 @@ import { } from "@copilotkit/runtime/v2"; import type { Auth } from "./auth.ts"; import type { Config } from "./config.ts"; +import type { AgentRunner } from "@copilotkit/runtime/v2"; import { ConversationAgent } from "./engine/conversation.ts"; import type { AgentService } from "./engine/service.ts"; import { createJevAdapter, type JevAdapter } from "./jev/adapter.ts"; @@ -29,7 +30,8 @@ export function makeRuntime( config: Config, service: AgentService, auth: Auth, - intelligence: CopilotKitIntelligence, + intelligence: CopilotKitIntelligence | undefined, + runner?: AgentRunner, ) { // Built on first use, then shared so live mode reuses one TypeSafe client across requests. let jevAdapter: JevAdapter | undefined; @@ -55,14 +57,26 @@ export function makeRuntime( sharedJevAdapter(), ), }); - const runtime = new CopilotRuntime({ - agents, - intelligence, - identifyUser: async (request) => ({ - id: await auth.owner(request.headers.get("authorization") ?? undefined), - name: "OpenMuse user", - }), - generateThreadNames: false, - }); + // The two branches map to the runtime's two official modes: + // - with Intelligence → cloud (or self-hosted shim) threads + realtime WS + // - without → standard AG-UI/SSE runtime with local thread endpoints + const runtime = intelligence + ? new CopilotRuntime({ + agents, + intelligence, + identifyUser: async (request) => ({ + id: await auth.owner(request.headers.get("authorization") ?? undefined), + name: "OpenMuse user", + }), + generateThreadNames: false, + }) + : new CopilotRuntime({ + agents, + ...(runner ? { runner } : {}), + identifyUser: async (request: Request) => ({ + id: await auth.owner(request.headers.get("authorization") ?? undefined), + name: "OpenMuse user", + }), + }); return createCopilotHonoHandler({ runtime, basePath: "/api/copilotkit" }); } diff --git a/apps/server/src/app.ts b/apps/server/src/app.ts index f5fdd91b2..981b099b7 100644 --- a/apps/server/src/app.ts +++ b/apps/server/src/app.ts @@ -12,7 +12,12 @@ import { createAuth } from "./auth.ts"; import { BrowserService } from "./browser.ts"; import { ComputerService, type DockerRunner } from "./computer.ts"; import { computerRoutes } from "./computer-routes.ts"; -import { assertApiDeploymentConfig, type Config } from "./config.ts"; +import { DurableAgentRunner } from "./engine/durable-runner.ts"; +import { + assertApiDeploymentConfig, + type Config, + intelligenceConfigured, +} from "./config.ts"; import type { Store } from "./db.ts"; import { agentRoutes } from "./engine/routes.ts"; import { AgentService } from "./engine/service.ts"; @@ -41,8 +46,17 @@ 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); - const intelligence = new CopilotKitIntelligence({ apiKey: config.intelligenceApiKey }); - const runtime = makeRuntime(config, agent, auth, intelligence); + // Intelligence is optional: with the key (or a self-hosted shim) the runtime + // uses cloud threads and the managed realtime WS; without it, the official + // SSE runtime serves local thread endpoints, persisted by DurableAgentRunner. + const intelligence = intelligenceConfigured(config) + ? new CopilotKitIntelligence({ + apiKey: config.intelligenceApiKey ?? "local-shim", + ...(config.intelligenceApiUrl ? { apiUrl: config.intelligenceApiUrl } : {}), + ...(config.intelligenceWsUrl ? { wsUrl: config.intelligenceWsUrl } : {}), + }) + : undefined; + const runtime = makeRuntime(config, agent, auth, intelligence, new DurableAgentRunner(db)); const app = new Hono<{ Variables: { owner: string } }>(); const origins = new Set([...config.allowedOrigins, new URL(config.publicUrl).origin]); app.use("*", async (c, next) => { @@ -203,6 +217,9 @@ export async function createApp( ); }); app.get("/api/main-thread", async (c) => { + // Local (no Intelligence) mode: no platform thread exists; a stable local id + // is all the client needs — history comes from the local thread endpoints. + if (!intelligence) return c.json({ threadId: "local-main", existing: true }); const owner = c.get("owner"); await db.insertIfAbsent(owner, "conversation-settings", { id: "main", diff --git a/apps/server/src/config.ts b/apps/server/src/config.ts index c8efda0b3..f3e7dec3b 100644 --- a/apps/server/src/config.ts +++ b/apps/server/src/config.ts @@ -43,6 +43,8 @@ export interface Config { agentUrl?: string; agentToken?: string; intelligenceApiKey?: string; + intelligenceApiUrl?: string; + intelligenceWsUrl?: string; googleClientId?: string; googleClientSecret?: string; googleRedirectUri: string; @@ -69,14 +71,22 @@ export function required(name: string, message: string, value = process.env[name return value.trim(); } -export function assertApiDeploymentConfig( - config: Config, -): asserts config is Config & { intelligenceApiKey: string } { - required( - "CPK_INTELLIGENCE_API_KEY", - intelligenceKeyRequiredMessage, - config.intelligenceApiKey ?? "", - ); +/** + * Intelligence is optional: with the official cloud key (or a self-hosted + * Intelligence shim) the runtime persists threads and streams over its + * realtime WS; without one it falls back to the SSE runtime with local thread + * endpoints. A self-hosted endpoint must set apiUrl and wsUrl together — + * configuring only one would silently leave the other half pointed at the + * managed cloud. + */ +export function assertApiDeploymentConfig(config: Config): void { + if (Boolean(config.intelligenceApiUrl) !== Boolean(config.intelligenceWsUrl)) + throw new Error("INTELLIGENCE_API_URL and INTELLIGENCE_WS_URL must be set together"); +} + +/** Whether Intelligence-backed threads/realtime are enabled at all. */ +export function intelligenceConfigured(config: Config): boolean { + return Boolean(config.intelligenceApiKey ?? config.intelligenceApiUrl); } /** Accept a full worker URL, or host:port from a platform that omits the scheme. */ @@ -127,7 +137,9 @@ export function readConfig(): Config { agentBackend: backend, agentUrl: process.env.AGENT_URL, agentToken: process.env.AGENT_TOKEN, - intelligenceApiKey: required("CPK_INTELLIGENCE_API_KEY", intelligenceKeyRequiredMessage), + intelligenceApiKey: process.env.CPK_INTELLIGENCE_API_KEY?.trim() || undefined, + intelligenceApiUrl: process.env.INTELLIGENCE_API_URL?.trim() || undefined, + intelligenceWsUrl: process.env.INTELLIGENCE_WS_URL?.trim() || undefined, googleClientId: process.env.GOOGLE_CLIENT_ID, googleClientSecret: process.env.GOOGLE_CLIENT_SECRET, googleRedirectUri: `${publicUrl}/api/google/callback`, diff --git a/apps/server/src/engine/durable-runner.ts b/apps/server/src/engine/durable-runner.ts new file mode 100644 index 000000000..613410040 --- /dev/null +++ b/apps/server/src/engine/durable-runner.ts @@ -0,0 +1,241 @@ +import type { BaseEvent, Message } from "@ag-ui/core"; +import { + AgentRunner, + type AgentRunnerConnectRequest, + type AgentRunnerIsRunningRequest, + type AgentRunnerRunRequest, + type AgentRunnerStopRequest, + InMemoryAgentRunner, + type LocalThreadEndpointRecord, + type LocalThreadEndpointRunner, + ɵGLOBAL_STORE, +} from "@copilotkit/runtime/v2"; +import { defer, from, mergeAll, type Observable, tap } from "rxjs"; +import type { Store } from "../db.ts"; +import { backgroundFailure } from "../log.ts"; + +/** + * Keeps the runtime's local thread endpoints (/threads list, /threads/:id + * messages, /state) alive across restarts. + * + * The official SSE mode stores threads in a process-level in-memory runner, + * so a restart loses every conversation. This wraps InMemoryAgentRunner: + * - reads are served from our own PGlite store (threads/messages/events/state); + * - writes snapshot before each run starts and again when it completes, so a + * crash mid-run still leaves the previous run intact. + * + * Rows are stored under a fixed owner: the app is a single-user deployment and + * the runtime's runner interface carries no owner. + */ +const OWNER = "local-user"; +const THREADS = "chat-threads"; +const MESSAGES = "chat-messages"; +const EVENTS = "chat-events"; +const STATE = "chat-state"; + +type StoredMessages = { id: string; messages: Message[] }; +type StoredEvents = { id: string; events: BaseEvent[] }; +type StoredState = { id: string; state: Record | null }; + +export class DurableAgentRunner extends AgentRunner implements LocalThreadEndpointRunner { + readonly ɵsupportsLocalThreadEndpoints = true as const; + private readonly inner = new InMemoryAgentRunner(); + private readonly threads = new Map(); + private readonly messages = new Map(); + private readonly events = new Map(); + private readonly states = new Map | null>(); + private readonly hydrating: Promise; + + constructor(private readonly store: Store) { + super(); + this.hydrating = this.hydrate(); + } + + /** Read persisted threads back into the in-memory cache at boot. Reads are synchronous, so writes wait on `hydrating`. */ + private async hydrate(): Promise { + try { + for (const thread of await this.store.list(OWNER, THREADS)) + this.threads.set(thread.id, thread); + for (const row of await this.store.list(OWNER, MESSAGES)) + this.messages.set(row.id, row.messages ?? []); + for (const row of await this.store.list(OWNER, EVENTS)) + this.events.set(row.id, row.events ?? []); + for (const row of await this.store.list(OWNER, STATE)) + this.states.set(row.id, row.state ?? null); + } catch (error) { + backgroundFailure("durable-runner-hydrate", error); + } + } + + /** Merge event streams by runId: stored history first, same-runId runs resolved to the live (newest) copy. */ + private mergeEvents(stored: BaseEvent[], live: BaseEvent[]): BaseEvent[] { + if (!stored.length) return live; + if (!live.length) return stored; + const group = (events: BaseEvent[], tag: string) => { + const groups = new Map(); + let current = ""; + for (const [index, event] of events.entries()) { + const runId = (event as { runId?: string }).runId; + if ((event as { type?: string }).type === "RUN_STARTED" || !current) + current = runId ?? `${tag}:${index}`; + const bucket = groups.get(current) ?? []; + bucket.push(event); + groups.set(current, bucket); + } + return groups; + }; + const storedGroups = group(stored, "stored"); + const liveGroups = group(live, "live"); + const merged: BaseEvent[] = []; + for (const key of new Set([...storedGroups.keys(), ...liveGroups.keys()])) { + const events = liveGroups.get(key) ?? storedGroups.get(key); + if (events?.length) merged.push(...events); + } + return merged; + } + + /** Merge messages by id: first-seen order wins, same id resolved to the live copy. */ + private mergeMessages(stored: Message[], live: Message[]): Message[] { + if (!stored.length) return live; + if (!live.length) return stored; + const byId = new Map(); + for (const message of stored) byId.set(message.id, message); + for (const message of live) byId.set(message.id, message); + return [...byId.values()]; + } + + private async persist(threadId: string): Promise { + try { + const live = this.inner.listThreads().find((thread) => thread.id === threadId); + const previous = this.threads.get(threadId); + // After a restart the in-memory runner starts empty; overwriting directly + // would erase history — merge by runId/message id instead. + const messages = this.mergeMessages( + this.messages.get(threadId) ?? [], + this.inner.getThreadMessages(threadId), + ); + const events = this.mergeEvents( + this.events.get(threadId) ?? [], + this.inner.getThreadEvents(threadId), + ); + const state = this.inner.getThreadState(threadId) ?? this.states.get(threadId) ?? null; + const now = new Date().toISOString(); + const thread: LocalThreadEndpointRecord = { + id: threadId, + name: live?.name ?? previous?.name ?? null, + agentId: live?.agentId ?? previous?.agentId ?? "default", + organizationId: "", + createdById: "", + archived: false, + createdAt: previous?.createdAt ?? live?.createdAt ?? now, + updatedAt: now, + }; + this.threads.set(threadId, thread); + this.messages.set(threadId, messages); + this.events.set(threadId, events); + this.states.set(threadId, state); + await this.store.put(OWNER, THREADS, thread); + await this.store.put(OWNER, MESSAGES, { id: threadId, messages } satisfies StoredMessages); + await this.store.put(OWNER, EVENTS, { id: threadId, events } satisfies StoredEvents); + await this.store.put(OWNER, STATE, { id: threadId, state } satisfies StoredState); + } catch (error) { + // Persistence failures must never break the conversation itself. + backgroundFailure("durable-runner-persist", error); + } + } + + /** + * After a restart the in-memory runner is empty: replay persisted events and + * messages as one "historical run" so connect can rehydrate the conversation + * and the next run carries full context. Idempotent — skipped when the + * in-memory store already has events for the thread. + */ + private seed(threadId: string): void { + if (this.inner.getThreadEvents(threadId).length) return; + const events = this.events.get(threadId) ?? []; + const messages = this.messages.get(threadId) ?? []; + if (!events.length && !messages.length) return; + // appendRun silently drops unknown threads, so getOrCreate first. + ɵGLOBAL_STORE.getOrCreate(threadId); + ɵGLOBAL_STORE.appendRun(threadId, { + threadId, + runId: `seed:${threadId}`, + agentId: this.threads.get(threadId)?.agentId ?? "default", + parentRunId: null, + events, + messages, + createdAt: Date.now(), + }); + } + + run(request: AgentRunnerRunRequest): Observable { + return defer(() => + from( + (async () => { + await this.hydrating; + this.seed(request.threadId); + void this.persist(request.threadId); + return this.inner + .run(request) + .pipe(tap({ complete: () => void this.persist(request.threadId) })); + })(), + ).pipe(mergeAll()), + ); + } + + connect(request: AgentRunnerConnectRequest): Observable { + return defer(() => + from( + (async () => { + await this.hydrating; + if (request.threadId) this.seed(request.threadId); + return this.inner.connect(request); + })(), + ).pipe(mergeAll()), + ); + } + + isRunning(request: AgentRunnerIsRunningRequest): Promise { + return this.inner.isRunning(request); + } + + async stop(request: AgentRunnerStopRequest): Promise { + const stopped = await this.inner.stop(request); + await this.persist(request.threadId); + return stopped; + } + + listThreads(): LocalThreadEndpointRecord[] { + return [...this.threads.values()].sort((a, b) => (a.updatedAt < b.updatedAt ? 1 : -1)); + } + + getThreadMessages(threadId: string): Message[] { + return this.messages.get(threadId) ?? []; + } + + getThreadEvents(threadId: string): BaseEvent[] { + return this.events.get(threadId) ?? []; + } + + getThreadState(threadId: string): Record | null { + return this.states.get(threadId) ?? null; + } + + clearThreads(): void { + this.inner.clearThreads(); + this.threads.clear(); + this.messages.clear(); + this.events.clear(); + this.states.clear(); + void (async () => { + try { + for (const kind of [THREADS, MESSAGES, EVENTS, STATE]) { + for (const row of await this.store.list<{ id: string }>(OWNER, kind)) + await this.store.remove(OWNER, kind, row.id); + } + } catch (error) { + backgroundFailure("durable-runner-clear", error); + } + })(); + } +} diff --git a/apps/server/src/workspace.ts b/apps/server/src/workspace.ts index 6d5aafd22..efc154791 100644 --- a/apps/server/src/workspace.ts +++ b/apps/server/src/workspace.ts @@ -323,6 +323,9 @@ export class WorkspaceService { provider: this.config.agentBackend === "sample" ? "sample" : "model", configured: agentConfigured(this.config), openbotConfigured: false, + // Rich threads (multi-session sidebar / history hydration) work in both + // modes: with Intelligence they use the cloud, without it they use the + // runtime's local thread endpoints backed by DurableAgentRunner. richThreads: true, }, }; diff --git a/tests/config.test.ts b/tests/config.test.ts index 148720749..35b9405aa 100644 --- a/tests/config.test.ts +++ b/tests/config.test.ts @@ -27,23 +27,36 @@ function liveConfig(intelligenceApiKey?: string): Config { }; } -const missingKeyMessage = - "OpenMuse requires CPK_INTELLIGENCE_API_KEY. " + - "Run `npx copilotkit@latest login` and `npx copilotkit@latest project select`, " + - "then set the generated server-only key. " + - "See https://docs.copilotkit.ai/intelligence/connect-your-runtime"; - -test("every API mode rejects a missing or blank Intelligence key", () => { +test("every API mode runs without an Intelligence key (local SSE mode)", () => { for (const mode of [sampleConfig, liveConfig()]) { for (const key of [undefined, "", " \t\n"]) { - assert.throws(() => assertApiDeploymentConfig({ ...mode, intelligenceApiKey: key }), { - name: "Error", - message: missingKeyMessage, - }); + assert.doesNotThrow(() => assertApiDeploymentConfig({ ...mode, intelligenceApiKey: key })); } } }); +test("a self-hosted Intelligence endpoint must set apiUrl and wsUrl together", () => { + for (const mode of [sampleConfig, liveConfig()]) { + assert.throws( + () => + assertApiDeploymentConfig({ + ...mode, + intelligenceApiKey: "local", + intelligenceApiUrl: "https://shim.local", + }), + { name: "Error" }, + ); + assert.doesNotThrow(() => + assertApiDeploymentConfig({ + ...mode, + intelligenceApiKey: "local", + intelligenceApiUrl: "https://shim.local", + intelligenceWsUrl: "wss://shim.local", + }), + ); + } +}); + test("every API mode accepts a non-empty Intelligence key", () => { for (const mode of [sampleConfig, liveConfig()]) { assert.doesNotThrow(() =>