feat(coding-agent): added MCP server auto-reconnect with SSE monitoring
- Added auto-reconnect capability for MCP servers with SSE stream monitoring and exponential retry backoff. - Added tool-level reconnect handling for retriable connection errors (ECONNREFUSED, ECONNRESET, 404/502/503). - Added `/mcp reconnect <name>` command for manual MCP server recovery. - Improved reconnect robustness by aborting retries when MCP configuration changes via epoch checking. - Extended transport reconnect handling to all transport types (stdio, HTTP/SSE) with unified onClose logic. - Added comprehensive test coverage for MCPManager reconnect behavior and tool-level abort propagation.
This commit is contained in:
@@ -1,14 +1,6 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
### Changed
|
||||
|
||||
- Updated explore agent thinking level from off to med for improved reasoning
|
||||
- Simplified explore agent output schema: consolidated file references into single `ref` field with optional line ranges instead of separate `path`, `line_start`, `line_end` fields
|
||||
- Removed `code` section from explore agent output (critical code excerpts no longer extracted)
|
||||
- Removed `dependencies` section from explore agent output
|
||||
- Removed `risks` section from explore agent output
|
||||
- Removed `start_here` section from explore agent output
|
||||
|
||||
### Added
|
||||
|
||||
@@ -16,8 +8,20 @@
|
||||
- Tool-level reconnect: retriable connection errors (ECONNREFUSED, ECONNRESET, stale session 404/502/503) trigger automatic reconnection and single retry
|
||||
- `/mcp reconnect <name>` command for manual server recovery after extended outages
|
||||
|
||||
### Changed
|
||||
|
||||
- Extended transport reconnect handling to all transport types (not just HTTP/SSE), ensuring stdio and other transports trigger automatic reconnection on connection loss
|
||||
- Improved reconnect robustness by aborting retry attempts when MCP server configuration changes during reconnection sequence
|
||||
- Updated explore agent thinking level from off to med for improved reasoning
|
||||
- Simplified explore agent output schema: consolidated file references into single `ref` field with optional line ranges instead of separate `path`, `line_start`, `line_end` fields
|
||||
- Removed `code` section from explore agent output (critical code excerpts no longer extracted)
|
||||
- Removed `dependencies` section from explore agent output
|
||||
- Removed `risks` section from explore agent output
|
||||
- Removed `start_here` section from explore agent output
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed reconnect retry loop continuing after configuration changes by checking epoch before each reconnection attempt
|
||||
- `roots/list` timeout on MCP server initialization: `connectToServer` now always installs a default handler for `ping` and `roots/list`
|
||||
- Fixed resumed GitHub Copilot conversations that could fail with `401 input item does not belong to this connection` on the first follow-up after process restart ([#488](https://github.com/can1357/oh-my-pi/issues/488))
|
||||
- Fixed STT Alt+H mic cursor rendering to measure the actual microphone glyph width, preventing one-column TUI overflow crashes when the active symbol preset uses a wide icon ([#484](https://github.com/can1357/oh-my-pi/issues/484))
|
||||
|
||||
@@ -715,6 +715,14 @@ export class MCPManager {
|
||||
// Retry with backoff — the server may still be starting up.
|
||||
const delays = [500, 1000, 2000, 4000];
|
||||
for (let attempt = 0; attempt <= delays.length; attempt++) {
|
||||
if (this.#epoch !== reconnectEpoch) {
|
||||
logger.debug("MCP reconnect aborted before attempt after configuration changed", {
|
||||
path: `mcp:${name}`,
|
||||
storedEpoch: reconnectEpoch,
|
||||
currentEpoch: this.#epoch,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
const connection = await this.#connectAndWireServer(name, config, source, reconnectEpoch);
|
||||
logger.debug("MCP reconnected", { path: `mcp:${name}`, tools: connection.tools?.length ?? 0 });
|
||||
@@ -778,23 +786,20 @@ export class MCPManager {
|
||||
|
||||
this.#connections.set(name, connection);
|
||||
|
||||
// Wire auth refresh and SSE reconnect for HTTP transports
|
||||
if (connection.transport instanceof HttpTransport) {
|
||||
if (config.auth?.type === "oauth") {
|
||||
connection.transport.onAuthError = async () => {
|
||||
const refreshed = await this.#resolveAuthConfig(config, true);
|
||||
if (refreshed.type === "http" || refreshed.type === "sse") {
|
||||
return refreshed.headers ?? null;
|
||||
}
|
||||
return null;
|
||||
};
|
||||
}
|
||||
connection.transport.onClose = () => {
|
||||
logger.debug("MCP SSE stream lost, triggering reconnect", { path: `mcp:${name}` });
|
||||
void this.reconnectServer(name);
|
||||
// Wire auth refresh for HTTP transports, and reconnect for any transport.
|
||||
if (connection.transport instanceof HttpTransport && config.auth?.type === "oauth") {
|
||||
connection.transport.onAuthError = async () => {
|
||||
const refreshed = await this.#resolveAuthConfig(config, true);
|
||||
if (refreshed.type === "http" || refreshed.type === "sse") {
|
||||
return refreshed.headers ?? null;
|
||||
}
|
||||
return null;
|
||||
};
|
||||
}
|
||||
|
||||
connection.transport.onClose = () => {
|
||||
logger.debug("MCP transport lost, triggering reconnect", { path: `mcp:${name}` });
|
||||
void this.reconnectServer(name);
|
||||
};
|
||||
try {
|
||||
const serverTools = await listTools(connection);
|
||||
const reconnect = () => this.reconnectServer(name);
|
||||
@@ -805,10 +810,8 @@ export class MCPManager {
|
||||
void this.#loadServerResourcesAndPrompts(name, connection);
|
||||
return connection;
|
||||
} catch (error) {
|
||||
// Clean up the connection to avoid zombie SSE streams
|
||||
if (connection.transport instanceof HttpTransport) {
|
||||
connection.transport.onClose = undefined;
|
||||
}
|
||||
// Clean up the connection to avoid zombie transports
|
||||
connection.transport.onClose = undefined;
|
||||
await connection.transport.close().catch(() => {});
|
||||
this.#connections.delete(name);
|
||||
throw error;
|
||||
|
||||
@@ -4,8 +4,8 @@
|
||||
* Converts MCP tool definitions to CustomTool format for the agent.
|
||||
*/
|
||||
import type { AgentToolUpdateCallback } from "@oh-my-pi/pi-agent-core";
|
||||
import { untilAborted } from "@oh-my-pi/pi-utils";
|
||||
import { sanitizeSchemaForMCP } from "@oh-my-pi/pi-ai/utils/schema";
|
||||
import { untilAborted } from "@oh-my-pi/pi-utils";
|
||||
import type { TSchema } from "@sinclair/typebox";
|
||||
import type { SourceMeta } from "../capability/types";
|
||||
import type {
|
||||
@@ -20,7 +20,6 @@ import { callTool } from "./client";
|
||||
import { renderMCPCall, renderMCPResult } from "./render";
|
||||
import type { MCPContent, MCPServerConnection, MCPToolCallParams, MCPToolCallResult, MCPToolDefinition } from "./types";
|
||||
|
||||
|
||||
/** Reconnect callback: tears down stale connection, returns new one or null. */
|
||||
export type MCPReconnect = () => Promise<MCPServerConnection | null>;
|
||||
|
||||
@@ -145,6 +144,14 @@ function rethrowIfAborted(error: unknown, signal?: AbortSignal): void {
|
||||
if (signal?.aborted) throw new ToolAbortError();
|
||||
}
|
||||
|
||||
async function reconnectWithAbort(reconnect: MCPReconnect, signal?: AbortSignal): Promise<MCPServerConnection | null> {
|
||||
try {
|
||||
return await untilAborted(signal, reconnect);
|
||||
} catch (error) {
|
||||
rethrowIfAborted(error, signal);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a unique tool name for an MCP tool.
|
||||
@@ -255,7 +262,7 @@ export class MCPTool implements CustomTool<TSchema, MCPToolDetails> {
|
||||
} catch (error) {
|
||||
rethrowIfAborted(error, signal);
|
||||
if (this.reconnect && isRetriableConnectionError(error)) {
|
||||
const newConn = await untilAborted(signal, this.reconnect).catch(() => null);
|
||||
const newConn = await reconnectWithAbort(this.reconnect, signal);
|
||||
if (newConn) {
|
||||
// Rebind so subsequent calls on this instance use the fresh connection
|
||||
this.connection = newConn;
|
||||
@@ -359,7 +366,7 @@ export class DeferredMCPTool implements CustomTool<TSchema, MCPToolDetails> {
|
||||
} catch (callError) {
|
||||
rethrowIfAborted(callError, signal);
|
||||
if (this.reconnect && isRetriableConnectionError(callError)) {
|
||||
const newConn = await untilAborted(signal, this.reconnect).catch(() => null);
|
||||
const newConn = await reconnectWithAbort(this.reconnect, signal);
|
||||
if (newConn) {
|
||||
const retryProvider = newConn._source?.provider ?? provider;
|
||||
const retryProviderName = newConn._source?.providerName ?? providerName;
|
||||
@@ -386,7 +393,7 @@ export class DeferredMCPTool implements CustomTool<TSchema, MCPToolDetails> {
|
||||
// error ("MCP server not connected") isn't a network error from callTool.
|
||||
rethrowIfAborted(connError, signal);
|
||||
if (this.reconnect) {
|
||||
const newConn = await untilAborted(signal, this.reconnect).catch(() => null);
|
||||
const newConn = await reconnectWithAbort(this.reconnect, signal);
|
||||
if (newConn) {
|
||||
try {
|
||||
const result = await callTool(newConn, this.tool.name, args, { signal });
|
||||
|
||||
@@ -913,10 +913,7 @@ describe("stripNewLinePrefixes", () => {
|
||||
|
||||
it("strips plus hashline prefixes in mixed +/ - change style", () => {
|
||||
const lines = ["-**Storage location TBD:**", "+MW:**Storage location TBD:**"];
|
||||
expect(stripNewLinePrefixes(lines)).toEqual([
|
||||
"-**Storage location TBD:**",
|
||||
"**Storage location TBD:**",
|
||||
]);
|
||||
expect(stripNewLinePrefixes(lines)).toEqual(["-**Storage location TBD:**", "**Storage location TBD:**"]);
|
||||
});
|
||||
|
||||
it("does NOT strip hashline prefixes when any non-empty line is plain content", () => {
|
||||
|
||||
@@ -0,0 +1,121 @@
|
||||
import { beforeEach, describe, expect, it, mock, vi } from "bun:test";
|
||||
import type { SourceMeta } from "../src/capability/types";
|
||||
import type { MCPServerConfig, MCPServerConnection, MCPToolDefinition, MCPTransport } from "../src/mcp/types";
|
||||
|
||||
const connectToServerMock = vi.fn();
|
||||
const disconnectServerMock = vi.fn();
|
||||
const listToolsMock = vi.fn();
|
||||
|
||||
mock.module("../src/mcp/client", () => ({
|
||||
connectToServer: connectToServerMock,
|
||||
disconnectServer: disconnectServerMock,
|
||||
getPrompt: vi.fn(),
|
||||
listPrompts: vi.fn(),
|
||||
listResources: vi.fn(),
|
||||
listResourceTemplates: vi.fn(),
|
||||
listTools: listToolsMock,
|
||||
readResource: vi.fn(),
|
||||
serverSupportsPrompts: vi.fn(() => false),
|
||||
serverSupportsResources: vi.fn(() => false),
|
||||
subscribeToResources: vi.fn(),
|
||||
unsubscribeFromResources: vi.fn(),
|
||||
}));
|
||||
|
||||
import { MCPManager } from "../src/mcp/manager";
|
||||
|
||||
function createTransport(): MCPTransport {
|
||||
return {
|
||||
connected: true,
|
||||
async request() {
|
||||
throw new Error("request not implemented");
|
||||
},
|
||||
async notify() {},
|
||||
async close() {},
|
||||
};
|
||||
}
|
||||
|
||||
function createConnection(name: string, config: MCPServerConfig, transport: MCPTransport): MCPServerConnection {
|
||||
return {
|
||||
name,
|
||||
config,
|
||||
transport,
|
||||
serverInfo: { name: "mock", version: "1.0.0" },
|
||||
capabilities: { tools: {} },
|
||||
};
|
||||
}
|
||||
|
||||
function createSource(path: string): SourceMeta {
|
||||
return {
|
||||
provider: "mcp-json",
|
||||
providerName: "MCP JSON",
|
||||
path,
|
||||
level: "project",
|
||||
};
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
connectToServerMock.mockReset();
|
||||
disconnectServerMock.mockReset();
|
||||
listToolsMock.mockReset();
|
||||
|
||||
disconnectServerMock.mockImplementation(async (connection: MCPServerConnection) => {
|
||||
await connection.transport.close();
|
||||
});
|
||||
listToolsMock.mockResolvedValue([] satisfies MCPToolDefinition[]);
|
||||
});
|
||||
|
||||
describe("MCPManager reconnect behavior", () => {
|
||||
it("wires stdio transport onClose to reconnectServer", async () => {
|
||||
const manager = new MCPManager("/tmp");
|
||||
const serverName = "stdio-server";
|
||||
const config: MCPServerConfig = { type: "stdio", command: "mock-server" };
|
||||
const transport = createTransport();
|
||||
connectToServerMock.mockResolvedValueOnce(createConnection(serverName, config, transport));
|
||||
|
||||
await manager.connectServers({ [serverName]: config }, { [serverName]: createSource("/tmp/.mcp.json") });
|
||||
|
||||
const connection = manager.getConnection(serverName);
|
||||
expect(connection).toBeDefined();
|
||||
expect(typeof connection?.transport.onClose).toBe("function");
|
||||
|
||||
const reconnectSpy = vi.spyOn(manager, "reconnectServer").mockResolvedValue(null);
|
||||
connection?.transport.onClose?.();
|
||||
expect(reconnectSpy).toHaveBeenCalledWith(serverName);
|
||||
});
|
||||
|
||||
it("stops reconnect retries after disconnectAll increments epoch", async () => {
|
||||
const manager = new MCPManager("/tmp");
|
||||
const serverName = "epoch-server";
|
||||
const config: MCPServerConfig = { type: "stdio", command: "mock-server" };
|
||||
const firstReconnectAttempt = Promise.withResolvers<void>();
|
||||
let connectCalls = 0;
|
||||
|
||||
connectToServerMock.mockImplementation(async () => {
|
||||
connectCalls += 1;
|
||||
if (connectCalls === 1) {
|
||||
return createConnection(serverName, config, createTransport());
|
||||
}
|
||||
if (connectCalls === 2) {
|
||||
firstReconnectAttempt.resolve();
|
||||
throw new Error("ECONNREFUSED");
|
||||
}
|
||||
return createConnection(serverName, config, createTransport());
|
||||
});
|
||||
|
||||
await manager.connectServers({ [serverName]: config }, { [serverName]: createSource("/tmp/.mcp.json") });
|
||||
|
||||
const sleepGate = Promise.withResolvers<void>();
|
||||
const sleepSpy = vi.spyOn(Bun, "sleep").mockImplementation(async () => sleepGate.promise);
|
||||
|
||||
const reconnectPromise = manager.reconnectServer(serverName);
|
||||
await firstReconnectAttempt.promise;
|
||||
await manager.disconnectAll();
|
||||
sleepGate.resolve();
|
||||
|
||||
const result = await reconnectPromise;
|
||||
expect(result).toBeNull();
|
||||
expect(connectCalls).toBe(2);
|
||||
|
||||
sleepSpy.mockRestore();
|
||||
});
|
||||
});
|
||||
@@ -1,7 +1,8 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import type { MCPReconnect } from "../src/mcp/tool-bridge";
|
||||
import { isRetriableConnectionError, MCPTool } from "../src/mcp/tool-bridge";
|
||||
import { DeferredMCPTool, isRetriableConnectionError, MCPTool } from "../src/mcp/tool-bridge";
|
||||
import type { MCPServerConnection, MCPToolCallResult, MCPTransport } from "../src/mcp/types";
|
||||
import { ToolAbortError } from "../src/tools/tool-errors";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Helpers
|
||||
@@ -271,3 +272,39 @@ describe("MCPTool.execute retry on connection error", () => {
|
||||
expect(result.details?.providerName).toBe("Original");
|
||||
});
|
||||
});
|
||||
|
||||
describe("reconnect abort propagation", () => {
|
||||
const noop = () => {};
|
||||
const noCtx = {} as Parameters<MCPTool["execute"]>[3];
|
||||
const noDeferredCtx = {} as Parameters<DeferredMCPTool["execute"]>[3];
|
||||
|
||||
it("throws ToolAbortError when MCPTool reconnect is aborted", async () => {
|
||||
const failTransport = mockTransport(async () => {
|
||||
throw new Error("ECONNRESET");
|
||||
});
|
||||
const { promise } = Promise.withResolvers<MCPServerConnection | null>();
|
||||
const reconnect: MCPReconnect = async () => promise;
|
||||
|
||||
const tool = new MCPTool(makeConnection(failTransport), TOOL_DEF, reconnect);
|
||||
const controller = new AbortController();
|
||||
const pending = tool.execute("call-1", {}, noop, noCtx, controller.signal);
|
||||
controller.abort();
|
||||
|
||||
await expect(pending).rejects.toBeInstanceOf(ToolAbortError);
|
||||
});
|
||||
|
||||
it("throws ToolAbortError when DeferredMCPTool reconnect is aborted", async () => {
|
||||
const getConnection = async () => {
|
||||
throw new Error("MCP server not connected");
|
||||
};
|
||||
const { promise } = Promise.withResolvers<MCPServerConnection | null>();
|
||||
const reconnect: MCPReconnect = async () => promise;
|
||||
|
||||
const tool = new DeferredMCPTool("test-server", TOOL_DEF, getConnection, undefined, reconnect);
|
||||
const controller = new AbortController();
|
||||
const pending = tool.execute("call-1", {}, noop, noDeferredCtx, controller.signal);
|
||||
controller.abort();
|
||||
|
||||
await expect(pending).rejects.toBeInstanceOf(ToolAbortError);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user