mirror of
https://github.com/getpaseo/paseo.git
synced 2026-07-29 12:01:31 +00:00
voice: switch internal MCP path to stdio unix-socket bridge
This commit is contained in:
@@ -14,6 +14,7 @@ import { runInspectCommand } from './commands/agent/inspect.js'
|
||||
import { runWaitCommand } from './commands/agent/wait.js'
|
||||
import { runAttachCommand } from './commands/agent/attach.js'
|
||||
import { withOutput } from './output/index.js'
|
||||
import { runVoiceMcpBridgeCommand } from './commands/voice-mcp-bridge.js'
|
||||
|
||||
const VERSION = '0.1.0'
|
||||
|
||||
@@ -141,5 +142,13 @@ export function createCli(): Command {
|
||||
// Worktree commands
|
||||
program.addCommand(createWorktreeCommand())
|
||||
|
||||
// Internal voice MCP stdio bridge command (hidden).
|
||||
program
|
||||
.command('__paseo_voice_mcp_bridge')
|
||||
.description('Internal voice MCP bridge command')
|
||||
.argument('<callerAgentId>')
|
||||
.requiredOption('--socket <path>')
|
||||
.action(runVoiceMcpBridgeCommand)
|
||||
|
||||
return program
|
||||
}
|
||||
|
||||
17
packages/cli/src/commands/voice-mcp-bridge.ts
Normal file
17
packages/cli/src/commands/voice-mcp-bridge.ts
Normal file
@@ -0,0 +1,17 @@
|
||||
import { runVoiceMcpBridgeCli } from "@getpaseo/server";
|
||||
|
||||
type VoiceBridgeOptions = {
|
||||
socket: string;
|
||||
};
|
||||
|
||||
export async function runVoiceMcpBridgeCommand(
|
||||
callerAgentId: string,
|
||||
options: VoiceBridgeOptions
|
||||
): Promise<void> {
|
||||
await runVoiceMcpBridgeCli([
|
||||
"--socket",
|
||||
options.socket,
|
||||
"--caller-agent-id",
|
||||
callerAgentId,
|
||||
]);
|
||||
}
|
||||
@@ -3,6 +3,7 @@ import { createServer as createHTTPServer } from "http";
|
||||
import { createReadStream, unlinkSync, existsSync } from "fs";
|
||||
import { stat } from "fs/promises";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import path from "node:path";
|
||||
import { StreamableHTTPServerTransport } from "@modelcontextprotocol/sdk/server/streamableHttp.js";
|
||||
import { InMemoryTransport } from "@modelcontextprotocol/sdk/inMemory.js";
|
||||
import { isInitializeRequest } from "@modelcontextprotocol/sdk/types.js";
|
||||
@@ -75,11 +76,35 @@ import type {
|
||||
} from "./agent/agent-sdk-types.js";
|
||||
import { acquirePidLock, releasePidLock } from "./pid-lock.js";
|
||||
import { isHostAllowed, type AllowedHostsConfig } from "./allowed-hosts.js";
|
||||
import { createVoiceMcpBridgeSocketServer, type VoiceMcpBridgeSocketServer } from "./voice-mcp-bridge.js";
|
||||
|
||||
type AgentMcpTransportMap = Map<string, StreamableHTTPServerTransport>;
|
||||
type VoiceAgentProvider = "claude" | "codex" | "opencode";
|
||||
const VOICE_AGENT_FALLBACK_ORDER: VoiceAgentProvider[] = ["claude", "codex", "opencode"];
|
||||
|
||||
function resolveVoiceMcpBridgeCommand(logger: Logger): { command: string; baseArgs: string[] } {
|
||||
const explicit = process.env.PASEO_BIN_PATH?.trim();
|
||||
if (explicit) {
|
||||
return { command: explicit, baseArgs: ["__paseo_voice_mcp_bridge"] };
|
||||
}
|
||||
|
||||
const argv1 = process.argv[1]?.trim();
|
||||
if (!argv1) {
|
||||
logger.warn("Could not resolve argv[1] for voice MCP bridge; falling back to 'paseo'");
|
||||
return { command: "paseo", baseArgs: ["__paseo_voice_mcp_bridge"] };
|
||||
}
|
||||
|
||||
const base = path.basename(argv1).toLowerCase();
|
||||
if (base.includes("tsx") && process.argv[2]) {
|
||||
return {
|
||||
command: process.execPath,
|
||||
baseArgs: [argv1, process.argv[2], "__paseo_voice_mcp_bridge"],
|
||||
};
|
||||
}
|
||||
|
||||
return { command: argv1, baseArgs: ["__paseo_voice_mcp_bridge"] };
|
||||
}
|
||||
|
||||
export type PaseoOpenAIConfig = {
|
||||
apiKey?: string;
|
||||
stt?: Partial<STTConfig> & { apiKey?: string };
|
||||
@@ -350,14 +375,8 @@ export async function createPaseoDaemon(
|
||||
},
|
||||
"Voice LLM provider reconciliation completed"
|
||||
);
|
||||
if (listenTarget.type !== "tcp" && requestedVoiceLlmProvider !== "openrouter") {
|
||||
logger.error(
|
||||
{ listen: config.listen, requestedVoiceLlmProvider },
|
||||
"Local voice agent mode requires TCP listen target for HTTP MCP bridge"
|
||||
);
|
||||
throw new Error("Local voice agent mode requires TCP listen target");
|
||||
}
|
||||
let wsServer: VoiceAssistantWebSocketServer | null = null;
|
||||
let voiceMcpBridgeServer: VoiceMcpBridgeSocketServer | null = null;
|
||||
|
||||
// Create in-memory transport for Session's Agent MCP client (voice assistant tools)
|
||||
const createInMemoryAgentMcpTransport = async (): Promise<InMemoryTransport> => {
|
||||
@@ -497,6 +516,25 @@ export async function createPaseoDaemon(
|
||||
logger.info("Agent MCP HTTP endpoint disabled");
|
||||
}
|
||||
|
||||
const voiceMcpSocketPath = path.join(config.paseoHome, "runtime", "voice-mcp.sock");
|
||||
const voiceMcpBridgeCommand = resolveVoiceMcpBridgeCommand(logger);
|
||||
voiceMcpBridgeServer = createVoiceMcpBridgeSocketServer({
|
||||
socketPath: voiceMcpSocketPath,
|
||||
logger,
|
||||
createAgentMcpServerForCaller: async (callerAgentId) => {
|
||||
return createAgentMcpServer({
|
||||
agentManager,
|
||||
agentStorage,
|
||||
paseoHome: config.paseoHome,
|
||||
callerAgentId,
|
||||
enableVoiceTools: false,
|
||||
resolveSpeakHandler: (agentId) => wsServer?.resolveVoiceSpeakHandler(agentId) ?? null,
|
||||
resolveCallerContext: (agentId) => wsServer?.resolveVoiceCallerContext(agentId) ?? null,
|
||||
logger,
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
|
||||
let sttService: SpeechToTextProvider | null = null;
|
||||
let ttsService: TextToSpeechProvider | null = null;
|
||||
@@ -875,11 +913,6 @@ export async function createPaseoDaemon(
|
||||
);
|
||||
}
|
||||
|
||||
const voiceAgentMcpUrl =
|
||||
listenTarget.type === "tcp"
|
||||
? `http://127.0.0.1:${listenTarget.port}/mcp/agents`
|
||||
: null;
|
||||
|
||||
wsServer = new VoiceAssistantWebSocketServer(
|
||||
httpServer,
|
||||
logger,
|
||||
@@ -899,7 +932,17 @@ export async function createPaseoDaemon(
|
||||
voiceLlmDefaultProvider,
|
||||
voiceLlmModel: config.voiceLlmModel ?? null,
|
||||
voiceLlmAvailability,
|
||||
voiceAgentMcpUrl,
|
||||
voiceAgentMcpStdio: {
|
||||
command: voiceMcpBridgeCommand.command,
|
||||
baseArgs: [
|
||||
...voiceMcpBridgeCommand.baseArgs,
|
||||
"--socket",
|
||||
voiceMcpSocketPath,
|
||||
],
|
||||
env: {
|
||||
PASEO_HOME: config.paseoHome,
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
finalTimeoutMs: config.dictationFinalTimeoutMs,
|
||||
@@ -994,6 +1037,9 @@ export async function createPaseoDaemon(
|
||||
httpServer.listen(listenTarget.path);
|
||||
}
|
||||
});
|
||||
if (voiceMcpBridgeServer) {
|
||||
await voiceMcpBridgeServer.start();
|
||||
}
|
||||
};
|
||||
|
||||
const stop = async () => {
|
||||
@@ -1012,6 +1058,9 @@ export async function createPaseoDaemon(
|
||||
if (wsServer) {
|
||||
await wsServer.close();
|
||||
}
|
||||
if (voiceMcpBridgeServer) {
|
||||
await voiceMcpBridgeServer.stop().catch(() => undefined);
|
||||
}
|
||||
await new Promise<void>((resolve) => {
|
||||
httpServer.close(() => resolve());
|
||||
});
|
||||
|
||||
@@ -5,6 +5,7 @@ export { resolvePaseoHome } from "./paseo-home.js";
|
||||
export { createRootLogger, type LogLevel, type LogFormat } from "./logger.js";
|
||||
export { loadPersistedConfig, type PersistedConfig } from "./persisted-config.js";
|
||||
export { DaemonClient, type DaemonClientConfig, type ConnectionState, type DaemonEvent } from "../client/daemon-client.js";
|
||||
export { runVoiceMcpBridgeCli } from "./voice-mcp-bridge.js";
|
||||
|
||||
// Agent SDK types for CLI commands
|
||||
export type {
|
||||
|
||||
@@ -10,8 +10,15 @@ import { resolvePaseoHome } from "./paseo-home.js";
|
||||
import { createRootLogger } from "./logger.js";
|
||||
import { loadPersistedConfig } from "./persisted-config.js";
|
||||
import { PidLockError } from "./pid-lock.js";
|
||||
import { runVoiceMcpBridgeCli } from "./voice-mcp-bridge.js";
|
||||
|
||||
async function main() {
|
||||
const bridgeArgIndex = process.argv.findIndex((arg) => arg === "__paseo_voice_mcp_bridge");
|
||||
if (bridgeArgIndex >= 0) {
|
||||
await runVoiceMcpBridgeCli(process.argv.slice(bridgeArgIndex + 1));
|
||||
return;
|
||||
}
|
||||
|
||||
let paseoHome: string;
|
||||
let logger: ReturnType<typeof createRootLogger>;
|
||||
let config: ReturnType<typeof loadConfig>;
|
||||
|
||||
@@ -57,6 +57,11 @@ type VoiceCallerContext = {
|
||||
allowCustomCwd?: boolean;
|
||||
enableVoiceTools?: boolean;
|
||||
};
|
||||
type VoiceMcpStdioConfig = {
|
||||
command: string;
|
||||
baseArgs: string[];
|
||||
env?: Record<string, string>;
|
||||
};
|
||||
import { buildProviderRegistry } from "./agent/provider-registry.js";
|
||||
import { AgentManager } from "./agent/agent-manager.js";
|
||||
import type { ManagedAgent } from "./agent/agent-manager.js";
|
||||
@@ -158,6 +163,25 @@ const VOICE_AGENT_SYSTEM_INSTRUCTION = [
|
||||
"Only use the paseo MCP tools.",
|
||||
].join(" ");
|
||||
|
||||
export function buildVoiceAgentMcpServerConfig(params: {
|
||||
callerAgentId: string;
|
||||
command: string;
|
||||
baseArgs: string[];
|
||||
env?: Record<string, string>;
|
||||
}): {
|
||||
type: "stdio";
|
||||
command: string;
|
||||
args: string[];
|
||||
env?: Record<string, string>;
|
||||
} {
|
||||
return {
|
||||
type: "stdio",
|
||||
command: params.command,
|
||||
args: [...params.baseArgs, "--caller-agent-id", params.callerAgentId],
|
||||
...(params.env ? { env: params.env } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
type ProcessingPhase = "idle" | "transcribing" | "llm";
|
||||
|
||||
type NormalizedGitOptions = {
|
||||
@@ -350,7 +374,7 @@ export class Session {
|
||||
private readonly voiceLlmDefaultProvider: VoiceAgentProvider | null;
|
||||
private readonly voiceLlmModel: string | null;
|
||||
private readonly voiceLlmAvailability: Record<VoiceAgentProvider, boolean> | null;
|
||||
private readonly voiceAgentMcpUrl: string | null;
|
||||
private readonly voiceAgentMcpStdio: VoiceMcpStdioConfig | null;
|
||||
private readonly registerVoiceSpeakHandler?: (
|
||||
agentId: string,
|
||||
handler: VoiceSpeakHandler
|
||||
@@ -383,7 +407,7 @@ export class Session {
|
||||
voiceLlmDefaultProvider?: VoiceAgentProvider | null;
|
||||
voiceLlmModel?: string | null;
|
||||
voiceLlmAvailability?: Record<VoiceAgentProvider, boolean> | null;
|
||||
voiceAgentMcpUrl?: string | null;
|
||||
voiceAgentMcpStdio?: VoiceMcpStdioConfig | null;
|
||||
},
|
||||
voiceBridge?: {
|
||||
registerVoiceSpeakHandler?: (agentId: string, handler: VoiceSpeakHandler) => void;
|
||||
@@ -412,7 +436,7 @@ export class Session {
|
||||
this.voiceLlmDefaultProvider = voice?.voiceLlmDefaultProvider ?? null;
|
||||
this.voiceLlmModel = voice?.voiceLlmModel ?? null;
|
||||
this.voiceLlmAvailability = voice?.voiceLlmAvailability ?? null;
|
||||
this.voiceAgentMcpUrl = voice?.voiceAgentMcpUrl ?? null;
|
||||
this.voiceAgentMcpStdio = voice?.voiceAgentMcpStdio ?? null;
|
||||
this.registerVoiceSpeakHandler = voiceBridge?.registerVoiceSpeakHandler;
|
||||
this.unregisterVoiceSpeakHandler = voiceBridge?.unregisterVoiceSpeakHandler;
|
||||
this.registerVoiceCallerContext = voiceBridge?.registerVoiceCallerContext;
|
||||
@@ -4687,9 +4711,9 @@ export class Session {
|
||||
const cwd = join(this.paseoHome, "voice-agent-workspace");
|
||||
await mkdir(cwd, { recursive: true });
|
||||
|
||||
const mcpUrl = this.voiceAgentMcpUrl;
|
||||
if (!mcpUrl) {
|
||||
throw new Error("Voice MCP URL is not configured");
|
||||
const mcpStdio = this.voiceAgentMcpStdio;
|
||||
if (!mcpStdio) {
|
||||
throw new Error("Voice MCP stdio bridge is not configured");
|
||||
}
|
||||
|
||||
const model = this.getVoiceAgentModel(provider);
|
||||
@@ -4699,10 +4723,12 @@ export class Session {
|
||||
modeId: VOICE_AGENT_DEFAULT_MODE[provider],
|
||||
...(model ? { model } : {}),
|
||||
mcpServers: {
|
||||
paseo: {
|
||||
type: "http",
|
||||
url: `${mcpUrl}?callerAgentId=${encodeURIComponent(voiceAgentId)}`,
|
||||
},
|
||||
paseo: buildVoiceAgentMcpServerConfig({
|
||||
callerAgentId: voiceAgentId,
|
||||
command: mcpStdio.command,
|
||||
baseArgs: mcpStdio.baseArgs,
|
||||
env: mcpStdio.env,
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
|
||||
29
packages/server/src/server/session.voice-mcp-config.test.ts
Normal file
29
packages/server/src/server/session.voice-mcp-config.test.ts
Normal file
@@ -0,0 +1,29 @@
|
||||
import { describe, expect, test } from "vitest";
|
||||
|
||||
import { buildVoiceAgentMcpServerConfig } from "./session.js";
|
||||
|
||||
describe("voice MCP stdio config", () => {
|
||||
test("builds stdio MCP config for voice agent", () => {
|
||||
const config = buildVoiceAgentMcpServerConfig({
|
||||
callerAgentId: "voice-agent-123",
|
||||
command: "/usr/local/bin/paseo",
|
||||
baseArgs: ["__paseo_voice_mcp_bridge", "--socket", "/tmp/paseo-voice.sock"],
|
||||
env: {
|
||||
PASEO_HOME: "/tmp/paseo-home",
|
||||
},
|
||||
});
|
||||
|
||||
expect(config.type).toBe("stdio");
|
||||
expect(config.command).toBe("/usr/local/bin/paseo");
|
||||
expect(config.args).toEqual([
|
||||
"__paseo_voice_mcp_bridge",
|
||||
"--socket",
|
||||
"/tmp/paseo-voice.sock",
|
||||
"--caller-agent-id",
|
||||
"voice-agent-123",
|
||||
]);
|
||||
expect(config.env).toEqual({
|
||||
PASEO_HOME: "/tmp/paseo-home",
|
||||
});
|
||||
});
|
||||
});
|
||||
305
packages/server/src/server/voice-mcp-bridge.ts
Normal file
305
packages/server/src/server/voice-mcp-bridge.ts
Normal file
@@ -0,0 +1,305 @@
|
||||
import net from "node:net";
|
||||
import path from "node:path";
|
||||
import { mkdir, rm } from "node:fs/promises";
|
||||
import type { Logger } from "pino";
|
||||
import pino from "pino";
|
||||
import { InMemoryTransport } from "@modelcontextprotocol/sdk/inMemory.js";
|
||||
import type { JSONRPCMessage } from "@modelcontextprotocol/sdk/types.js";
|
||||
import { StdioServerTransport } from "@modelcontextprotocol/sdk/server/stdio.js";
|
||||
|
||||
|
||||
type BridgeEnvelope =
|
||||
| { type: "init"; callerAgentId: string }
|
||||
| { type: "mcp"; message: JSONRPCMessage }
|
||||
| { type: "ready" }
|
||||
| { type: "error"; message: string };
|
||||
|
||||
function isRecord(value: unknown): value is Record<string, unknown> {
|
||||
return typeof value === "object" && value !== null;
|
||||
}
|
||||
|
||||
function parseEnvelope(raw: string): BridgeEnvelope {
|
||||
const parsed = JSON.parse(raw);
|
||||
if (!isRecord(parsed) || typeof parsed.type !== "string") {
|
||||
throw new Error("Invalid bridge envelope");
|
||||
}
|
||||
if (parsed.type === "init") {
|
||||
const callerAgentId = typeof parsed.callerAgentId === "string" ? parsed.callerAgentId.trim() : "";
|
||||
if (!callerAgentId) {
|
||||
throw new Error("Invalid init payload: callerAgentId is required");
|
||||
}
|
||||
return { type: "init", callerAgentId };
|
||||
}
|
||||
if (parsed.type === "mcp") {
|
||||
if (!("message" in parsed)) {
|
||||
throw new Error("Invalid mcp payload: message is required");
|
||||
}
|
||||
return { type: "mcp", message: parsed.message as JSONRPCMessage };
|
||||
}
|
||||
if (parsed.type === "ready") return { type: "ready" };
|
||||
if (parsed.type === "error") {
|
||||
return { type: "error", message: typeof parsed.message === "string" ? parsed.message : "Unknown error" };
|
||||
}
|
||||
throw new Error(`Unknown envelope type: ${parsed.type}`);
|
||||
}
|
||||
|
||||
function encodeEnvelope(envelope: BridgeEnvelope): string {
|
||||
return `${JSON.stringify(envelope)}\n`;
|
||||
}
|
||||
|
||||
export type VoiceMcpBridgeSocketServer = {
|
||||
socketPath: string;
|
||||
start: () => Promise<void>;
|
||||
stop: () => Promise<void>;
|
||||
};
|
||||
|
||||
export function createVoiceMcpBridgeSocketServer(params: {
|
||||
socketPath: string;
|
||||
logger: Logger;
|
||||
createAgentMcpServerForCaller: (callerAgentId: string) => Promise<{ connect: (transport: InMemoryTransport) => Promise<void>; close?: () => Promise<void> }>;
|
||||
}): VoiceMcpBridgeSocketServer {
|
||||
const logger = params.logger.child({ module: "voice-mcp-bridge" });
|
||||
const sockets = new Set<net.Socket>();
|
||||
|
||||
const server = net.createServer((socket) => {
|
||||
sockets.add(socket);
|
||||
const connectionLogger = logger.child({ component: "connection" });
|
||||
let readBuffer = "";
|
||||
let initialized = false;
|
||||
let clientTransport: InMemoryTransport | null = null;
|
||||
let mcpServer: { close?: () => Promise<void> } | null = null;
|
||||
|
||||
const send = (payload: BridgeEnvelope) => {
|
||||
socket.write(encodeEnvelope(payload));
|
||||
};
|
||||
|
||||
const fail = (message: string) => {
|
||||
send({ type: "error", message });
|
||||
socket.end();
|
||||
};
|
||||
|
||||
const cleanup = async () => {
|
||||
sockets.delete(socket);
|
||||
const closeTasks: Promise<unknown>[] = [];
|
||||
if (clientTransport) {
|
||||
closeTasks.push(clientTransport.close().catch(() => undefined));
|
||||
}
|
||||
if (mcpServer?.close) {
|
||||
closeTasks.push(mcpServer.close().catch(() => undefined));
|
||||
}
|
||||
await Promise.all(closeTasks);
|
||||
};
|
||||
|
||||
socket.on("data", (chunk) => {
|
||||
readBuffer += chunk.toString("utf8");
|
||||
while (true) {
|
||||
const newlineIndex = readBuffer.indexOf("\n");
|
||||
if (newlineIndex < 0) break;
|
||||
const line = readBuffer.slice(0, newlineIndex).trim();
|
||||
readBuffer = readBuffer.slice(newlineIndex + 1);
|
||||
if (!line) continue;
|
||||
|
||||
let message: BridgeEnvelope;
|
||||
try {
|
||||
message = parseEnvelope(line);
|
||||
} catch (error) {
|
||||
fail(error instanceof Error ? error.message : String(error));
|
||||
return;
|
||||
}
|
||||
|
||||
if (message.type === "init") {
|
||||
if (initialized) {
|
||||
fail("Bridge already initialized");
|
||||
return;
|
||||
}
|
||||
initialized = true;
|
||||
void (async () => {
|
||||
try {
|
||||
const [proxyClient, proxyServer] = InMemoryTransport.createLinkedPair();
|
||||
const serverInstance = await params.createAgentMcpServerForCaller(message.callerAgentId);
|
||||
await serverInstance.connect(proxyServer);
|
||||
await proxyClient.start();
|
||||
proxyClient.onmessage = (jsonrpcMessage) => {
|
||||
send({ type: "mcp", message: jsonrpcMessage });
|
||||
};
|
||||
clientTransport = proxyClient;
|
||||
mcpServer = serverInstance;
|
||||
send({ type: "ready" });
|
||||
} catch (error) {
|
||||
connectionLogger.error({ err: error }, "Failed to initialize voice MCP bridge connection");
|
||||
fail(error instanceof Error ? error.message : String(error));
|
||||
}
|
||||
})();
|
||||
continue;
|
||||
}
|
||||
|
||||
if (message.type === "mcp") {
|
||||
if (!clientTransport) {
|
||||
fail("Bridge is not initialized");
|
||||
return;
|
||||
}
|
||||
void clientTransport.send(message.message).catch((error) => {
|
||||
connectionLogger.error({ err: error }, "Failed to forward MCP message");
|
||||
fail(error instanceof Error ? error.message : String(error));
|
||||
});
|
||||
continue;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
socket.on("error", (error) => {
|
||||
connectionLogger.error({ err: error }, "Voice MCP bridge socket error");
|
||||
});
|
||||
socket.on("close", () => {
|
||||
void cleanup();
|
||||
});
|
||||
});
|
||||
|
||||
return {
|
||||
socketPath: params.socketPath,
|
||||
async start() {
|
||||
await mkdir(path.dirname(params.socketPath), { recursive: true });
|
||||
await rm(params.socketPath, { force: true }).catch(() => undefined);
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
server.once("error", reject);
|
||||
server.listen(params.socketPath, () => {
|
||||
server.off("error", reject);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
logger.info({ socketPath: params.socketPath }, "Voice MCP bridge socket server listening");
|
||||
},
|
||||
async stop() {
|
||||
for (const socket of sockets) {
|
||||
socket.destroy();
|
||||
}
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
server.close((error) => {
|
||||
if (error) reject(error);
|
||||
else resolve();
|
||||
});
|
||||
});
|
||||
await rm(params.socketPath, { force: true }).catch(() => undefined);
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function parseBridgeCliArgs(argv: string[]): { socketPath: string; callerAgentId: string } {
|
||||
let socketPath: string | null = null;
|
||||
let callerAgentId: string | null = null;
|
||||
|
||||
for (let index = 0; index < argv.length; index += 1) {
|
||||
const arg = argv[index];
|
||||
if (arg === "--socket") {
|
||||
socketPath = argv[index + 1] ?? null;
|
||||
index += 1;
|
||||
continue;
|
||||
}
|
||||
if (arg === "--caller-agent-id") {
|
||||
callerAgentId = argv[index + 1] ?? null;
|
||||
index += 1;
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
if (!socketPath?.trim()) {
|
||||
throw new Error("Missing required --socket <path>");
|
||||
}
|
||||
if (!callerAgentId?.trim()) {
|
||||
throw new Error("Missing required --caller-agent-id <id>");
|
||||
}
|
||||
|
||||
return {
|
||||
socketPath: socketPath.trim(),
|
||||
callerAgentId: callerAgentId.trim(),
|
||||
};
|
||||
}
|
||||
|
||||
export async function runVoiceMcpBridgeCli(argv: string[], logger?: Logger): Promise<void> {
|
||||
const bridgeLogger = logger ?? pino({ level: "error" });
|
||||
const parsed = parseBridgeCliArgs(argv);
|
||||
|
||||
const socket = net.createConnection(parsed.socketPath);
|
||||
const stdioTransport = new StdioServerTransport(process.stdin, process.stdout);
|
||||
let ready = false;
|
||||
let socketBuffer = "";
|
||||
|
||||
const sendEnvelope = (payload: BridgeEnvelope) => {
|
||||
socket.write(encodeEnvelope(payload));
|
||||
};
|
||||
|
||||
const closeWithError = (message: string): never => {
|
||||
throw new Error(message);
|
||||
};
|
||||
|
||||
socket.on("data", (chunk) => {
|
||||
socketBuffer += chunk.toString("utf8");
|
||||
while (true) {
|
||||
const newlineIndex = socketBuffer.indexOf("\n");
|
||||
if (newlineIndex < 0) break;
|
||||
const line = socketBuffer.slice(0, newlineIndex).trim();
|
||||
socketBuffer = socketBuffer.slice(newlineIndex + 1);
|
||||
if (!line) continue;
|
||||
|
||||
const envelope = parseEnvelope(line);
|
||||
if (envelope.type === "ready") {
|
||||
ready = true;
|
||||
continue;
|
||||
}
|
||||
if (envelope.type === "error") {
|
||||
socket.destroy();
|
||||
closeWithError(`Voice MCP bridge error: ${envelope.message}`);
|
||||
}
|
||||
if (envelope.type === "mcp") {
|
||||
void stdioTransport.send(envelope.message);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
socket.on("error", (error) => {
|
||||
bridgeLogger.error({ err: error }, "Voice MCP bridge socket client error");
|
||||
});
|
||||
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
socket.once("connect", () => resolve());
|
||||
socket.once("error", reject);
|
||||
});
|
||||
|
||||
sendEnvelope({ type: "init", callerAgentId: parsed.callerAgentId });
|
||||
|
||||
// Wait for bridge readiness before forwarding stdio messages.
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
const deadline = setTimeout(() => {
|
||||
reject(new Error("Timed out waiting for voice MCP bridge initialization"));
|
||||
}, 10000);
|
||||
const checkReady = () => {
|
||||
if (ready) {
|
||||
clearTimeout(deadline);
|
||||
resolve();
|
||||
} else {
|
||||
setTimeout(checkReady, 25);
|
||||
}
|
||||
};
|
||||
checkReady();
|
||||
});
|
||||
|
||||
stdioTransport.onmessage = (message) => {
|
||||
sendEnvelope({ type: "mcp", message });
|
||||
};
|
||||
stdioTransport.onerror = (error) => {
|
||||
bridgeLogger.error({ err: error }, "Voice MCP stdio transport error");
|
||||
socket.destroy();
|
||||
};
|
||||
stdioTransport.onclose = () => {
|
||||
socket.end();
|
||||
};
|
||||
|
||||
await stdioTransport.start();
|
||||
|
||||
await new Promise<void>((resolve) => {
|
||||
socket.once("close", () => resolve());
|
||||
process.stdin.once("end", () => resolve());
|
||||
});
|
||||
|
||||
await stdioTransport.close().catch(() => undefined);
|
||||
}
|
||||
@@ -23,6 +23,11 @@ import type { SpeechToTextProvider, TextToSpeechProvider } from "./speech/speech
|
||||
|
||||
export type AgentMcpTransportFactory = () => Promise<Transport>;
|
||||
type VoiceAgentProvider = "claude" | "codex" | "opencode";
|
||||
type VoiceMcpStdioConfig = {
|
||||
command: string;
|
||||
baseArgs: string[];
|
||||
env?: Record<string, string>;
|
||||
};
|
||||
|
||||
type WebSocketServerConfig = {
|
||||
allowedOrigins: Set<string>;
|
||||
@@ -80,7 +85,7 @@ export class VoiceAssistantWebSocketServer {
|
||||
voiceLlmDefaultProvider?: VoiceAgentProvider | null;
|
||||
voiceLlmModel?: string | null;
|
||||
voiceLlmAvailability?: Record<VoiceAgentProvider, boolean> | null;
|
||||
voiceAgentMcpUrl?: string | null;
|
||||
voiceAgentMcpStdio?: VoiceMcpStdioConfig | null;
|
||||
} | null;
|
||||
private readonly voiceSpeakHandlers = new Map<
|
||||
string,
|
||||
@@ -115,7 +120,7 @@ export class VoiceAssistantWebSocketServer {
|
||||
voiceLlmDefaultProvider?: VoiceAgentProvider | null;
|
||||
voiceLlmModel?: string | null;
|
||||
voiceLlmAvailability?: Record<VoiceAgentProvider, boolean> | null;
|
||||
voiceAgentMcpUrl?: string | null;
|
||||
voiceAgentMcpStdio?: VoiceMcpStdioConfig | null;
|
||||
},
|
||||
dictation?: {
|
||||
finalTimeoutMs?: number;
|
||||
|
||||
Reference in New Issue
Block a user