From a13b19fa45df33b5e5e74bada3ada343a974af85 Mon Sep 17 00:00:00 2001 From: Johnny Huynh <27847622+johnnyhuy@users.noreply.github.com> Date: Thu, 6 Aug 2026 08:41:35 +1000 Subject: [PATCH] fix(protocol,server): rpc_response discriminator + server-side rpc routing RpcResponseSchema gains kind: "rpc_response" so SupaplaneClient.rpc() can actually match replies. Adds an RpcRouter (provider.list/models/ modes/diagnostic, workspace.list, session.list) wired into the daemon, plus E2E coverage of the rpc round-trip. Co-authored-by: opencode-agent --- packages/protocol/src/messages.ts | 1 + packages/server/src/daemon.ts | 4 + .../server/src/server/agent/agent-manager.ts | 16 ++++ .../server/src/server/daemon-e2e/rpc.test.ts | 58 ++++++++++++++ packages/server/src/server/rpc-router.ts | 75 +++++++++++++++++++ packages/server/src/websocket-server.ts | 19 ++++- 6 files changed, 172 insertions(+), 1 deletion(-) create mode 100644 packages/server/src/server/daemon-e2e/rpc.test.ts create mode 100644 packages/server/src/server/rpc-router.ts diff --git a/packages/protocol/src/messages.ts b/packages/protocol/src/messages.ts index 788c66b..b8e5d65 100644 --- a/packages/protocol/src/messages.ts +++ b/packages/protocol/src/messages.ts @@ -303,6 +303,7 @@ export const RpcRequestSchema = z.object({ export type RpcRequest = z.infer; export const RpcResponseSchema = z.object({ + kind: z.literal("rpc_response"), rpc: z.string(), requestId: z.string(), ok: z.boolean(), diff --git a/packages/server/src/daemon.ts b/packages/server/src/daemon.ts index 95a4e73..e251e84 100644 --- a/packages/server/src/daemon.ts +++ b/packages/server/src/daemon.ts @@ -10,6 +10,7 @@ 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 { RpcRouter } from "./server/rpc-router.js"; import { WorkspaceRegistry } from "./server/workspace-registry.js"; import { SupaplaneWebsocketServer } from "./websocket-server.js"; @@ -97,6 +98,9 @@ export async function startDaemon(args?: { }), ); + const rpcRouter = new RpcRouter({ workspaces, agents: agentManager, logger }); + wsServer.setRpcHandler((req) => rpcRouter.handle(req)); + await new Promise((resolve, reject) => { const onError = (err: Error) => { httpServer.off("listening", onListening); diff --git a/packages/server/src/server/agent/agent-manager.ts b/packages/server/src/server/agent/agent-manager.ts index 24427a0..5dd3098 100644 --- a/packages/server/src/server/agent/agent-manager.ts +++ b/packages/server/src/server/agent/agent-manager.ts @@ -88,6 +88,22 @@ export class AgentManager { return provider; } + listModels(providerId: string) { + return this.getProvider(providerId).listModels(); + } + + listModes(providerId: string) { + return this.getProvider(providerId).listModes(); + } + + async getDiagnostic(providerId: string): Promise<{ diagnostic: string }> { + const provider = this.getProvider(providerId); + if (!provider.getDiagnostic) { + return { diagnostic: "no diagnostic available" }; + } + return provider.getDiagnostic(); + } + async startSession(args: StartSessionArgs): Promise { const provider = this.getProvider(args.providerId); const sessionId = newSessionId(); diff --git a/packages/server/src/server/daemon-e2e/rpc.test.ts b/packages/server/src/server/daemon-e2e/rpc.test.ts new file mode 100644 index 0000000..25b68d4 --- /dev/null +++ b/packages/server/src/server/daemon-e2e/rpc.test.ts @@ -0,0 +1,58 @@ +import { mkdtemp } from "node:fs/promises"; +import type { AddressInfo } from "node:net"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +import { SupaplaneClient } from "@echohello/client"; +import { afterEach, beforeEach, describe, expect, it } from "vitest"; + +import { startDaemon, type DaemonHandle } from "../../daemon.js"; + +describe("daemon e2e: rpc", () => { + let daemon: DaemonHandle; + let client: SupaplaneClient; + + beforeEach(async () => { + const supaplaneHome = await mkdtemp(join(tmpdir(), "supaplane-e2e-rpc-")); + daemon = await startDaemon({ + config: { listenPort: 0, logLevel: "error" }, + supaplaneHome, + }); + const { port } = daemon.httpServer.address() as AddressInfo; + client = new SupaplaneClient({ + endpoint: `ws://127.0.0.1:${port}`, + clientId: "e2e-rpc-client", + clientType: "cli", + reconnect: false, + }); + await client.connect(); + }); + + afterEach(async () => { + client.close(); + await daemon.stop(); + }); + + it("provider.list round-trips", async () => { + const result = await client.rpc("provider.list"); + expect(result.providers).toContain("claude"); + }); + + it("provider.models returns the claude model list", async () => { + const result = await client.rpc<{ providerId: string }, { models: { id: string }[] }>( + "provider.models", + { providerId: "claude" }, + ); + expect(result.models.map((m) => m.id)).toContain("sonnet"); + }); + + it("unknown rpc rejects with an error", async () => { + await expect(client.rpc("no.such.rpc")).rejects.toThrow("Unknown rpc"); + }); + + it("unknown provider rejects with an error", async () => { + await expect(client.rpc("provider.models", { providerId: "does-not-exist" })).rejects.toThrow( + "Unknown provider", + ); + }); +}); diff --git a/packages/server/src/server/rpc-router.ts b/packages/server/src/server/rpc-router.ts new file mode 100644 index 0000000..4964504 --- /dev/null +++ b/packages/server/src/server/rpc-router.ts @@ -0,0 +1,75 @@ +import { z } from "zod"; +import type { Logger } from "pino"; +import { SupaplaneError, type RpcRequest, type RpcResponse } from "@echohello/protocol"; + +import type { AgentManager } from "./agent/agent-manager.js"; +import type { WorkspaceRegistry } from "./workspace-registry.js"; + +export interface RpcRouterOptions { + workspaces: WorkspaceRegistry; + agents: AgentManager; + logger: Logger; +} + +const ProviderArgsSchema = z.object({ providerId: z.string().min(1) }); +const SessionListArgsSchema = z.object({ workspaceId: z.string().optional() }); + +/** + * Answers one-shot client RPCs (`SupaplaneClient.rpc()`). Read-only views + * over the provider registry, workspace registry, and session table. + * Unknown RPCs get an `ok: false` response, never silence. + */ +export class RpcRouter { + readonly #workspaces: WorkspaceRegistry; + readonly #agents: AgentManager; + readonly #logger: Logger; + + constructor(options: RpcRouterOptions) { + this.#workspaces = options.workspaces; + this.#agents = options.agents; + this.#logger = options.logger.child({ module: "rpc-router" }); + } + + async handle(req: RpcRequest): Promise { + try { + const result = await this.#route(req); + return { kind: "rpc_response", rpc: req.rpc, requestId: req.requestId, ok: true, result }; + } catch (err) { + const code = err instanceof SupaplaneError ? err.code : "internal"; + const message = err instanceof Error ? err.message : String(err); + this.#logger.warn({ rpc: req.rpc, code, err: message }, "rpc failed"); + return { + kind: "rpc_response", + rpc: req.rpc, + requestId: req.requestId, + ok: false, + error: { code, message }, + }; + } + } + + async #route(req: RpcRequest): Promise { + switch (req.rpc) { + case "provider.list": + return { providers: this.#agents.providerIds() }; + case "provider.models": + return { models: await this.#agents.listModels(this.#providerId(req)) }; + case "provider.modes": + return { modes: await this.#agents.listModes(this.#providerId(req)) }; + case "provider.diagnostic": + return this.#agents.getDiagnostic(this.#providerId(req)); + case "workspace.list": + return { workspaces: this.#workspaces.list() }; + case "session.list": { + const args = SessionListArgsSchema.parse(req.args ?? {}); + return { sessions: this.#agents.listSessions(args.workspaceId) }; + } + default: + throw new SupaplaneError({ code: "not_found", message: `Unknown rpc: ${req.rpc}` }); + } + } + + #providerId(req: RpcRequest): string { + return ProviderArgsSchema.parse(req.args).providerId; + } +} diff --git a/packages/server/src/websocket-server.ts b/packages/server/src/websocket-server.ts index 53c2948..e3feeb5 100644 --- a/packages/server/src/websocket-server.ts +++ b/packages/server/src/websocket-server.ts @@ -8,6 +8,8 @@ import { type ClientType, EnvelopeSchema, type HelloAckMessage, + type RpcRequest, + type RpcResponse, type ServerEvent, } from "@echohello/protocol"; @@ -29,6 +31,7 @@ export interface WebsocketServerOptions { daemonLabel?: string; serverVersion: string; onCommand?: (cmd: ClientCommand, session: SessionRecord) => void; + onRpc?: (req: RpcRequest, session: SessionRecord) => Promise; /** Provider ids advertised in the `hello_ack` capabilities. */ providers?: readonly string[]; } @@ -51,6 +54,7 @@ export class SupaplaneWebsocketServer { #serverVersion: string; #daemonLabel: string | undefined; #onCommand?: (cmd: ClientCommand, session: SessionRecord) => void; + #onRpc?: (req: RpcRequest, session: SessionRecord) => Promise; #providers: readonly string[]; constructor(options: WebsocketServerOptions) { @@ -86,6 +90,11 @@ export class SupaplaneWebsocketServer { this.#onCommand = handler; } + /** Install the RPC handler after construction. */ + setRpcHandler(handler: (req: RpcRequest, session: SessionRecord) => Promise): void { + this.#onRpc = handler; + } + /** Broadcast a server event to all connected sessions that have subscribed to its topic. */ broadcast(event: ServerEvent): void { const payload = JSON.stringify({ event }); @@ -175,8 +184,16 @@ export class SupaplaneWebsocketServer { return; } this.#onCommand?.(cmd, session); + return; + } + if ("rpc" in envelope && !("kind" in envelope) && this.#onRpc) { + void this.#onRpc(envelope, session) + .then((response) => this.sendTo(socket, response)) + .catch((err: unknown) => { + this.#logger.warn({ err, rpc: envelope.rpc }, "rpc handler failed"); + }); + return; } - // RPC request/response envelopes are routed via SupaplaneClient.rpc(), not here. // Server events flow outbound via broadcast(). `hello`/`hello_ack` are // handshake-only and are caught at the connection boundary above. } catch (err) {