From 2e17493b618c39666878cdbae33fbc9c0b469a0d Mon Sep 17 00:00:00 2001 From: kugouming Date: Tue, 1 Sep 2026 18:06:03 +0800 Subject: [PATCH] =?UTF-8?q?fix(routing):=20=E5=90=8E=E7=AB=AF=E4=BC=9A?= =?UTF-8?q?=E8=AF=9D=E8=BF=87=E6=9C=9F=E6=97=B6=E9=87=8D=E5=BB=BA=E8=BF=9E?= =?UTF-8?q?=E6=8E=A5=EF=BC=8C=E9=81=BF=E5=85=8D=20tools/list=20=E6=B0=B8?= =?UTF-8?q?=E4=B9=85=E4=B8=BA=E7=A9=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit SSE/HTTP 后端空闲后会话过期,其结果不是 transport 进入 ERROR,而是以 JSON-RPC -32001 错误返回。此前该错误被当作普通 Error,携带过期会话的 连接被原样放回池中反复复用,导致 tools/list 长期返回空。 - 识别 -32001 及 "session not found/expired" 消息,标记为会话过期 - 连接级错误判断纳入会话过期,失效陈旧连接而非放回池中 - 工具发现阶段对会话过期做有界重试,透明重建后端会话 - 修复 HttpTransport.doReceive 缺失 continue 导致的队列消息丢失 - 新增会话过期恢复的单元与端到端回归用例 验证:typecheck / lint / format 通过,951 个测试通过。 --- src/pool/connection-pool.ts | 7 + src/routing/tool-router.ts | 150 +++++--- src/transport/http.ts | 5 + .../backend-session-expiry.test.ts | 219 +++++++++++ tests/unit/routing/tool-router.test.ts | 360 ++++++++++++++++++ tests/unit/transport/http.test.ts | 42 ++ 6 files changed, 735 insertions(+), 48 deletions(-) create mode 100644 tests/integration/backend-session-expiry.test.ts diff --git a/src/pool/connection-pool.ts b/src/pool/connection-pool.ts index 50e8b8a..0fed2f3 100644 --- a/src/pool/connection-pool.ts +++ b/src/pool/connection-pool.ts @@ -347,6 +347,13 @@ export class ConnectionPool extends EventEmitter { this.emit('closed'); } + /** + * Get the configured maximum number of concurrent connections. + */ + public get maxConnections(): number { + return this.config.maxConnections; + } + /** * Get pool statistics * diff --git a/src/routing/tool-router.ts b/src/routing/tool-router.ts index 6586800..fa2ccde 100644 --- a/src/routing/tool-router.ts +++ b/src/routing/tool-router.ts @@ -519,10 +519,46 @@ export class ToolRouter extends EventEmitter { * @returns Promise resolving to array of tools from the service * @private */ + /** + * Whether an error indicates the backend session has expired or is no longer valid. + * + * SSE/HTTP backends often report an expired session as a JSON-RPC *error* + * (HTTP 200, code `-32001` or a "session ... not found/expired" message) rather + * than dropping the connection. A connection whose backend session is gone must + * be invalidated and re-initialized — releasing it back into the pool would keep + * reusing the stale session and fail every subsequent tools/list / tools/call. + */ + private isSessionExpiryError(error: unknown): boolean { + if (error instanceof TransportError && error.code === 'SESSION_EXPIRED') { + return true; + } + + if (error && typeof error === 'object' && 'code' in error) { + const code = (error as { code?: unknown }).code; + if (code === -32001) { + return true; + } + } + + if (error instanceof Error) { + // Ordered match only: "session" must precede the expiry signal so an + // unrelated error like "tool X not found in session Y" isn't misclassified. + if (/session\b.*\b(not found|expired)/i.test(error.message)) { + return true; + } + } + + return false; + } + /** * Whether an error indicates the connection is dead and should be removed from the pool. */ private isConnectionLevelError(error: unknown): boolean { + if (this.isSessionExpiryError(error)) { + return true; + } + if (!(error instanceof Error)) { return false; } @@ -633,58 +669,76 @@ export class ToolRouter extends EventEmitter { service: ServiceDefinition, pool: ConnectionPool ): Promise { - const connection = await pool.acquire(); - let connectionHandled = false; - try { - const timeoutMs = DEFAULT_DISCOVERY_TIMEOUT_MS; - const rawTools: unknown[] = await this.queryToolsViaMCP(connection, timeoutMs); - - const tools: Tool[] = rawTools.map((rawTool: unknown) => { - const toolObj = rawTool as { - name: string; - description?: string; - inputSchema?: { - type: 'object'; - properties: Record; - required?: string[]; + // A backend session (SSE/HTTP) can expire after idle. On a session-expiry + // failure we invalidate the stale connection and retry so discovery + // transparently re-initializes the backend session instead of surfacing an + // empty tool list to the client. A connection whose backend session expired + // still looks healthy (isConnected stays true for HTTP transports), so + // acquire() keeps handing back stale idle connections until each has been + // invalidated. Allow up to maxConnections invalidations plus one final + // attempt that forces a fresh connection. Non-session connection failures + // (backend down, timeout) still fail fast — a fresh connection won't help. + const maxAttempts = Math.max(2, pool.maxConnections + 1); + for (let attempt = 0; ; attempt++) { + const connection = await pool.acquire(); + let connectionHandled = false; + try { + const timeoutMs = DEFAULT_DISCOVERY_TIMEOUT_MS; + const rawTools: unknown[] = await this.queryToolsViaMCP(connection, timeoutMs); + + const tools: Tool[] = rawTools.map((rawTool: unknown) => { + const toolObj = rawTool as { + name: string; + description?: string; + inputSchema?: { + type: 'object'; + properties: Record; + required?: string[]; + }; }; - }; - const namespacedName = this.namespaceManager.generateNamespacedName( - service.name, - toolObj.name - ); - const enabled = this.isToolEnabled(service, toolObj.name); - - return { - name: toolObj.name, - namespacedName, - serviceName: service.name, - description: enhanceDescription(toolObj.name, service.name, toolObj.description || ''), - inputSchema: toolObj.inputSchema || { - type: 'object', - properties: {}, - required: [], - }, - enabled, - }; - }); + const namespacedName = this.namespaceManager.generateNamespacedName( + service.name, + toolObj.name + ); + const enabled = this.isToolEnabled(service, toolObj.name); - return tools; - } catch (error) { - if (this.isConnectionLevelError(error)) { - await pool.markConnectionFailed( - connection, - error instanceof Error ? error : new Error(String(error)) - ); - connectionHandled = true; - } else { + return { + name: toolObj.name, + namespacedName, + serviceName: service.name, + description: enhanceDescription(toolObj.name, service.name, toolObj.description || ''), + inputSchema: toolObj.inputSchema || { + type: 'object', + properties: {}, + required: [], + }, + enabled, + }; + }); + + return tools; + } catch (error) { + if (this.isConnectionLevelError(error)) { + await pool.markConnectionFailed( + connection, + error instanceof Error ? error : new Error(String(error)) + ); + connectionHandled = true; + // Retry only for a stale backend session: a fresh connection + // re-establishes it. Any other connection failure (backend down, + // timeout) would just fail again on a fresh connection, so fail fast. + if (this.isSessionExpiryError(error) && attempt + 1 < maxAttempts) { + continue; + } + throw error; + } pool.release(connection); connectionHandled = true; - } - throw error; - } finally { - if (!connectionHandled) { - pool.release(connection); + throw error; + } finally { + if (!connectionHandled) { + pool.release(connection); + } } } } diff --git a/src/transport/http.ts b/src/transport/http.ts index 270803f..65fb07d 100644 --- a/src/transport/http.ts +++ b/src/transport/http.ts @@ -379,6 +379,11 @@ export class HttpTransport extends BaseTransport { if (message) { yield message; } + // Loop back to drain any remaining queued messages before waiting for + // a new one. Without this, a receiver that skips a queued message (e.g. + // a notification response) would await a fresh message and never see the + // actual response already sitting in the queue. + continue; } // If receive is closed, we're done diff --git a/tests/integration/backend-session-expiry.test.ts b/tests/integration/backend-session-expiry.test.ts new file mode 100644 index 0000000..9b74b62 --- /dev/null +++ b/tests/integration/backend-session-expiry.test.ts @@ -0,0 +1,219 @@ +/** + * Integration tests: backend session expiry must not permanently break tool discovery. + * + * These tests exercise the REAL ConnectionPool + HttpTransport against a real, + * in-process Streamable HTTP MCP backend. They cover the user-facing scenario + * the unit tests mock away: + * + * backend session expires (backend answers tools/list with -32001 on the now + * stale session) → onemcp invalidates the stale connection and re-initializes + * a fresh backend session transparently → tools/list succeeds again. + * + * The mock backend expires a session deterministically (one tools/list per + * session) instead of using a clock, so the test is not timing-sensitive. + */ + +import { describe, it, expect, afterEach } from 'vitest'; +import http from 'node:http'; +import { AddressInfo } from 'node:net'; +import { ConnectionPool } from '../../src/pool/connection-pool.js'; +import { ServiceRegistry } from '../../src/registry/service-registry.js'; +import { NamespaceManager } from '../../src/namespace/manager.js'; +import { HealthMonitor } from '../../src/health/health-monitor.js'; +import { ToolRouter } from '../../src/routing/tool-router.js'; +import type { ServiceDefinition, ConnectionPoolConfig } from '../../src/types/service.js'; +import type { ConfigProvider, SystemConfig } from '../../src/types/config.js'; + +const TOOLS = [ + { name: 'alpha', description: 'Alpha tool', inputSchema: { type: 'object', properties: {} } }, + { name: 'beta', description: 'Beta tool', inputSchema: { type: 'object', properties: {} } }, +]; + +function createMockConfigProvider(): ConfigProvider { + const storedConfig: SystemConfig = { + mode: 'cli', + logLevel: 'INFO', + configDir: '/test/config', + mcpServers: {}, + connectionPool: { maxConnections: 5, idleTimeout: 60000, connectionTimeout: 30000 }, + healthCheck: { enabled: true, interval: 30000, failureThreshold: 3, autoUnload: true }, + audit: { + enabled: true, + level: 'standard', + logInput: false, + logOutput: false, + retention: { days: 30, maxSize: '1GB' }, + }, + security: { dataMasking: { enabled: true, patterns: [] } }, + }; + return { + load: async () => ({ ...storedConfig }), + save: async () => {}, + validate: () => ({ valid: true, errors: [] }), + watch: () => () => {}, + }; +} + +/** + * Start a minimal Streamable HTTP MCP backend. + * + * Each session may serve exactly one non-initialize request before it is + * reported as expired (`-32001 Session not found or expired`), modelling the + * idle-timeout behaviour of the real jymcp backend without wall-clock timing. + */ +function startMockBackend(): Promise<{ + url: string; + close: () => Promise; + initializeCount: () => number; +}> { + const sessions = new Map(); // sid -> remaining allowed requests + let sidCounter = 0; + let initializeCount = 0; + + const server = http.createServer((req, res) => { + let raw = ''; + req.on('data', (chunk) => (raw += chunk)); + req.on('end', () => { + let msg: Record; + try { + msg = JSON.parse(raw) as Record; + } catch { + res.writeHead(400).end(); + return; + } + + // Notifications: no id -> 202 empty (MCP Streamable HTTP spec) + if (msg['id'] === undefined || msg['id'] === null) { + res.writeHead(202, { 'Content-Type': 'text/event-stream' }); + res.end(); + return; + } + + const sendJson = (status: number, body: unknown, headers: Record = {}) => { + const json = JSON.stringify(body); + res.writeHead(status, { + 'Content-Type': 'application/json', + 'Content-Length': Buffer.byteLength(json), + ...headers, + }); + res.end(json); + }; + + if (msg['method'] === 'initialize') { + initializeCount++; + const sessionId = `sess-${++sidCounter}`; + sessions.set(sessionId, 1); + sendJson( + 200, + { + jsonrpc: '2.0', + id: msg['id'], + result: { + protocolVersion: '2024-11-05', + capabilities: { tools: { listChanged: true } }, + serverInfo: { name: 'mock-backend', version: '1.0.0' }, + }, + }, + { 'mcp-session-id': sessionId } + ); + return; + } + + const sessionId = Array.isArray(req.headers['mcp-session-id']) + ? String(req.headers['mcp-session-id'][0]) + : req.headers['mcp-session-id']; + const remaining = typeof sessionId === 'string' ? sessions.get(sessionId) : undefined; + + if (remaining !== undefined && remaining > 0) { + sessions.set(sessionId as string, remaining - 1); + if (msg['method'] === 'tools/list') { + sendJson(200, { jsonrpc: '2.0', id: msg['id'], result: { tools: TOOLS } }); + } else { + sendJson(200, { jsonrpc: '2.0', id: msg['id'], result: {} }); + } + return; + } + + sendJson(200, { + jsonrpc: '2.0', + id: msg['id'], + error: { + code: -32001, + message: 'Session not found or expired. Please send initialize again.', + }, + }); + }); + }); + + return new Promise((resolve) => { + server.listen(0, '127.0.0.1', () => { + const port = (server.address() as AddressInfo).port; + resolve({ + url: `http://127.0.0.1:${port}/mcp`, + initializeCount: () => initializeCount, + close: () => new Promise((r) => server.close(() => r())), + }); + }); + }); +} + +describe('Backend session expiry recovery (integration)', () => { + let backend: Awaited> | undefined; + + afterEach(async () => { + if (backend) { + await backend.close().catch(() => {}); + backend = undefined; + } + }); + + it('recovers tools/list after the backend session expires', async () => { + backend = await startMockBackend(); + + const service: ServiceDefinition = { + name: 'mock-http', + enabled: true, + tags: [], + transport: 'http', + url: backend.url, + connectionPool: { maxConnections: 2, idleTimeout: 60000, connectionTimeout: 10000 }, + } as ServiceDefinition; + const poolConfig: ConnectionPoolConfig = { + maxConnections: 2, + idleTimeout: 60000, + connectionTimeout: 10000, + }; + + const configProvider = createMockConfigProvider(); + const serviceRegistry = new ServiceRegistry(configProvider); + await serviceRegistry.initialize(); + await serviceRegistry.register(service); + + const namespaceManager = new NamespaceManager(); + const healthMonitor = new HealthMonitor(serviceRegistry); + const toolRouter = new ToolRouter(serviceRegistry, namespaceManager, healthMonitor); + + const pool = new ConnectionPool(service, poolConfig); + pool.on('error', () => {}); + toolRouter.registerConnectionPool(service.name, pool); + + try { + // First discovery uses a fresh session and succeeds. + const tools1 = await toolRouter.discoverTools(); + expect(tools1.map((t) => t.name)).toEqual(['alpha', 'beta']); + + // The mock now considers that session expired. + toolRouter.invalidateCache(); + + // Second discovery must transparently re-initialize a fresh session and + // still return the tools (previously this returned an empty list forever). + const tools2 = await toolRouter.discoverTools(); + expect(tools2.map((t) => t.name)).toEqual(['alpha', 'beta']); + + // Exactly two sessions were created: the original + one re-initialized. + expect(backend.initializeCount()).toBe(2); + } finally { + await pool.closeAll().catch(() => {}); + } + }, 30000); +}); diff --git a/tests/unit/routing/tool-router.test.ts b/tests/unit/routing/tool-router.test.ts index 58ae51f..6d2fad7 100644 --- a/tests/unit/routing/tool-router.test.ts +++ b/tests/unit/routing/tool-router.test.ts @@ -7,6 +7,7 @@ import { ToolRouter } from '../../../src/routing/tool-router'; import { ServiceRegistry } from '../../../src/registry/service-registry'; import { NamespaceManager } from '../../../src/namespace/manager'; import { HealthMonitor } from '../../../src/health/health-monitor'; +import { TransportError } from '../../../src/transport/base'; import type { ServiceDefinition } from '../../../src/types/service'; import type { ConnectionPool } from '../../../src/pool/connection-pool'; import type { ConfigProvider, SystemConfig } from '../../../src/types/config'; @@ -1838,6 +1839,365 @@ describe('ToolRouter', () => { }); }); + describe('session expiry recovery', () => { + it('should re-establish the backend session when discovery hits an expired session', async () => { + // Register an enabled service + const service: ServiceDefinition = { + name: 'test-service', + enabled: true, + tags: [], + transport: 'http', + url: 'http://127.0.0.1:9999/mcp', + connectionPool: { + maxConnections: 5, + idleTimeout: 60000, + connectionTimeout: 30000, + }, + }; + + await serviceRegistry.register(service); + + const mockTool = { + name: 'real_tool', + description: 'Real tool', + inputSchema: { type: 'object' as const, properties: {} }, + }; + + // Build a transport that answers tools/list with either a session-expired + // error (-32001, the JSON-RPC error SSE/HTTP backends return) or success. + const makeTransport = (respondWithExpiredSession: boolean) => { + let requestId = ''; + return { + send: vi.fn(async (request: { id?: string | number }) => { + requestId = String(request.id ?? ''); + }), + receive: vi.fn().mockReturnValue({ + next: vi.fn().mockImplementation(async () => ({ + done: false as const, + value: respondWithExpiredSession + ? { + jsonrpc: '2.0', + id: requestId, + error: { + code: -32001, + message: 'Session not found or expired. Please send initialize again.', + }, + } + : { + jsonrpc: '2.0', + id: requestId, + result: { tools: [mockTool] }, + }, + })), + return: vi.fn().mockResolvedValue({ done: true }), + }), + close: vi.fn().mockResolvedValue(undefined), + getType: vi.fn().mockReturnValue('http'), + isConnected: vi.fn().mockReturnValue(true), + }; + }; + + const staleConnection = { + id: 'conn-stale', + transport: makeTransport(true), + state: 'idle' as const, + lastUsed: new Date(), + createdAt: new Date(), + }; + const freshConnection = { + id: 'conn-fresh', + transport: makeTransport(false), + state: 'idle' as const, + lastUsed: new Date(), + createdAt: new Date(), + }; + + const acquire = vi + .fn() + .mockResolvedValueOnce(staleConnection) + .mockResolvedValueOnce(freshConnection); + const mockPool = { + acquire, + release: vi.fn(), + markConnectionFailed: vi.fn().mockResolvedValue(undefined), + maxConnections: 1, + } as any; + + toolRouter.registerConnectionPool('test-service', mockPool); + + const tools = await toolRouter.discoverTools(); + + // The stale session connection must be invalidated and a fresh one acquired. + expect(acquire).toHaveBeenCalledTimes(2); + expect(mockPool.markConnectionFailed).toHaveBeenCalledTimes(1); + expect(mockPool.markConnectionFailed).toHaveBeenCalledWith( + staleConnection, + expect.any(Error) + ); + expect(mockPool.release).not.toHaveBeenCalledWith(staleConnection); + + // Discovery must recover transparently instead of returning an empty list. + expect(tools).toHaveLength(1); + expect(tools[0]?.namespacedName).toBe('test-service__real_tool'); + }); + + it('should detect session expiry from the error message alone (no numeric code)', async () => { + const service: ServiceDefinition = { + name: 'test-service', + enabled: true, + tags: [], + transport: 'http', + url: 'http://127.0.0.1:9999/mcp', + connectionPool: { maxConnections: 5, idleTimeout: 60000, connectionTimeout: 30000 }, + }; + await serviceRegistry.register(service); + + const mockTool = { + name: 'real_tool', + description: 'Real tool', + inputSchema: { type: 'object' as const, properties: {} }, + }; + const makeTransport = (expired: boolean) => { + let requestId = ''; + return { + send: vi.fn(async (request: { id?: string | number }) => { + requestId = String(request.id ?? ''); + }), + receive: vi.fn().mockReturnValue({ + next: vi.fn().mockImplementation(async () => ({ + done: false as const, + value: expired + ? { + jsonrpc: '2.0', + id: requestId, + // No `code` field — detection must rely on the message text. + error: { + message: 'Session not found or expired. Please send initialize again.', + }, + } + : { jsonrpc: '2.0', id: requestId, result: { tools: [mockTool] } }, + })), + return: vi.fn().mockResolvedValue({ done: true }), + }), + close: vi.fn().mockResolvedValue(undefined), + getType: vi.fn().mockReturnValue('http'), + isConnected: vi.fn().mockReturnValue(true), + }; + }; + const staleConnection = { + id: 'conn-stale', + transport: makeTransport(true), + state: 'idle' as const, + lastUsed: new Date(), + createdAt: new Date(), + }; + const freshConnection = { + id: 'conn-fresh', + transport: makeTransport(false), + state: 'idle' as const, + lastUsed: new Date(), + createdAt: new Date(), + }; + + const acquire = vi + .fn() + .mockResolvedValueOnce(staleConnection) + .mockResolvedValueOnce(freshConnection); + const mockPool = { + acquire, + release: vi.fn(), + markConnectionFailed: vi.fn().mockResolvedValue(undefined), + maxConnections: 1, + } as any; + toolRouter.registerConnectionPool('test-service', mockPool); + + const tools = await toolRouter.discoverTools(); + + expect(acquire).toHaveBeenCalledTimes(2); + expect(mockPool.markConnectionFailed).toHaveBeenCalledTimes(1); + expect(tools).toHaveLength(1); + }); + + it('should fail fast (no retry) on a non-session connection error', async () => { + const service: ServiceDefinition = { + name: 'test-service', + enabled: true, + tags: [], + transport: 'http', + url: 'http://127.0.0.1:9999/mcp', + connectionPool: { maxConnections: 5, idleTimeout: 60000, connectionTimeout: 30000 }, + }; + await serviceRegistry.register(service); + + const staleConnection = { + id: 'conn-1', + transport: { + send: vi + .fn() + .mockRejectedValue(new TransportError('Response timeout', 'RESPONSE_TIMEOUT')), + receive: vi.fn(), + close: vi.fn(), + getType: vi.fn().mockReturnValue('http'), + isConnected: vi.fn().mockReturnValue(true), + }, + state: 'idle' as const, + lastUsed: new Date(), + createdAt: new Date(), + }; + + const acquire = vi.fn().mockResolvedValue(staleConnection); + const mockPool = { + acquire, + release: vi.fn(), + markConnectionFailed: vi.fn().mockResolvedValue(undefined), + maxConnections: 2, + } as any; + toolRouter.registerConnectionPool('test-service', mockPool); + + const errorSpy = vi.fn(); + toolRouter.on('toolDiscoveryError', errorSpy); + + const tools = await toolRouter.discoverTools(); + + // A real transport failure should not trigger a second backend spawn. + expect(acquire).toHaveBeenCalledTimes(1); + expect(mockPool.markConnectionFailed).toHaveBeenCalledTimes(1); + expect(tools).toHaveLength(0); + expect(errorSpy).toHaveBeenCalled(); + }); + + it('should not treat a generic "not found in session" error as session expiry', async () => { + const service: ServiceDefinition = { + name: 'test-service', + enabled: true, + tags: [], + transport: 'http', + url: 'http://127.0.0.1:9999/mcp', + connectionPool: { maxConnections: 5, idleTimeout: 60000, connectionTimeout: 30000 }, + }; + await serviceRegistry.register(service); + + let requestId = ''; + const connection = { + id: 'conn-1', + transport: { + send: vi.fn(async (request: { id?: string | number }) => { + requestId = String(request.id ?? ''); + }), + receive: vi.fn().mockReturnValue({ + next: vi.fn().mockImplementation(async () => ({ + done: false as const, + value: { + jsonrpc: '2.0', + id: requestId, + // "not found" precedes "session" — this is NOT a session-expiry + // signal and must not invalidate the connection. + error: { message: "Tool 'foo' not found in session local" }, + }, + })), + return: vi.fn().mockResolvedValue({ done: true }), + }), + close: vi.fn().mockResolvedValue(undefined), + getType: vi.fn().mockReturnValue('http'), + isConnected: vi.fn().mockReturnValue(true), + }, + state: 'idle' as const, + lastUsed: new Date(), + createdAt: new Date(), + }; + + const mockPool = { + acquire: vi.fn().mockResolvedValue(connection), + release: vi.fn(), + markConnectionFailed: vi.fn().mockResolvedValue(undefined), + maxConnections: 1, + } as any; + toolRouter.registerConnectionPool('test-service', mockPool); + + await toolRouter.discoverTools(); + + // Not session expiry: the connection is released intact, not dropped. + expect(mockPool.release).toHaveBeenCalledWith(connection); + expect(mockPool.markConnectionFailed).not.toHaveBeenCalled(); + }); + + it('should invalidate the connection when a tool call hits an expired session', async () => { + const service: ServiceDefinition = { + name: 'test-service', + enabled: true, + tags: [], + transport: 'http', + url: 'http://127.0.0.1:9999/mcp', + connectionPool: { maxConnections: 5, idleTimeout: 60000, connectionTimeout: 30000 }, + }; + await serviceRegistry.register(service); + + let requestId = ''; + const mockTransport = { + send: vi.fn(async (request: { id?: string | number }) => { + requestId = String(request.id ?? ''); + }), + receive: vi.fn().mockReturnValue({ + next: vi.fn().mockImplementation(async () => ({ + done: false as const, + value: { + jsonrpc: '2.0', + id: requestId, + error: { + code: -32001, + message: 'Session not found or expired. Please send initialize again.', + }, + }, + })), + return: vi.fn().mockResolvedValue({ done: true }), + }), + close: vi.fn().mockResolvedValue(undefined), + getType: vi.fn().mockReturnValue('http'), + isConnected: vi.fn().mockReturnValue(true), + }; + const mockConnection = { + id: 'conn-1', + transport: mockTransport, + state: 'idle' as const, + lastUsed: new Date(), + createdAt: new Date(), + }; + const mockPool = { + acquire: vi.fn().mockResolvedValue(mockConnection), + release: vi.fn(), + markConnectionFailed: vi.fn().mockResolvedValue(undefined), + maxConnections: 1, + } as any; + toolRouter.registerConnectionPool('test-service', mockPool); + + const mockTool: Tool = { + name: 'test_tool', + namespacedName: 'test-service__test_tool', + serviceName: 'test-service', + description: 'Test tool', + inputSchema: { type: 'object', properties: {} }, + enabled: true, + }; + const findToolSpy = vi.spyOn(toolRouter as any, 'findTool').mockResolvedValue(mockTool); + + const context: RequestContext = { + requestId: 'req-1', + correlationId: 'corr-1', + timestamp: new Date(), + }; + + await expect(toolRouter.callTool('test-service__test_tool', {}, context)).rejects.toThrow( + 'Session not found' + ); + + // The stale session connection must be dropped so the next call reconnects. + expect(mockPool.markConnectionFailed).toHaveBeenCalledTimes(1); + expect(mockPool.markConnectionFailed).toHaveBeenCalledWith(mockConnection, expect.any(Error)); + + findToolSpy.mockRestore(); + }); + }); + describe('discoverTools - cache invalidation events', () => { it('should emit cacheInvalidated event when cache is populated from empty state', async () => { // Register an enabled service diff --git a/tests/unit/transport/http.test.ts b/tests/unit/transport/http.test.ts index f73f734..6c2d4e5 100644 --- a/tests/unit/transport/http.test.ts +++ b/tests/unit/transport/http.test.ts @@ -355,6 +355,48 @@ describe('HttpTransport', () => { await transport.close(); }); + it('should drain multiple queued responses in order', async () => { + // Two sends complete before any receive: both responses are buffered in + // messageQueue. The receiver must yield them BOTH in order — a missing + // `continue` in doReceive() would leave the second response stranded and + // the second next() would hang forever. + const notificationResponse = { jsonrpc: '2.0', result: {} }; + const toolResponse = { jsonrpc: '2.0', id: 42, result: { tools: [] } }; + + const makeResponse = (body: unknown) => ({ + ok: true, + status: 200, + statusText: 'OK', + headers: { + get: vi.fn().mockReturnValue(null), + }, + text: vi.fn().mockResolvedValue(JSON.stringify(body)), + }); + + vi.mocked(fetch) + .mockResolvedValueOnce(makeResponse(notificationResponse) as any) + .mockResolvedValueOnce(makeResponse(toolResponse) as any); + + const transport = new HttpTransport({ + url: 'http://localhost:300/rpc', + mode: 'http', + }); + + await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized', params: {} }); + await transport.send({ jsonrpc: '2.0', id: 42, method: 'tools/list', params: {} }); + + const iterator = transport.receive(); + const first = await iterator.next(); + const second = await iterator.next(); + + expect(first.done).toBe(false); + expect(first.value).toEqual(notificationResponse); + expect(second.done).toBe(false); + expect(second.value).toEqual(toolResponse); + + await transport.close(); + }); + it('should handle network errors', async () => { vi.mocked(fetch).mockRejectedValue(new Error('Network error'));