Skip to content
Merged
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
68 changes: 68 additions & 0 deletions bun.lock

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions packages/server/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
"lint": "oxlint src --config ../../oxlint.jsonc"
},
"dependencies": {
"@anthropic-ai/claude-agent-sdk": "^0.3.222",
"@echohello/client": "workspace:*",
"@echohello/protocol": "workspace:*",
"dotenv": "^17.2.3",
Expand Down
38 changes: 38 additions & 0 deletions packages/server/src/daemon.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,11 @@ import { loadOrCreateIdentity } from "./handshake.js";
import { createHttpApp } from "./http-app.js";
import { createLogger } from "./logger.js";
import { resolveSupaplaneHome, SUPAPLANE_VERSION } from "./paths.js";
import { AgentManager } from "./server/agent/agent-manager.js";
import { HandleStore } from "./server/agent/handle-store.js";
import { ClaudeAgentClient } from "./server/agent/providers/claude/claude-provider.js";
import { CommandDispatcher } from "./server/command-dispatcher.js";
import { WorkspaceRegistry } from "./server/workspace-registry.js";
import { SupaplaneWebsocketServer } from "./websocket-server.js";

export interface DaemonHandle {
Expand All @@ -15,6 +20,8 @@ export interface DaemonHandle {
stop: () => Promise<void>;
httpServer: ReturnType<typeof createServer>;
wsServer: SupaplaneWebsocketServer;
agentManager: AgentManager;
workspaces: WorkspaceRegistry;
}

/**
Expand Down Expand Up @@ -54,14 +61,42 @@ export async function startDaemon(args?: {
});

const httpServer = createServer(httpApp);

const handleStore = new HandleStore(supaplaneHome);
const agentManager = new AgentManager({ handleStore, logger });
agentManager.registerProvider(new ClaudeAgentClient());
const workspaces = new WorkspaceRegistry();

const wsServer = new SupaplaneWebsocketServer({
httpServer,
logger,
identity,
...(config.daemonAuthToken ? { authToken: config.daemonAuthToken } : {}),
serverVersion: SUPAPLANE_VERSION,
providers: agentManager.providerIds(),
});

agentManager.onAgentEvent = (event) => wsServer.broadcast({ kind: "event", event });
agentManager.onSessionState = (session) => wsServer.broadcast({ kind: "session_state", session });

const dispatcher = new CommandDispatcher({
workspaces,
agents: agentManager,
broadcast: (event) => wsServer.broadcast(event),
logger,
});
wsServer.setCommandHandler((cmd, session) =>
dispatcher.handle(cmd, {
clientId: session.clientId,
sendError: (error) =>
wsServer.sendTo(session.socket, {
type: "error",
code: error.code,
message: error.message,
}),
}),
);

await new Promise<void>((resolve, reject) => {
const onError = (err: Error) => {
httpServer.off("listening", onListening);
Expand All @@ -84,8 +119,11 @@ export async function startDaemon(args?: {
supaplaneHome,
httpServer,
wsServer,
agentManager,
workspaces,
async stop(): Promise<void> {
logger.info("stopping daemon");
await agentManager.disposeAll();
await new Promise<void>((resolve, reject) => {
httpServer.close((err) => (err ? reject(err) : resolve()));
});
Expand Down
4 changes: 3 additions & 1 deletion packages/server/src/handshake.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@ import { getOrCreateServerId } from "./server-id.js";
export interface WebsocketSessionContext {
serverId: string;
serverVersion: string;
/** Provider ids advertised in the `hello_ack` capabilities. */
providers: readonly string[];
daemonLabel?: string;
logger: Logger;
}
Expand Down Expand Up @@ -51,7 +53,7 @@ export function handleHello(args: {
serverVersion: args.ctx.serverVersion,
protocolVersion: SUPAPLANE_PROTOCOL_VERSION,
capabilities: {
providers: ["opencode", "claude", "cursor"],
providers: [...args.ctx.providers],
relay: false,
worktrees: false,
scheduling: false,
Expand Down
252 changes: 252 additions & 0 deletions packages/server/src/server/agent/agent-manager.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,252 @@
import { randomBytes } from "node:crypto";

import type { Logger } from "pino";
import { AgentError, type AgentEvent, type SessionState } from "@echohello/protocol";

import type {
AgentClient,
AgentEventSink,
AgentSession,
PersistenceHandle,
SessionScopedEvent,
} from "./agent-sdk-types.js";
import { HandleStore } from "./handle-store.js";

export interface AgentManagerOptions {
handleStore: HandleStore;
logger: Logger;
}

interface ManagedSession {
state: SessionState;
session: AgentSession;
handle: PersistenceHandle;
}

export interface StartSessionArgs {
workspaceId: string;
cwd: string;
providerId: string;
modelId?: string;
modeId?: string;
initialPrompt?: string;
}

export interface ResumeSessionArgs {
workspaceId: string;
cwd: string;
handle: PersistenceHandle;
overrides?: {
modelId?: string | undefined;
modeId?: string | undefined;
};
}

/**
* Owns agent session lifecycle: provider registry, Supaplane-side session ids,
* event stamping/forwarding, status bookkeeping, and handle persistence.
*
* Wire events out via `onAgentEvent` / `onSessionState` (set by the daemon
* before any session can be created).
*/
export class AgentManager {
readonly #providers = new Map<string, AgentClient>();
readonly #sessions = new Map<string, ManagedSession>();
readonly #handleStore: HandleStore;
readonly #logger: Logger;

onAgentEvent: ((event: AgentEvent) => void) | undefined;
onSessionState: ((state: SessionState) => void) | undefined;

constructor(options: AgentManagerOptions) {
this.#handleStore = options.handleStore;
this.#logger = options.logger.child({ module: "agent-manager" });
}

registerProvider(client: AgentClient): void {
if (this.#providers.has(client.providerId)) {
throw new AgentError({
code: "conflict",
message: `Provider already registered: ${client.providerId}`,
});
}
this.#providers.set(client.providerId, client);
}

providerIds(): string[] {
return [...this.#providers.keys()];
}

getProvider(providerId: string): AgentClient {
const provider = this.#providers.get(providerId);
if (!provider) {
throw new AgentError({
code: "provider_unavailable",
message: `Unknown provider: ${providerId}`,
});
}
return provider;
}

async startSession(args: StartSessionArgs): Promise<SessionState> {
const provider = this.getProvider(args.providerId);
const sessionId = newSessionId();
const emit = this.#makeSink(sessionId);

const { session, handle } = await provider.createSession(
{
cwd: args.cwd,
...(args.modelId !== undefined ? { modelId: args.modelId } : {}),
...(args.modeId !== undefined ? { modeId: args.modeId } : {}),
},
emit,
);

const now = Date.now();
const state: SessionState = {
sessionId,
workspaceId: args.workspaceId,
providerId: args.providerId,
...(args.modelId !== undefined ? { modelId: args.modelId } : {}),
...(args.modeId !== undefined ? { modeId: args.modeId } : {}),
status: "idle",
startedAt: now,
updatedAt: now,
forkCount: 0,
};
this.#sessions.set(sessionId, { state, session, handle });
await this.#persistHandle(args.cwd, handle);
this.#emitSessionState(state);

if (args.initialPrompt !== undefined && args.initialPrompt.length > 0) {
void this.send(sessionId, args.initialPrompt).catch((err: unknown) => {
this.#logger.warn({ err, sessionId }, "initial prompt failed");
});
}
return state;
}

async resumeSession(args: ResumeSessionArgs): Promise<SessionState> {
const provider = this.getProvider(args.handle.provider);
const sessionId = newSessionId();
const emit = this.#makeSink(sessionId);

const stored = await this.#handleStore
.load(args.cwd, args.handle.provider, args.handle.sessionId)
.catch(() => null);
const metadata = { ...stored?.metadata, ...args.handle.metadata };

const { session, handle } = await provider.resumeSession(
{
handle: {
provider: args.handle.provider,
sessionId: args.handle.sessionId,
...(metadata !== undefined && Object.keys(metadata).length > 0 ? { metadata } : {}),
},
cwd: args.cwd,
...(args.overrides !== undefined ? { overrides: args.overrides } : {}),
},
emit,
);

const now = Date.now();
const state: SessionState = {
sessionId,
workspaceId: args.workspaceId,
providerId: args.handle.provider,
...(args.overrides?.modelId !== undefined ? { modelId: args.overrides.modelId } : {}),
...(args.overrides?.modeId !== undefined ? { modeId: args.overrides.modeId } : {}),
status: "idle",
startedAt: now,
updatedAt: now,
forkCount: 0,
};
this.#sessions.set(sessionId, { state, session, handle });
await this.#persistHandle(args.cwd, handle);
this.#emitSessionState(state);
return state;
}

async send(sessionId: string, prompt: string, attachments?: unknown[]): Promise<void> {
const managed = this.#requireSession(sessionId);
this.#setStatus(sessionId, "running");
try {
await managed.session.send(prompt, attachments);
} catch (err) {
this.#setStatus(sessionId, "error");
throw err;
}
}

async abort(sessionId: string): Promise<void> {
const managed = this.#requireSession(sessionId);
await managed.session.abort();
this.#setStatus(sessionId, "idle");
}

async disposeSession(sessionId: string): Promise<void> {
const managed = this.#sessions.get(sessionId);
if (!managed) return;
this.#sessions.delete(sessionId);
await managed.session.dispose().catch((err: unknown) => {
this.#logger.warn({ err, sessionId }, "session dispose failed");
});
}

async disposeAll(): Promise<void> {
const ids = [...this.#sessions.keys()];
await Promise.all(ids.map((id) => this.disposeSession(id)));
}

getSession(sessionId: string): SessionState | undefined {
return this.#sessions.get(sessionId)?.state;
}

listSessions(workspaceId?: string): SessionState[] {
const all = [...this.#sessions.values()].map((s) => s.state);
return workspaceId === undefined ? all : all.filter((s) => s.workspaceId === workspaceId);
}

#makeSink(sessionId: string): AgentEventSink {
return (event: SessionScopedEvent) => {
const stamped = { ...event, sessionId } as AgentEvent;
if (stamped.type === "status") {
this.#setStatus(sessionId, stamped.status);
} else if (stamped.type === "error") {
this.#setStatus(sessionId, "error");
}
this.onAgentEvent?.(stamped);
};
}

#setStatus(sessionId: string, status: SessionState["status"]): void {
const managed = this.#sessions.get(sessionId);
if (!managed || managed.state.status === status) return;
managed.state = { ...managed.state, status, updatedAt: Date.now() };
this.#emitSessionState(managed.state);
}

#requireSession(sessionId: string): ManagedSession {
const managed = this.#sessions.get(sessionId);
if (!managed) {
throw new AgentError({ code: "not_found", message: `Unknown session: ${sessionId}` });
}
return managed;
}

async #persistHandle(cwd: string, handle: PersistenceHandle): Promise<void> {
try {
await this.#handleStore.save(cwd, handle);
} catch (err) {
this.#logger.warn({ err, cwd }, "failed to persist session handle");
}
}

#emitSessionState(state: SessionState): void {
this.onSessionState?.(state);
}
}

function newSessionId(): string {
return `ses_${randomBytes(9).toString("base64url")}`;
}
Loading
Loading