diff --git a/src/cli/src/commands/mcp.ts b/src/cli/src/commands/mcp.ts index d3b4f9a..f6c62ec 100644 --- a/src/cli/src/commands/mcp.ts +++ b/src/cli/src/commands/mcp.ts @@ -11,6 +11,14 @@ export interface McpBridgeOptions { stderr: NodeJS.WritableStream; } +/** + * Per-request socket-inactivity timeout for the bridge. A stalled upstream must + * become a JSON-RPC error, not silence: the client cannot distinguish "still + * working" from "wedged", and its own idle limit may be half an hour away. + * Raise it for a project with genuinely long tool calls. + */ +export const BRIDGE_TIMEOUT_MS = Number(process.env['MCPCTL_MCP_TIMEOUT_MS']) || 30_000; + export function postJsonRpc( url: string, body: string, @@ -37,7 +45,7 @@ export function postJsonRpc( path: parsed.pathname, method: 'POST', headers, - timeout: 30_000, + timeout: BRIDGE_TIMEOUT_MS, }, (res) => { const chunks: Buffer[] = []; @@ -54,7 +62,7 @@ export function postJsonRpc( req.on('error', reject); req.on('timeout', () => { req.destroy(); - reject(new Error('Request timed out')); + reject(new Error(`Request timed out after ${BRIDGE_TIMEOUT_MS}ms (set MCPCTL_MCP_TIMEOUT_MS to change)`)); }); req.write(body); req.end(); @@ -128,19 +136,15 @@ export async function runMcpBridge(opts: McpBridgeOptions): Promise { const rl = createInterface({ input: stdin, crlfDelay: Infinity }); - for await (const line of rl) { - const trimmed = line.trim(); - if (!trimmed) continue; - - // Parse request ID for error responses - let requestId: unknown = null; - try { - const parsed = JSON.parse(trimmed) as Record; - requestId = parsed.id ?? null; - } catch { - // Non-JSON or notification — no id to respond to - } + /** + * In-flight requests. Dispatch is CONCURRENT after the session is + * established — see the head-of-line note below — so stdin can keep being + * read while a slow call is outstanding. + */ + const inFlight = new Set>(); + /** POST one JSON-RPC message and write whatever comes back to stdout. */ + const dispatch = async (trimmed: string, requestId: unknown): Promise => { try { const result = await postJsonRpc(endpointUrl, trimmed, sessionId, token); @@ -178,9 +182,48 @@ export async function runMcpBridge(opts: McpBridgeOptions): Promise { stdout.write(errorResponse + '\n'); } } + }; + + for await (const line of rl) { + const trimmed = line.trim(); + if (!trimmed) continue; + + // Parse request ID for error responses + let requestId: unknown = null; + try { + const parsed = JSON.parse(trimmed) as Record; + requestId = parsed.id ?? null; + } catch { + // Non-JSON or notification — no id to respond to + } + + // HEAD-OF-LINE BLOCKING: this loop used to `await` every request, so the + // bridge handled exactly one at a time. A single slow call stalled EVERY + // later request — the client saw silence rather than an error, because the + // queued requests were never even sent, so nothing could time them out. + // Observed 2026-08-05: two gitea calls sat mute until the client aborted + // them at its own 1800s idle limit, while the upstream server was healthy + // and answering other sessions in milliseconds. JSON-RPC ids exist exactly + // so responses can come back out of order; nothing here needs a queue. + // + // We still serialise until the session id exists: it comes back on the + // first response, and firing later requests without it would open a second + // upstream session. In practice a client sends `initialize` first and waits + // for its reply anyway, so this costs one round trip, not throughput. + if (sessionId === undefined) { + await dispatch(trimmed, requestId); + continue; + } + + const p = dispatch(trimmed, requestId).finally(() => inFlight.delete(p)); + inFlight.add(p); } - // stdin closed — cleanup session + // stdin closed — let outstanding work finish before tearing the session down, + // otherwise a concurrent call races the DELETE and dies with a 404. + if (inFlight.size > 0) { + await Promise.allSettled([...inFlight]); + } if (sessionId) { await sendDelete(endpointUrl, sessionId, token); } diff --git a/src/cli/tests/commands/mcp.test.ts b/src/cli/tests/commands/mcp.test.ts index 293c990..7e52ebd 100644 --- a/src/cli/tests/commands/mcp.test.ts +++ b/src/cli/tests/commands/mcp.test.ts @@ -483,3 +483,105 @@ describe('createMcpCommand', () => { expect(parsed.opts().project).toBe('my-project'); }); }); + +// --------------------------------------------------------------------------- +// Regression: head-of-line blocking (2026-08-05) +// +// The bridge used to `await` every request inside its stdin loop, so it handled +// exactly one at a time. A single slow call stalled every later request, and +// because those requests were never even sent, nothing could time them out — +// the client just saw silence until its own idle limit fired (30 min, in the +// incident that prompted this). These pin both halves of the fix: later +// requests must not queue behind a slow one, and a stalled request must produce +// a JSON-RPC error rather than nothing. +// --------------------------------------------------------------------------- +describe('MCP bridge concurrency', () => { + let srv: http.Server; + let port: number; + + function sse(id: number | string, result: unknown) { + return `event: message\ndata: ${JSON.stringify({ jsonrpc: '2.0', id, result })}\n\n`; + } + + beforeAll(async () => { + srv = http.createServer((req, res) => { + let body = ''; + req.on('data', (c) => (body += c)); + req.on('end', () => { + const msg = JSON.parse(body || '{}') as { id?: number | string; method?: string; params?: any }; + const headers: Record = { 'Content-Type': 'text/event-stream' }; + if (msg.method === 'initialize') headers['mcp-session-id'] = 'sess-1'; + // `slow` blocks far longer than `fast`, so a serial bridge would force + // fast's response to arrive second. + const delay = msg.params?.name === 'slow' ? 400 : 0; + setTimeout(() => { + res.writeHead(200, headers); + res.end(sse(msg.id ?? 0, { ok: msg.params?.name ?? msg.method })); + }, delay); + }); + }); + await new Promise((r) => srv.listen(0, '127.0.0.1', r)); + port = (srv.address() as any).port; + }); + + afterAll(async () => { + await new Promise((r) => srv.close(() => r())); + }); + + function bridge(lines: string[], out: string[]) { + const stdin = Readable.from(lines.map((l) => l + '\n')); + const stdout = new Writable({ + write(chunk, _enc, cb) { + out.push(chunk.toString().trim()); + cb(); + }, + }); + const stderr = new Writable({ write(_c, _e, cb) { cb(); } }); + return runMcpBridge({ + projectName: 'p', + mcplocalUrl: `http://127.0.0.1:${port}`, + stdin, + stdout, + stderr, + }); + } + + it('does not let a slow request block a later fast one', async () => { + const out: string[] = []; + await bridge( + [ + JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'initialize', params: {} }), + JSON.stringify({ jsonrpc: '2.0', id: 2, method: 'tools/call', params: { name: 'slow' } }), + JSON.stringify({ jsonrpc: '2.0', id: 3, method: 'tools/call', params: { name: 'fast' } }), + ], + out, + ); + const ids = out.map((l) => (JSON.parse(l) as { id: number }).id); + expect(ids).toContain(2); + expect(ids).toContain(3); + // The whole point: fast (id 3) overtakes slow (id 2). Serially it could not. + expect(ids.indexOf(3)).toBeLessThan(ids.indexOf(2)); + }); + + it('answers with a JSON-RPC error instead of silence when a request fails', async () => { + const out: string[] = []; + const stdin = Readable.from([ + JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'initialize', params: {} }) + '\n', + ]); + const stdout = new Writable({ + write(chunk, _enc, cb) { out.push(chunk.toString().trim()); cb(); }, + }); + const stderr = new Writable({ write(_c, _e, cb) { cb(); } }); + // Port with nothing on it: the POST fails fast, and the bridge must still + // emit a response carrying the original id. + await runMcpBridge({ + projectName: 'p', + mcplocalUrl: 'http://127.0.0.1:1', + stdin, stdout, stderr, + }); + expect(out.length).toBeGreaterThan(0); + const msg = JSON.parse(out[0]!) as { id: number; error?: { code: number } }; + expect(msg.id).toBe(1); + expect(msg.error?.code).toBe(-32603); + }); +});