diff --git a/src/mcplocal/src/http/index.ts b/src/mcplocal/src/http/index.ts index 274a655..ffdb7d4 100644 --- a/src/mcplocal/src/http/index.ts +++ b/src/mcplocal/src/http/index.ts @@ -2,7 +2,7 @@ export { createHttpServer } from './server.js'; export type { HttpServerDeps } from './server.js'; export { loadHttpConfig } from './config.js'; export type { HttpConfig } from './config.js'; -export { McpdClient, AuthenticationError, ConnectionError } from './mcpd-client.js'; +export { McpdClient, AuthenticationError, ConnectionError, UpstreamTimeoutError } from './mcpd-client.js'; export { registerProxyRoutes } from './routes/proxy.js'; export { registerMcpEndpoint } from './mcp-endpoint.js'; export { registerProjectMcpEndpoint } from './project-mcp-endpoint.js'; diff --git a/src/mcplocal/src/http/mcpd-client.ts b/src/mcplocal/src/http/mcpd-client.ts index 53511e5..8e755e4 100644 --- a/src/mcplocal/src/http/mcpd-client.ts +++ b/src/mcplocal/src/http/mcpd-client.ts @@ -20,9 +20,41 @@ export class ConnectionError extends Error { } } +/** + * Thrown when mcpd was reachable but did not finish in time. + * + * Deliberately NOT a ConnectionError. Folding timeouts into "cannot connect" + * is what made this class of failure so expensive to diagnose: mcpd answered + * /healthz in 32ms while the proxy insisted the daemon was down. A timeout and + * an unreachable daemon need different messages and different status codes. + */ +export class UpstreamTimeoutError extends Error { + constructor(readonly url: string, readonly timeoutMs: number) { + super(`mcpd did not respond within ${String(timeoutMs)}ms: ${url}`); + this.name = 'UpstreamTimeoutError'; + } +} + +/** True when `err` is an AbortSignal.timeout() firing. */ +function isTimeout(err: unknown): boolean { + return err instanceof DOMException && err.name === 'TimeoutError'; +} + /** Default timeout for mcpd requests (ms). Prevents indefinite hangs on slow upstream tool calls. */ export const DEFAULT_TIMEOUT_MS = 30_000; +/** + * Budget for routes that are *expected* to run long: agent/project chat and + * raw inference. An agent turn is a multi-turn tool-use loop and legitimately + * runs for minutes, so the 30s default is not a safety net there — it is a + * guaranteed failure. Matches `STREAM_TIMEOUT_MS` in the CLI's chat command + * (src/cli/src/commands/chat.ts), which already allowed 10 minutes; mcplocal + * sitting in the middle with 30s was the binding constraint. + * + * Override with `MCPLOCAL_LONG_TIMEOUT_MS`. + */ +export const LONG_RUNNING_TIMEOUT_MS = Number(process.env['MCPLOCAL_LONG_TIMEOUT_MS']) || 600_000; + /** * Discovery-class operations (tools/list, resources/list, prompts/list) should not share * the full tool-call timeout budget — a single dead upstream would stall session init for @@ -121,9 +153,7 @@ export class McpdClient { try { res = await fetch(url, init); } catch (err: unknown) { - if (err instanceof DOMException && err.name === 'TimeoutError') { - throw new ConnectionError(this.baseUrl, new Error(`Request timed out after ${this.timeoutMs}ms`)); - } + if (isTimeout(err)) throw new UpstreamTimeoutError(this.baseUrl, this.timeoutMs); throw new ConnectionError(this.baseUrl, err); } @@ -131,7 +161,18 @@ export class McpdClient { throw new AuthenticationError(); } - const text = await res.text(); + // The body read MUST be inside a try. mcpd writes SSE headers immediately + // on chat routes, so fetch() resolves long before the turn finishes and the + // abort lands here instead — previously escaping as a raw DOMException and + // surfacing to the user as an opaque `500 code:23`. + let text: string; + try { + text = await res.text(); + } catch (err: unknown) { + if (isTimeout(err)) throw new UpstreamTimeoutError(this.baseUrl, this.timeoutMs); + throw new ConnectionError(this.baseUrl, err); + } + let parsed: unknown; try { parsed = JSON.parse(text); @@ -142,6 +183,51 @@ export class McpdClient { return { status: res.status, body: parsed }; } + /** + * Forward a request and hand back the raw Response, body unread. + * + * `forward()` buffers through `res.text()`, which is fine for CRUD but + * defeats streaming entirely: an SSE chat arrives at the client as one blob + * after the turn ends, so the token-by-token output the CLI draws never + * appears. Streaming routes use this instead and pipe the body straight + * through. + */ + async forwardStream( + method: string, + path: string, + query: string, + body: unknown | undefined, + authOverride?: string, + ): Promise { + const url = `${this.baseUrl}${path}${query ? `?${query}` : ''}`; + const headers: Record = { + ...this.extraHeaders, + 'Authorization': `Bearer ${authOverride ?? this.token}`, + // Accept both: mcpd picks SSE or JSON based on the request's `stream` flag. + 'Accept': 'text/event-stream, application/json', + }; + + const init: RequestInit = { + method, + headers, + signal: AbortSignal.timeout(this.timeoutMs), + }; + if (body !== undefined && body !== null && method !== 'GET' && method !== 'HEAD') { + headers['Content-Type'] = 'application/json'; + init.body = JSON.stringify(body); + } + + try { + const res = await fetch(url, init); + if (res.status === 401) throw new AuthenticationError(); + return res; + } catch (err: unknown) { + if (err instanceof AuthenticationError) throw err; + if (isTimeout(err)) throw new UpstreamTimeoutError(this.baseUrl, this.timeoutMs); + throw new ConnectionError(this.baseUrl, err); + } + } + private async request(method: string, path: string, body?: unknown): Promise { const result = await this.forward(method, path, '', body); diff --git a/src/mcplocal/src/http/routes/proxy.ts b/src/mcplocal/src/http/routes/proxy.ts index 985f6e6..9b2bebd 100644 --- a/src/mcplocal/src/http/routes/proxy.ts +++ b/src/mcplocal/src/http/routes/proxy.ts @@ -1,10 +1,62 @@ /** * Catch-all proxy route that forwards /api/v1/* requests to mcpd. */ -import type { FastifyInstance } from 'fastify'; -import { AuthenticationError, ConnectionError } from '../mcpd-client.js'; +import { Readable } from 'node:stream'; + +import type { FastifyInstance, FastifyReply } from 'fastify'; + +import { AuthenticationError, ConnectionError, UpstreamTimeoutError, LONG_RUNNING_TIMEOUT_MS } from '../mcpd-client.js'; import type { McpdClient } from '../mcpd-client.js'; +/** + * Routes that are expected to run long and/or stream. + * + * An agent turn is a multi-turn tool-use loop — minutes, not seconds — so the + * 30s default budget guarantees failure rather than guarding against it. These + * also stream SSE, which must be piped rather than buffered or the client sees + * one blob at the end instead of live output. + */ +const LONG_RUNNING = [ + /^\/api\/v1\/agents\/[^/]+\/chat\b/, + /^\/api\/v1\/projects\/[^/]+\/chat\b/, + /^\/api\/v1\/llms\/[^/]+\/infer\b/, + /^\/api\/v1\/inference-tasks\/[^/]+\/stream\b/, +]; + +function isLongRunning(path: string): boolean { + return LONG_RUNNING.some((re) => re.test(path)); +} + +/** Headers worth preserving from mcpd; everything else is re-derived by Fastify. */ +const PASSTHROUGH_HEADERS = ['content-type', 'cache-control', 'x-accel-buffering']; + +function sendUpstreamError(reply: FastifyReply, err: unknown): FastifyReply | undefined { + if (err instanceof AuthenticationError) { + return reply.code(401).send({ + error: 'unauthorized', + message: 'Authentication with mcpd failed. Run `mcpctl login` to refresh your token.', + }); + } + if (err instanceof UpstreamTimeoutError) { + // 504, not 503 — mcpd was reachable, it just did not finish. Reporting this + // as "cannot reach mcpd" sent a previous debugging session chasing a + // network fault while /healthz answered in 32ms. + return reply.code(504).send({ + error: 'upstream_timeout', + message: + `mcpd did not respond within ${String(err.timeoutMs)}ms. The daemon is reachable — the ` + + 'request itself ran long. Raise MCPLOCAL_LONG_TIMEOUT_MS if this is a legitimately slow turn.', + }); + } + if (err instanceof ConnectionError) { + return reply.code(503).send({ + error: 'service_unavailable', + message: 'Cannot reach mcpd daemon. Is it running?', + }); + } + return undefined; +} + export function registerProxyRoutes(app: FastifyInstance, client: McpdClient): void { app.all('/api/v1/*', async (request, reply) => { const path = (request.url.split('?')[0]) ?? '/'; @@ -19,25 +71,78 @@ export function registerProxyRoutes(app: FastifyInstance, client: McpdClient): v // Forward the user's auth token to mcpd so RBAC applies per-user. // If no user token is present, mcpd will use its auth hook to reject. const authHeader = request.headers['authorization'] as string | undefined; - const userToken = authHeader?.startsWith('Bearer ') ? authHeader.slice(7) : undefined; + const userToken = authHeader !== undefined && authHeader.startsWith('Bearer ') + ? authHeader.slice(7) + : undefined; + + if (isLongRunning(path)) { + return proxyStreaming(reply, client, request.method, path, querystring, body, userToken); + } try { const result = await client.forward(request.method, path, querystring, body, userToken); return reply.code(result.status).send(result.body); } catch (err: unknown) { - if (err instanceof AuthenticationError) { - return reply.code(401).send({ - error: 'unauthorized', - message: 'Authentication with mcpd failed. Run `mcpctl login` to refresh your token.', - }); - } - if (err instanceof ConnectionError) { - return reply.code(503).send({ - error: 'service_unavailable', - message: 'Cannot reach mcpd daemon. Is it running?', - }); - } + const handled = sendUpstreamError(reply, err); + if (handled) return handled; throw err; } }); } + +/** + * Pipe a long-running response straight through, headers and all. + * + * Hijacks the reply so Fastify does not try to serialize a stream, then copies + * mcpd's status and content-type before piping. `x-accel-buffering` matters: + * mcpd sets it to `no` so intermediaries don't buffer SSE, and dropping it here + * would reintroduce the exact stall we are fixing. + */ +async function proxyStreaming( + reply: FastifyReply, + client: McpdClient, + method: string, + path: string, + querystring: string, + body: unknown, + userToken: string | undefined, +): Promise { + const longClient = client.withTimeout(LONG_RUNNING_TIMEOUT_MS); + + let res: Response; + try { + res = await longClient.forwardStream(method, path, querystring, body, userToken); + } catch (err: unknown) { + const handled = sendUpstreamError(reply, err); + if (handled) return; + throw err; + } + + const headers: Record = {}; + for (const name of PASSTHROUGH_HEADERS) { + const value = res.headers.get(name); + if (value !== null) headers[name] = value; + } + + reply.hijack(); + reply.raw.writeHead(res.status, headers); + + if (res.body === null) { + reply.raw.end(); + return; + } + + try { + // Node's Readable.fromWeb bridges the fetch ReadableStream onto the socket. + await new Promise((resolve, reject) => { + const upstream = Readable.fromWeb(res.body as Parameters[0]); + upstream.on('error', reject); + reply.raw.on('close', () => { upstream.destroy(); resolve(); }); + upstream.pipe(reply.raw).on('finish', resolve).on('error', reject); + }); + } catch { + // Headers are already on the wire, so there is no status left to change. + // Close the socket; the client surfaces the truncated stream. + if (!reply.raw.writableEnded) reply.raw.end(); + } +} diff --git a/src/mcplocal/src/index.ts b/src/mcplocal/src/index.ts index 489e0ba..dbc01d4 100644 --- a/src/mcplocal/src/index.ts +++ b/src/mcplocal/src/index.ts @@ -11,7 +11,7 @@ export type { MainResult } from './main.js'; export { ProviderRegistry } from './providers/index.js'; export type { LlmProvider, CompletionOptions, CompletionResult, ChatMessage } from './providers/index.js'; export { OpenAiProvider, AnthropicProvider, OllamaProvider, GeminiCliProvider, DeepSeekProvider } from './providers/index.js'; -export { createHttpServer, loadHttpConfig, McpdClient, AuthenticationError, ConnectionError, registerProxyRoutes } from './http/index.js'; +export { createHttpServer, loadHttpConfig, McpdClient, AuthenticationError, ConnectionError, UpstreamTimeoutError, registerProxyRoutes } from './http/index.js'; export type { HttpConfig, HttpServerDeps } from './http/index.js'; export type { JsonRpcRequest, diff --git a/src/mcplocal/tests/mcpd-client.test.ts b/src/mcplocal/tests/mcpd-client.test.ts index c9a50d4..7dac059 100644 --- a/src/mcplocal/tests/mcpd-client.test.ts +++ b/src/mcplocal/tests/mcpd-client.test.ts @@ -1,6 +1,6 @@ import { describe, it, expect, afterAll, afterEach } from 'vitest'; import http from 'node:http'; -import { McpdClient, ConnectionError } from '../src/http/mcpd-client.js'; +import { McpdClient, ConnectionError, UpstreamTimeoutError } from '../src/http/mcpd-client.js'; /** * Create a local HTTP server for testing McpdClient behavior. @@ -85,7 +85,7 @@ describe('McpdClient', () => { // ── Timeout behavior ── - it('times out on slow responses and throws ConnectionError', async () => { + it('times out on slow responses and throws UpstreamTimeoutError', async () => { const { server, url } = await createTestServer((_req, _res) => { // Never respond — simulates a hanging upstream tool call }); @@ -96,7 +96,7 @@ describe('McpdClient', () => { const start = Date.now(); await expect(client.post('/api/v1/mcp/proxy', { serverId: 's1' })).rejects.toThrow( - /timed out/, + /did not respond within/, ); const elapsed = Date.now() - start; @@ -105,7 +105,7 @@ describe('McpdClient', () => { expect(elapsed).toBeLessThan(3000); }); - it('timeout error is a ConnectionError with descriptive message', async () => { + it('timeout is NOT a ConnectionError — a slow daemon is not an absent one', async () => { const { server, url } = await createTestServer((_req, _res) => { // Never respond }); @@ -117,8 +117,12 @@ describe('McpdClient', () => { await client.get('/test'); expect.unreachable('Should have thrown'); } catch (err) { - expect(err).toBeInstanceOf(ConnectionError); - expect((err as Error).message).toContain('Request timed out after 200ms'); + // Reporting a timeout as "cannot connect" is what sent a previous + // debugging session chasing a network fault that did not exist. + expect(err).toBeInstanceOf(UpstreamTimeoutError); + expect(err).not.toBeInstanceOf(ConnectionError); + expect((err as UpstreamTimeoutError).timeoutMs).toBe(200); + expect((err as Error).message).toContain('did not respond within 200ms'); } }); @@ -146,7 +150,7 @@ describe('McpdClient', () => { const derived = client.withHeaders({ 'X-Custom': 'val' }); const start = Date.now(); - await expect(derived.get('/test')).rejects.toThrow(/timed out/); + await expect(derived.get('/test')).rejects.toThrow(/did not respond within/); const elapsed = Date.now() - start; expect(elapsed).toBeLessThan(2000); }); diff --git a/src/mcplocal/tests/proxy-long-running.test.ts b/src/mcplocal/tests/proxy-long-running.test.ts new file mode 100644 index 0000000..a793084 --- /dev/null +++ b/src/mcplocal/tests/proxy-long-running.test.ts @@ -0,0 +1,255 @@ +import http from 'node:http'; + +import Fastify, { type FastifyInstance } from 'fastify'; +import { describe, it, expect, afterEach } from 'vitest'; + +import { + McpdClient, + UpstreamTimeoutError, + ConnectionError, + LONG_RUNNING_TIMEOUT_MS, + DEFAULT_TIMEOUT_MS, +} from '../src/http/mcpd-client.js'; +import { registerProxyRoutes } from '../src/http/routes/proxy.js'; + +/** + * Regression cover for the 30s proxy timeout that made `mcpctl chat` fail with + * a misleading "Cannot reach mcpd daemon" 503 while mcpd was answering + * /healthz in 32ms. + * + * Three separate defects are pinned here: + * 1. chat routes inherited the 30s CRUD budget, so any turn longer than 30s + * failed — and an agent turn is a tool-use loop that routinely exceeds it; + * 2. a timeout was reported as a connection failure, sending diagnosis after + * a network fault that did not exist; + * 3. SSE was buffered through res.text(), so streaming never reached the + * client even when the turn finished in time. + */ +let app: FastifyInstance | null = null; +let upstream: FastifyInstance | null = null; + +afterEach(async () => { + if (app) { await app.close(); app = null; } + if (upstream) { await upstream.close(); upstream = null; } +}); + +/** A stand-in mcpd. Returns its base URL. */ +async function startUpstream(register: (a: FastifyInstance) => void): Promise { + upstream = Fastify(); + register(upstream); + await upstream.listen({ port: 0, host: '127.0.0.1' }); + const addr = upstream.server.address(); + if (addr === null || typeof addr === 'string') throw new Error('no address'); + return `http://127.0.0.1:${String(addr.port)}`; +} + +async function startProxy(baseUrl: string, timeoutMs?: number): Promise { + app = Fastify(); + registerProxyRoutes(app, new McpdClient(baseUrl, 'test-token', {}, timeoutMs)); + await app.ready(); + return app; +} + +describe('proxy — long-running route budget', () => { + it('gives chat routes the long budget, not the 30s CRUD default', () => { + // The constants themselves are the contract: a 30s cap on an agent turn is + // a guaranteed failure, not a safety net. + expect(DEFAULT_TIMEOUT_MS).toBe(30_000); + expect(LONG_RUNNING_TIMEOUT_MS).toBeGreaterThanOrEqual(600_000); + }); + + it('does not abort an agent chat that outlives the CRUD budget', async () => { + const base = await startUpstream((a) => { + a.post('/api/v1/agents/:name/chat', async () => { + // Longer than the (deliberately tiny) CRUD budget below. Before the + // fix this inherited that budget and 503'd. + await new Promise((r) => setTimeout(r, 250)); + return { answer: 'pong' }; + }); + }); + // CRUD budget of 50ms — a chat route must NOT inherit it. + const proxy = await startProxy(base, 50); + + const res = await proxy.inject({ + method: 'POST', + url: '/api/v1/agents/reviewer/chat', + payload: { message: 'hi' }, + }); + + expect(res.statusCode).toBe(200); + expect(res.json()).toEqual({ answer: 'pong' }); + }); + + it('still applies the short budget to ordinary CRUD routes', async () => { + const base = await startUpstream((a) => { + a.get('/api/v1/servers', async () => { + await new Promise((r) => setTimeout(r, 300)); + return []; + }); + }); + const proxy = await startProxy(base, 50); + + const res = await proxy.inject({ method: 'GET', url: '/api/v1/servers' }); + // Times out — and is now reported honestly as a timeout, not a connection fault. + expect(res.statusCode).toBe(504); + expect(res.json().error).toBe('upstream_timeout'); + }); + + it('reports a timeout as 504, never as "cannot reach mcpd"', async () => { + const base = await startUpstream((a) => { + a.get('/api/v1/servers', async () => { + await new Promise((r) => setTimeout(r, 300)); + return []; + }); + }); + const proxy = await startProxy(base, 50); + + const res = await proxy.inject({ method: 'GET', url: '/api/v1/servers' }); + const body = res.json(); + expect(body.message).toMatch(/did not respond within/); + expect(body.message).not.toMatch(/Cannot reach mcpd/); + expect(body.message).toMatch(/reachable/); + }); + + it('streams SSE through instead of buffering it', async () => { + const base = await startUpstream((a) => { + a.post('/api/v1/agents/:name/chat', async (_req, reply) => { + reply.raw.writeHead(200, { + 'Content-Type': 'text/event-stream', + 'Cache-Control': 'no-cache', + 'X-Accel-Buffering': 'no', + }); + reply.raw.write('data: {"type":"text","delta":"po"}\n\n'); + reply.raw.write('data: {"type":"text","delta":"ng"}\n\n'); + reply.raw.write('data: [DONE]\n\n'); + reply.raw.end(); + return reply; + }); + }); + const proxy = await startProxy(base, 50); + + const res = await proxy.inject({ + method: 'POST', + url: '/api/v1/agents/reviewer/chat', + payload: { message: 'hi', stream: true }, + }); + + expect(res.statusCode).toBe(200); + // Content-type must survive — a client that gets application/json will not + // parse the event stream. + expect(res.headers['content-type']).toMatch(/text\/event-stream/); + // x-accel-buffering=no must survive too, or intermediaries re-buffer the + // stream and reintroduce the stall. + expect(res.headers['x-accel-buffering']).toBe('no'); + expect(res.body).toContain('"delta":"po"'); + expect(res.body).toContain('"delta":"ng"'); + expect(res.body).toContain('[DONE]'); + }); + + it('delivers each SSE frame while the upstream is still generating', async () => { + // The buffering regression is invisible to the pass-through test above: + // `inject()` collects the whole body, so a proxy that buffers via + // res.text() still passes it. This test proves *progressive* delivery by + // making the upstream withhold its final frame until the client has + // observed the first one. A buffering proxy can never satisfy that + // ordering — the 3s guard resolves the gate so the run fails cleanly + // instead of deadlocking. + let openGate: (seen: boolean) => void = () => {}; + const clientSawFirstFrame = new Promise((r) => { openGate = r; }); + const guard = setTimeout(() => openGate(false), 3_000); + + const base = await startUpstream((a) => { + a.post('/api/v1/agents/:name/chat', async (_req, reply) => { + reply.raw.writeHead(200, { 'Content-Type': 'text/event-stream' }); + reply.raw.write('data: {"type":"text","delta":"live"}\n\n'); + await clientSawFirstFrame; + reply.raw.write('data: {"type":"final"}\n\n'); + reply.raw.write('data: [DONE]\n\n'); + reply.raw.end(); + return reply; + }); + }); + const proxy = await startProxy(base, 50); + await proxy.listen({ port: 0, host: '127.0.0.1' }); + const addr = proxy.server.address(); + if (addr === null || typeof addr === 'string') throw new Error('no address'); + + const body = await new Promise((resolve, reject) => { + const req = http.request({ + hostname: '127.0.0.1', + port: addr.port, + path: '/api/v1/agents/reviewer/chat', + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + }, (res) => { + let acc = ''; + res.setEncoding('utf-8'); + res.on('data', (chunk: string) => { + acc += chunk; + if (acc.includes('"delta":"live"')) openGate(true); + }); + res.on('end', () => resolve(acc)); + res.on('error', reject); + }); + req.on('error', reject); + req.end(JSON.stringify({ message: 'hi', stream: true })); + }); + clearTimeout(guard); + + // The ordering proof: the first frame reached the client while the + // upstream was still holding the stream open. + await expect(clientSawFirstFrame).resolves.toBe(true); + expect(body).toContain('"type":"final"'); + expect(body).toContain('[DONE]'); + }); + + it('relays a non-200 status from a streaming route', async () => { + const base = await startUpstream((a) => { + a.post('/api/v1/agents/:name/chat', async (_req, reply) => { + return reply.code(404).send({ error: 'Agent not found' }); + }); + }); + const proxy = await startProxy(base, 50); + + const res = await proxy.inject({ + method: 'POST', + url: '/api/v1/agents/ghost/chat', + payload: { message: 'hi' }, + }); + expect(res.statusCode).toBe(404); + expect(res.body).toContain('Agent not found'); + }); + + it('still reports a genuinely unreachable daemon as 503', async () => { + // Port 1 is reserved and refuses instantly. + const proxy = await startProxy('http://127.0.0.1:1', 500); + const res = await proxy.inject({ method: 'GET', url: '/api/v1/servers' }); + + expect(res.statusCode).toBe(503); + expect(res.json().error).toBe('service_unavailable'); + }); + + it('propagates 401 from a streaming route so login guidance still fires', async () => { + const base = await startUpstream((a) => { + a.post('/api/v1/agents/:name/chat', async (_req, reply) => reply.code(401).send({})); + }); + const proxy = await startProxy(base, 50); + + const res = await proxy.inject({ + method: 'POST', + url: '/api/v1/agents/reviewer/chat', + payload: { message: 'hi' }, + }); + expect(res.statusCode).toBe(401); + expect(res.json().message).toMatch(/mcpctl login/); + }); +}); + +describe('error taxonomy', () => { + it('keeps timeout and unreachable as distinct types', () => { + const timeout = new UpstreamTimeoutError('http://mcpd', 30_000); + expect(timeout).not.toBeInstanceOf(ConnectionError); + expect(timeout.timeoutMs).toBe(30_000); + expect(timeout.message).toMatch(/did not respond within 30000ms/); + }); +}); diff --git a/src/mcplocal/tests/smoke/agent-chat.smoke.test.ts b/src/mcplocal/tests/smoke/agent-chat.smoke.test.ts index e445511..3261c43 100644 --- a/src/mcplocal/tests/smoke/agent-chat.smoke.test.ts +++ b/src/mcplocal/tests/smoke/agent-chat.smoke.test.ts @@ -18,8 +18,12 @@ import { describe, it, expect, beforeAll, afterAll } from 'vitest'; import http from 'node:http'; import https from 'node:https'; import { spawnSync, execSync } from 'node:child_process'; +import { existsSync, readFileSync } from 'node:fs'; +import { join } from 'node:path'; +import { homedir } from 'node:os'; const MCPD_URL = process.env.MCPD_URL ?? 'https://mcpctl.ad.itaz.eu'; +const MCPLOCAL_URL = process.env.MCPLOCAL_URL ?? 'http://localhost:3200'; const LLM_URL = process.env.MCPCTL_SMOKE_LLM_URL; const LLM_MODEL = process.env.MCPCTL_SMOKE_LLM_MODEL ?? 'qwen3-thinking'; const LLM_KEY = process.env.MCPCTL_SMOKE_LLM_KEY; @@ -27,6 +31,10 @@ const SUFFIX = Date.now().toString(36); const SECRET_NAME = `smoke-chat-sec-${SUFFIX}`; const LLM_NAME = `smoke-chat-llm-${SUFFIX}`; const AGENT_NAME = `smoke-chat-agent-${SUFFIX}`; +// Dedicated agent for the streaming-timing test: the shared agent's system +// prompt pins the reply to a single token, which is too short to distinguish +// live streaming from an end-of-turn buffer dump. +const STREAM_AGENT_NAME = `smoke-stream-agent-${SUFFIX}`; interface CliResult { code: number; stdout: string; stderr: string } @@ -99,6 +107,7 @@ describe('agent chat smoke (live LLM)', () => { afterAll(() => { if (!liveLlmConfigured || !mcpdUp) return; run(`delete agent ${AGENT_NAME}`); + run(`delete agent ${STREAM_AGENT_NAME}`); run(`delete llm ${LLM_NAME}`); run(`delete secret ${SECRET_NAME}`); }); @@ -139,6 +148,92 @@ describe('agent chat smoke (live LLM)', () => { expect(result.stderr).toMatch(/thread:\s+c[a-z0-9]+/); }); + it('streams progressively THROUGH mcplocal — frames arrive during generation, not in one burst', async () => { + if (!liveLlmConfigured || !mcpdUp) return; + // The regression this pins: mcplocal's /api/v1/* proxy buffered SSE via + // res.text(), so the CLI showed nothing until the turn finished and then + // dumped the whole answer at once. The --direct tests above bypass + // mcplocal entirely and cannot catch that. This one posts to the local + // proxy (the path `mcpctl chat` actually takes) and asserts frames are + // spread across the generation window: with buffering, everything lands + // within a few ms of stream end. + if (!(await healthz(MCPLOCAL_URL))) { + // eslint-disable-next-line no-console + console.warn(`\n ○ mcplocal streaming smoke: skipped — ${MCPLOCAL_URL}/healthz unreachable.\n`); + return; + } + let token = ''; + try { + const credsPath = join(homedir(), '.mcpctl', 'credentials'); + if (existsSync(credsPath)) { + const creds = JSON.parse(readFileSync(credsPath, 'utf-8')) as { token?: string }; + if (creds.token !== undefined) token = creds.token; + } + } catch { /* unauthenticated — the request will 401 and fail loudly */ } + + run(`delete agent ${STREAM_AGENT_NAME}`); + const agent = run([ + `create agent ${STREAM_AGENT_NAME}`, + `--llm ${LLM_NAME}`, + `--description "mcplocal streaming smoke"`, + `--system-prompt "You are a smoke test. Follow the user's instructions exactly."`, + '--default-temperature 0', + '--default-max-tokens 512', + ].join(' ')); + expect(agent.code, agent.stderr).toBe(0); + + const url = new URL(`${MCPLOCAL_URL.replace(/\/$/, '')}/api/v1/agents/${STREAM_AGENT_NAME}/chat`); + const deltaTimes: number[] = []; + let endTime = 0; + let status = 0; + let raw = ''; + + await new Promise((resolve, reject) => { + const req = http.request({ + hostname: url.hostname, + port: url.port || 80, + path: url.pathname, + method: 'POST', + timeout: 120_000, + headers: { + 'Content-Type': 'application/json', + ...(token !== '' ? { Authorization: `Bearer ${token}` } : {}), + }, + }, (res) => { + status = res.statusCode ?? 0; + res.setEncoding('utf-8'); + let buf = ''; + res.on('data', (chunk: string) => { + raw += chunk; + buf += chunk; + let nl: number; + while ((nl = buf.indexOf('\n\n')) !== -1) { + const frame = buf.slice(0, nl); + buf = buf.slice(nl + 2); + if (/"type":"(text|thinking)"/.test(frame)) deltaTimes.push(Date.now()); + } + }); + res.on('end', () => { endTime = Date.now(); resolve(); }); + res.on('error', reject); + }); + req.on('error', reject); + req.on('timeout', () => { req.destroy(); reject(new Error('stream timed out')); }); + req.end(JSON.stringify({ + message: 'Count from 1 to 40, one number per line. No other text.', + stream: true, + max_tokens: 400, + })); + }); + + expect(status, raw.slice(0, 500)).toBe(200); + expect(deltaTimes.length).toBeGreaterThanOrEqual(2); + // The buffering signature: every frame lands in the same final burst as + // stream end. Live streaming puts the first delta well before the end — + // a 40-line generation spans seconds; 300ms is a conservative floor. + const firstDelta = deltaTimes[0]!; + expect(endTime - firstDelta).toBeGreaterThanOrEqual(300); + }, 150_000); + it('streaming `mcpctl chat` emits text deltas', () => { if (!liveLlmConfigured || !mcpdUp) return; // Default mode is streaming. Pipe stdout/stderr separately.