63b203b937
* feat(mcp): resource notifications, subscriptions, and read_resource builtin tool - Add MCP resource subscription lifecycle (subscribe/unsubscribe on connect/disconnect) - Wire mcp.notifications setting with live toggle support - Add debounced followUp injection for resource change notifications - Add global read_resource builtin tool with server resolution by URI/template scheme - Add MCP prompt commands (buildMCPPromptCommands) with array content support - Add server instructions injection into system prompt with attribution - Add mcp.notificationDebounceMs configurable setting Client (client.ts): listResources, listResourceTemplates, readResource with pagination subscribeToResources, unsubscribeFromResources listPrompts, getPrompt, serverSupportsPrompts serverSupportsResources, serverSupportsResourceSubscriptions Manager (manager.ts): Notification dispatch with subscribed-URI guard Concurrent refresh deduplication via pending promise map setNotificationsEnabled with subscribe/unsubscribe toggle Tests: client-resources.test.ts (31 tests) client-prompts.test.ts (20 tests) mcp-read-resource.test.ts (13 tests) * fix(mcp): address PR review - eager prompt init and stale subscription cleanup P1: Make setOnPromptsChanged eagerly fire for servers that already have prompts loaded. The callback is registered after MCP discovery has already loaded prompts and fired the hook, so without this the handler is never called on the common startup path. The fix is in the manager itself (not the caller), eliminating the race condition regardless of when the callback is wired. P2: Unsubscribe removed resource URIs on resource refresh. refreshServerResources was subscribing to the new URI set and overwriting #subscribedResources without unsubscribing URIs that were previously subscribed but no longer present, leaving stale subscriptions active on the server. * fix(mcp): add resources and prompts to /mcp help text and subcommand completions * feat(mcp): add /mcp notifications command Shows per-server notification capabilities with subscription state: - Lists supported notification types (tools/list_changed, resources/list_changed, prompts/list_changed) with check marks - Shows resources/subscribe status with active subscription count - Lists subscribed URIs with green ticks when notifications are enabled - Displays overall enabled/disabled state (mcp.notifications setting) * fix(mcp): address PR review comments on race conditions and stale state - Await subscribe/unsubscribe in refreshServerResources so the refresh promise doesn't resolve before subscriptions are settled, preventing a second refresh from racing and overwriting tracking state (P2 #3) - Guard setNotificationsEnabled subscribe .then() against a disable that happens while the subscribe request is in-flight (P2 #5) - Re-check mcp.notifications setting inside debounce setTimeout callback so toggling off mid-window actually suppresses the follow-up message (P2 #4) - Fire onToolsChanged and onPromptsChanged callbacks in disconnectServer so stale slash commands and tool registrations are cleaned up when a server is removed (P2 #2) --------- Co-authored-by: Miroslav Drbal <miroslav.drbal@gendigital.com>
428 lines
10 KiB
TypeScript
428 lines
10 KiB
TypeScript
/**
|
|
* MCP Client.
|
|
*
|
|
* Handles connection initialization, tool listing, and tool calling.
|
|
*/
|
|
import { 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",
|
|
};
|
|
|
|
/**
|
|
* 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 },
|
|
): 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>,
|
|
options,
|
|
);
|
|
|
|
if (options?.signal?.aborted) {
|
|
throw options.signal.reason instanceof Error ? options.signal.reason : new Error("Aborted");
|
|
}
|
|
|
|
// 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 },
|
|
): Promise<MCPServerConnection> {
|
|
const timeoutMs = config.timeout ?? CONNECTION_TIMEOUT_MS;
|
|
|
|
const connect = async (): Promise<MCPServerConnection> => {
|
|
const transport = await createTransport(config);
|
|
if (options?.onNotification) {
|
|
transport.onNotification = options.onNotification;
|
|
}
|
|
|
|
try {
|
|
const initResult = await initializeConnection(transport, options);
|
|
|
|
// Start SSE background listener for HTTP transports
|
|
if ("startSSEListener" in transport && typeof transport.startSSEListener === "function") {
|
|
void (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;
|
|
}
|
|
};
|
|
|
|
return withTimeout(
|
|
connect(),
|
|
timeoutMs,
|
|
`Connection to MCP server "${name}" timed out after ${timeoutMs}ms`,
|
|
options?.signal,
|
|
);
|
|
}
|
|
|
|
/**
|
|
* 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;
|
|
}
|