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
34 changes: 24 additions & 10 deletions apps/server/src/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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;
Expand All @@ -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" });
}
23 changes: 20 additions & 3 deletions apps/server/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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) => {
Expand Down Expand Up @@ -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",
Expand Down
30 changes: 21 additions & 9 deletions apps/server/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,8 @@ export interface Config {
agentUrl?: string;
agentToken?: string;
intelligenceApiKey?: string;
intelligenceApiUrl?: string;
intelligenceWsUrl?: string;
googleClientId?: string;
googleClientSecret?: string;
googleRedirectUri: string;
Expand All @@ -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. */
Expand Down Expand Up @@ -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`,
Expand Down
241 changes: 241 additions & 0 deletions apps/server/src/engine/durable-runner.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown> | null };

export class DurableAgentRunner extends AgentRunner implements LocalThreadEndpointRunner {
readonly ɵsupportsLocalThreadEndpoints = true as const;
private readonly inner = new InMemoryAgentRunner();
private readonly threads = new Map<string, LocalThreadEndpointRecord>();
private readonly messages = new Map<string, Message[]>();
private readonly events = new Map<string, BaseEvent[]>();
private readonly states = new Map<string, Record<string, unknown> | null>();
private readonly hydrating: Promise<void>;

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<void> {
try {
for (const thread of await this.store.list<LocalThreadEndpointRecord>(OWNER, THREADS))
this.threads.set(thread.id, thread);
for (const row of await this.store.list<StoredMessages>(OWNER, MESSAGES))
this.messages.set(row.id, row.messages ?? []);
for (const row of await this.store.list<StoredEvents>(OWNER, EVENTS))
this.events.set(row.id, row.events ?? []);
for (const row of await this.store.list<StoredState>(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<string, BaseEvent[]>();
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<string, Message>();
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<void> {
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<BaseEvent> {
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<BaseEvent> {
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<boolean> {
return this.inner.isRunning(request);
}

async stop(request: AgentRunnerStopRequest): Promise<boolean | undefined> {
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<string, unknown> | 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);
}
})();
}
}
Loading