8b028224e3
* feat: auto-reconnect MCP servers on connection loss
When an HTTP SSE stream drops (server restart, network interruption),
the transport fires onClose, and the manager proactively reconnects
with retry backoff (500ms, 1s, 2s, 4s). Tools are kept in the registry
during reconnection so they remain selected and available to the agent.
If proactive reconnection fails, stale tools stay registered. When the
agent calls one, the tool bridge detects the retriable connection error
(ECONNREFUSED, ECONNRESET, stale session 404/502/503, etc.), triggers
reconnectServer on the manager, and retries the call once on the fresh
connection. Concurrent reconnect attempts for the same server are deduped.
connectToServer now always installs a default onRequest handler for ping
and roots/list (using getProjectDir()), so all connections -- including
short-lived test/probe ones -- properly respond to server-initiated
requests during initialization.
Post-connection setup (resources, prompts, subscriptions) is extracted
into a shared #loadServerResourcesAndPrompts method used by both initial
connection and reconnection paths.
Add /mcp reconnect <name> command for manual recovery after extended
outages where both proactive and reactive reconnection have failed.
* docs: add changelog entry for MCP auto-reconnect
* fix: address P1 review findings in MCP reconnection
- Save server configs before connection attempt so deferred tools can
reconnect even when the initial connection timed out (P1-1)
- Make waitForConnection() and getConnectionStatus() aware of in-flight
reconnections so callers wait instead of failing immediately (P1-2)
- Add epoch counter incremented on disconnectAll() and checked in
connectAndWireServer() to invalidate stale reconnect attempts that
outlive a manager reset/reload (P1-3)
- Skip servers with pending reconnections in connectServers() to prevent
parallel connection attempts for the same server
* fix: deferred tool reconnect and non-blocking transport teardown
- DeferredMCPTool.execute now reconnects when getConnection() fails
("MCP server not connected"), not only on network errors from
callTool. Servers that missed the startup window can now be woken
by the first tool call against their cached tools. (P1-4)
- #doReconnect fire-and-forgets the old transport close instead of
awaiting it. HttpTransport.close() sends a DELETE with 30s timeout;
blocking here delayed the first reconnect attempt by that amount
on every server restart. (P1-5)
* fix: abort-aware reconnect waits and preserve tool selection on reconnect
- Wrap all reconnect() awaits with withAbort(signal) so user
cancellation (Esc) interrupts the reconnect backoff loop instead
of blocking for up to 7.5s. Applies to MCPTool (1 site) and
DeferredMCPTool (2 sites). (P2-1)
- Remove activateDiscoveredMCPTools call from /mcp reconnect handler.
refreshMCPTools already preserves the user's prior MCP tool
selection; the extra activation was silently opting into all
server tools including ones the user had not enabled. (P2-2)
* fix: rebind MCPTool connection after reconnect, add stdio retriable error
- MCPTool.connection is now mutable; after a successful reconnect retry,
this.connection is rebound to the fresh connection so subsequent calls
on the same instance (e.g. batched tool calls) use it instead of
triggering another reconnect cycle. (P2-3)
- Add "Transport closed" to RETRIABLE_PATTERNS. StdioTransport rejects
pending requests with this message when the subprocess dies, which
should trigger the reconnect path just like HTTP transport errors. (P2-4)
---------
Co-authored-by: Miroslav Drbal <miroslav.drbal@gendigital.com>
Co-authored-by: Can Bölük <can1357@users.noreply.github.com>
483 lines
12 KiB
TypeScript
483 lines
12 KiB
TypeScript
/**
|
|
* MCP Client.
|
|
*
|
|
* Handles connection initialization, tool listing, and tool calling.
|
|
*/
|
|
import * as path from "node:path";
|
|
import * as url from "node:url";
|
|
import { getProjectDir, logger, withTimeout } from "@oh-my-pi/pi-utils";
|
|
import { createHttpTransport } from "./transports/http";
|
|
import { createStdioTransport } from "./transports/stdio";
|
|
import type {
|
|
MCPGetPromptParams,
|
|
MCPGetPromptResult,
|
|
MCPHttpServerConfig,
|
|
MCPInitializeParams,
|
|
MCPInitializeResult,
|
|
MCPPrompt,
|
|
MCPPromptsListResult,
|
|
MCPRequestOptions,
|
|
MCPResource,
|
|
MCPResourceReadParams,
|
|
MCPResourceReadResult,
|
|
MCPResourceSubscribeParams,
|
|
MCPResourcesListResult,
|
|
MCPResourceTemplate,
|
|
MCPResourceTemplatesListResult,
|
|
MCPServerCapabilities,
|
|
MCPServerConfig,
|
|
MCPServerConnection,
|
|
MCPSseServerConfig,
|
|
MCPStdioServerConfig,
|
|
MCPToolCallParams,
|
|
MCPToolCallResult,
|
|
MCPToolDefinition,
|
|
MCPToolsListResult,
|
|
MCPTransport,
|
|
} from "./types";
|
|
|
|
/** MCP protocol version we support */
|
|
const PROTOCOL_VERSION = "2025-03-26";
|
|
|
|
/** Default connection timeout in ms */
|
|
const CONNECTION_TIMEOUT_MS = 30_000;
|
|
|
|
/** Client info sent during initialization */
|
|
const CLIENT_INFO = {
|
|
name: "omp-coding-agent",
|
|
version: "1.0.0",
|
|
};
|
|
|
|
/**
|
|
* Default handler for standard MCP server-to-client requests.
|
|
* Handles `ping` and `roots/list`; rejects unknown methods with -32601.
|
|
* Reads getProjectDir() at call time so the root stays stable even if
|
|
* the process cwd changes during tool execution.
|
|
*/
|
|
async function defaultRequestHandler(method: string, _params: unknown): Promise<unknown> {
|
|
switch (method) {
|
|
case "ping":
|
|
return {};
|
|
case "roots/list": {
|
|
const cwd = getProjectDir();
|
|
return {
|
|
roots: [{ uri: url.pathToFileURL(cwd).href, name: path.basename(cwd) }],
|
|
};
|
|
}
|
|
default:
|
|
throw Object.assign(new Error(`Unsupported server request: ${method}`), { code: -32601 });
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Create a transport for the given server config.
|
|
*/
|
|
async function createTransport(config: MCPServerConfig): Promise<MCPTransport> {
|
|
const serverType = config.type ?? "stdio";
|
|
|
|
switch (serverType) {
|
|
case "stdio":
|
|
return createStdioTransport(config as MCPStdioServerConfig);
|
|
case "http":
|
|
case "sse":
|
|
return createHttpTransport(config as MCPHttpServerConfig | MCPSseServerConfig);
|
|
default:
|
|
throw new Error(`Unknown server type: ${serverType}`);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Initialize connection with MCP server.
|
|
*/
|
|
async function initializeConnection(
|
|
transport: MCPTransport,
|
|
options?: {
|
|
signal?: AbortSignal;
|
|
/** Called after the initialize response (which sets the session ID) but before notifications/initialized. */
|
|
onInitialized?: () => void | Promise<void>;
|
|
},
|
|
): Promise<MCPInitializeResult> {
|
|
const params: MCPInitializeParams = {
|
|
protocolVersion: PROTOCOL_VERSION,
|
|
capabilities: {
|
|
roots: { listChanged: false },
|
|
},
|
|
clientInfo: CLIENT_INFO,
|
|
};
|
|
|
|
const result = await transport.request<MCPInitializeResult>(
|
|
"initialize",
|
|
params as unknown as Record<string, unknown>,
|
|
{ signal: options?.signal },
|
|
);
|
|
|
|
if (options?.signal?.aborted) {
|
|
throw options.signal.reason instanceof Error ? options.signal.reason : new Error("Aborted");
|
|
}
|
|
|
|
// Hook point: the transport now has the session ID from the initialize response.
|
|
// For HTTP, this is the moment to open the SSE stream so server-to-client requests
|
|
// triggered by notifications/initialized (e.g. roots/list) can be delivered.
|
|
await options?.onInitialized?.();
|
|
|
|
// Send initialized notification
|
|
await transport.notify("notifications/initialized");
|
|
|
|
return result;
|
|
}
|
|
|
|
/**
|
|
* Connect to an MCP server.
|
|
* Has a 30 second timeout to prevent blocking startup.
|
|
*/
|
|
export async function connectToServer(
|
|
name: string,
|
|
config: MCPServerConfig,
|
|
options?: {
|
|
signal?: AbortSignal;
|
|
onNotification?: (method: string, params: unknown) => void;
|
|
onRequest?: (method: string, params: unknown) => Promise<unknown>;
|
|
},
|
|
): Promise<MCPServerConnection> {
|
|
const timeoutMs = config.timeout ?? CONNECTION_TIMEOUT_MS;
|
|
let transport: MCPTransport | undefined;
|
|
|
|
const connect = async (): Promise<MCPServerConnection> => {
|
|
transport = await createTransport(config);
|
|
if (options?.onNotification) {
|
|
transport.onNotification = options.onNotification;
|
|
}
|
|
|
|
// Always handle standard MCP server-to-client requests (ping, roots/list).
|
|
// The initialize request declares roots capability, so we must respond to
|
|
// roots/list — even for short-lived test connections.
|
|
transport.onRequest = options?.onRequest ?? defaultRequestHandler;
|
|
|
|
try {
|
|
const initResult = await initializeConnection(transport, {
|
|
signal: options?.signal,
|
|
async onInitialized() {
|
|
// Open the SSE stream before sending initialized, so server-to-client
|
|
// requests triggered by on_initialized (e.g. roots/list) are delivered.
|
|
if ("startSSEListener" in transport! && typeof transport!.startSSEListener === "function") {
|
|
await (transport as { startSSEListener(): Promise<void> }).startSSEListener();
|
|
}
|
|
},
|
|
});
|
|
|
|
return {
|
|
name,
|
|
config,
|
|
transport,
|
|
serverInfo: initResult.serverInfo,
|
|
capabilities: initResult.capabilities,
|
|
instructions: initResult.instructions,
|
|
};
|
|
} catch (error) {
|
|
await transport.close();
|
|
throw error;
|
|
}
|
|
};
|
|
|
|
try {
|
|
return await withTimeout(
|
|
connect(),
|
|
timeoutMs,
|
|
`Connection to MCP server "${name}" timed out after ${timeoutMs}ms`,
|
|
options?.signal,
|
|
);
|
|
} catch (error) {
|
|
// If withTimeout rejected (timeout/abort) while connect() was still pending,
|
|
// the transport may be alive with an open SSE listener. Close it.
|
|
if (transport) {
|
|
void transport.close().catch(() => {});
|
|
}
|
|
throw error;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* List tools from a connected server.
|
|
*/
|
|
export async function listTools(
|
|
connection: MCPServerConnection,
|
|
options?: { signal?: AbortSignal },
|
|
): Promise<MCPToolDefinition[]> {
|
|
// Check if server supports tools
|
|
if (!connection.capabilities.tools) {
|
|
return [];
|
|
}
|
|
|
|
// Return cached tools if available
|
|
if (connection.tools) {
|
|
return connection.tools;
|
|
}
|
|
|
|
const allTools: MCPToolDefinition[] = [];
|
|
let cursor: string | undefined;
|
|
|
|
do {
|
|
const params: Record<string, unknown> = {};
|
|
if (cursor) {
|
|
params.cursor = cursor;
|
|
}
|
|
|
|
const result = await connection.transport.request<MCPToolsListResult>("tools/list", params, options);
|
|
allTools.push(...result.tools);
|
|
cursor = result.nextCursor;
|
|
} while (cursor);
|
|
|
|
// Cache tools
|
|
connection.tools = allTools;
|
|
|
|
return allTools;
|
|
}
|
|
|
|
/**
|
|
* Call a tool on a connected server.
|
|
*/
|
|
export async function callTool(
|
|
connection: MCPServerConnection,
|
|
toolName: string,
|
|
args: Record<string, unknown> = {},
|
|
options?: MCPRequestOptions,
|
|
): Promise<MCPToolCallResult> {
|
|
const params: MCPToolCallParams = {
|
|
name: toolName,
|
|
arguments: args,
|
|
};
|
|
|
|
return connection.transport.request<MCPToolCallResult>(
|
|
"tools/call",
|
|
params as unknown as Record<string, unknown>,
|
|
options,
|
|
);
|
|
}
|
|
|
|
/**
|
|
* Disconnect from a server.
|
|
*/
|
|
export async function disconnectServer(connection: MCPServerConnection): Promise<void> {
|
|
await connection.transport.close();
|
|
}
|
|
|
|
/**
|
|
* Check if a server supports tools.
|
|
*/
|
|
export function serverSupportsTools(capabilities: MCPServerCapabilities): boolean {
|
|
return capabilities.tools !== undefined;
|
|
}
|
|
|
|
/**
|
|
* List resources from a connected server.
|
|
*/
|
|
export async function listResources(
|
|
connection: MCPServerConnection,
|
|
options?: { signal?: AbortSignal },
|
|
): Promise<MCPResource[]> {
|
|
if (!connection.capabilities.resources) {
|
|
return [];
|
|
}
|
|
|
|
if (connection.resources) {
|
|
return connection.resources;
|
|
}
|
|
|
|
const allResources: MCPResource[] = [];
|
|
let cursor: string | undefined;
|
|
|
|
do {
|
|
const params: Record<string, unknown> = {};
|
|
if (cursor) {
|
|
params.cursor = cursor;
|
|
}
|
|
|
|
const result = await connection.transport.request<MCPResourcesListResult>("resources/list", params, options);
|
|
allResources.push(...result.resources);
|
|
cursor = result.nextCursor;
|
|
} while (cursor);
|
|
|
|
connection.resources = allResources;
|
|
return allResources;
|
|
}
|
|
|
|
/**
|
|
* List resource templates from a connected server.
|
|
*/
|
|
export async function listResourceTemplates(
|
|
connection: MCPServerConnection,
|
|
options?: { signal?: AbortSignal },
|
|
): Promise<MCPResourceTemplate[]> {
|
|
if (!connection.capabilities.resources) {
|
|
return [];
|
|
}
|
|
|
|
if (connection.resourceTemplates) {
|
|
return connection.resourceTemplates;
|
|
}
|
|
|
|
const allTemplates: MCPResourceTemplate[] = [];
|
|
let cursor: string | undefined;
|
|
|
|
do {
|
|
const params: Record<string, unknown> = {};
|
|
if (cursor) {
|
|
params.cursor = cursor;
|
|
}
|
|
|
|
const result = await connection.transport.request<MCPResourceTemplatesListResult>(
|
|
"resources/templates/list",
|
|
params,
|
|
options,
|
|
);
|
|
allTemplates.push(...result.resourceTemplates);
|
|
cursor = result.nextCursor;
|
|
} while (cursor);
|
|
|
|
connection.resourceTemplates = allTemplates;
|
|
return allTemplates;
|
|
}
|
|
|
|
/**
|
|
* Read a resource from a connected server.
|
|
*/
|
|
export async function readResource(
|
|
connection: MCPServerConnection,
|
|
uri: string,
|
|
options?: MCPRequestOptions,
|
|
): Promise<MCPResourceReadResult> {
|
|
const params: MCPResourceReadParams = { uri };
|
|
return connection.transport.request<MCPResourceReadResult>(
|
|
"resources/read",
|
|
params as unknown as Record<string, unknown>,
|
|
options,
|
|
);
|
|
}
|
|
|
|
/**
|
|
* Subscribe to resource update notifications.
|
|
*/
|
|
export async function subscribeToResources(
|
|
connection: MCPServerConnection,
|
|
uris: string[],
|
|
options?: MCPRequestOptions,
|
|
): Promise<void> {
|
|
if (uris.length === 0 || !connection.capabilities.resources?.subscribe) return;
|
|
const results = await Promise.allSettled(
|
|
uris.map(uri => {
|
|
const params: MCPResourceSubscribeParams = { uri };
|
|
return connection.transport.request(
|
|
"resources/subscribe",
|
|
params as unknown as Record<string, unknown>,
|
|
options,
|
|
);
|
|
}),
|
|
);
|
|
for (const result of results) {
|
|
if (result.status === "rejected") {
|
|
logger.warn("Failed to subscribe to MCP resource", { error: result.reason });
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Unsubscribe from resource update notifications.
|
|
*/
|
|
export async function unsubscribeFromResources(
|
|
connection: MCPServerConnection,
|
|
uris: string[],
|
|
options?: MCPRequestOptions,
|
|
): Promise<void> {
|
|
if (uris.length === 0 || !connection.capabilities.resources?.subscribe) return;
|
|
const results = await Promise.allSettled(
|
|
uris.map(uri => {
|
|
const params: MCPResourceSubscribeParams = { uri };
|
|
return connection.transport.request(
|
|
"resources/unsubscribe",
|
|
params as unknown as Record<string, unknown>,
|
|
options,
|
|
);
|
|
}),
|
|
);
|
|
for (const result of results) {
|
|
if (result.status === "rejected") {
|
|
logger.warn("Failed to unsubscribe from MCP resource", { error: result.reason });
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Check if a server supports resource subscriptions.
|
|
*/
|
|
export function serverSupportsResourceSubscriptions(capabilities: MCPServerCapabilities): boolean {
|
|
return capabilities.resources?.subscribe === true;
|
|
}
|
|
|
|
/**
|
|
* Check if a server supports resources.
|
|
*/
|
|
export function serverSupportsResources(capabilities: MCPServerCapabilities): boolean {
|
|
return capabilities.resources !== undefined;
|
|
}
|
|
|
|
/**
|
|
* List prompts from a connected server.
|
|
*/
|
|
export async function listPrompts(
|
|
connection: MCPServerConnection,
|
|
options?: { signal?: AbortSignal },
|
|
): Promise<MCPPrompt[]> {
|
|
if (!connection.capabilities.prompts) {
|
|
return [];
|
|
}
|
|
|
|
if (connection.prompts) {
|
|
return connection.prompts;
|
|
}
|
|
|
|
const allPrompts: MCPPrompt[] = [];
|
|
let cursor: string | undefined;
|
|
|
|
do {
|
|
const params: Record<string, unknown> = {};
|
|
if (cursor) {
|
|
params.cursor = cursor;
|
|
}
|
|
|
|
const result = await connection.transport.request<MCPPromptsListResult>("prompts/list", params, options);
|
|
allPrompts.push(...result.prompts);
|
|
cursor = result.nextCursor;
|
|
} while (cursor);
|
|
|
|
connection.prompts = allPrompts;
|
|
return allPrompts;
|
|
}
|
|
|
|
/**
|
|
* Get a specific prompt from a connected server.
|
|
*/
|
|
export async function getPrompt(
|
|
connection: MCPServerConnection,
|
|
name: string,
|
|
args?: Record<string, string>,
|
|
options?: MCPRequestOptions,
|
|
): Promise<MCPGetPromptResult> {
|
|
const params: MCPGetPromptParams = { name };
|
|
if (args && Object.keys(args).length > 0) {
|
|
params.arguments = args;
|
|
}
|
|
|
|
return connection.transport.request<MCPGetPromptResult>(
|
|
"prompts/get",
|
|
params as unknown as Record<string, unknown>,
|
|
options,
|
|
);
|
|
}
|
|
|
|
/**
|
|
* Check if a server supports prompts.
|
|
*/
|
|
export function serverSupportsPrompts(capabilities: MCPServerCapabilities): boolean {
|
|
return capabilities.prompts !== undefined;
|
|
}
|