Files
oh-my-pi/packages/coding-agent/src/mcp/client.ts
T
Miroslav Drbal [ApoC] 63b203b937 feat(mcp): resource notifications, subscriptions, and read_resource builtin tool (#254)
* 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>
2026-03-03 03:26:07 +01:00

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;
}