fix(mcplocal): stream chat SSE through the proxy instead of buffering it #109
@@ -2,7 +2,7 @@ export { createHttpServer } from './server.js';
|
|||||||
export type { HttpServerDeps } from './server.js';
|
export type { HttpServerDeps } from './server.js';
|
||||||
export { loadHttpConfig } from './config.js';
|
export { loadHttpConfig } from './config.js';
|
||||||
export type { HttpConfig } 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 { registerProxyRoutes } from './routes/proxy.js';
|
||||||
export { registerMcpEndpoint } from './mcp-endpoint.js';
|
export { registerMcpEndpoint } from './mcp-endpoint.js';
|
||||||
export { registerProjectMcpEndpoint } from './project-mcp-endpoint.js';
|
export { registerProjectMcpEndpoint } from './project-mcp-endpoint.js';
|
||||||
|
|||||||
@@ -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. */
|
/** Default timeout for mcpd requests (ms). Prevents indefinite hangs on slow upstream tool calls. */
|
||||||
export const DEFAULT_TIMEOUT_MS = 30_000;
|
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
|
* 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
|
* the full tool-call timeout budget — a single dead upstream would stall session init for
|
||||||
@@ -121,9 +153,7 @@ export class McpdClient {
|
|||||||
try {
|
try {
|
||||||
res = await fetch(url, init);
|
res = await fetch(url, init);
|
||||||
} catch (err: unknown) {
|
} catch (err: unknown) {
|
||||||
if (err instanceof DOMException && err.name === 'TimeoutError') {
|
if (isTimeout(err)) throw new UpstreamTimeoutError(this.baseUrl, this.timeoutMs);
|
||||||
throw new ConnectionError(this.baseUrl, new Error(`Request timed out after ${this.timeoutMs}ms`));
|
|
||||||
}
|
|
||||||
throw new ConnectionError(this.baseUrl, err);
|
throw new ConnectionError(this.baseUrl, err);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -131,7 +161,18 @@ export class McpdClient {
|
|||||||
throw new AuthenticationError();
|
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;
|
let parsed: unknown;
|
||||||
try {
|
try {
|
||||||
parsed = JSON.parse(text);
|
parsed = JSON.parse(text);
|
||||||
@@ -142,6 +183,51 @@ export class McpdClient {
|
|||||||
return { status: res.status, body: parsed };
|
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<Response> {
|
||||||
|
const url = `${this.baseUrl}${path}${query ? `?${query}` : ''}`;
|
||||||
|
const headers: Record<string, string> = {
|
||||||
|
...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<T>(method: string, path: string, body?: unknown): Promise<T> {
|
private async request<T>(method: string, path: string, body?: unknown): Promise<T> {
|
||||||
const result = await this.forward(method, path, '', body);
|
const result = await this.forward(method, path, '', body);
|
||||||
|
|
||||||
|
|||||||
@@ -1,10 +1,62 @@
|
|||||||
/**
|
/**
|
||||||
* Catch-all proxy route that forwards /api/v1/* requests to mcpd.
|
* Catch-all proxy route that forwards /api/v1/* requests to mcpd.
|
||||||
*/
|
*/
|
||||||
import type { FastifyInstance } from 'fastify';
|
import { Readable } from 'node:stream';
|
||||||
import { AuthenticationError, ConnectionError } from '../mcpd-client.js';
|
|
||||||
|
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';
|
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 {
|
export function registerProxyRoutes(app: FastifyInstance, client: McpdClient): void {
|
||||||
app.all('/api/v1/*', async (request, reply) => {
|
app.all('/api/v1/*', async (request, reply) => {
|
||||||
const path = (request.url.split('?')[0]) ?? '/';
|
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.
|
// 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.
|
// If no user token is present, mcpd will use its auth hook to reject.
|
||||||
const authHeader = request.headers['authorization'] as string | undefined;
|
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 {
|
try {
|
||||||
const result = await client.forward(request.method, path, querystring, body, userToken);
|
const result = await client.forward(request.method, path, querystring, body, userToken);
|
||||||
return reply.code(result.status).send(result.body);
|
return reply.code(result.status).send(result.body);
|
||||||
} catch (err: unknown) {
|
} catch (err: unknown) {
|
||||||
if (err instanceof AuthenticationError) {
|
const handled = sendUpstreamError(reply, err);
|
||||||
return reply.code(401).send({
|
if (handled) return handled;
|
||||||
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?',
|
|
||||||
});
|
|
||||||
}
|
|
||||||
throw err;
|
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<void> {
|
||||||
|
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<string, string> = {};
|
||||||
|
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<void>((resolve, reject) => {
|
||||||
|
const upstream = Readable.fromWeb(res.body as Parameters<typeof Readable.fromWeb>[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();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -11,7 +11,7 @@ export type { MainResult } from './main.js';
|
|||||||
export { ProviderRegistry } from './providers/index.js';
|
export { ProviderRegistry } from './providers/index.js';
|
||||||
export type { LlmProvider, CompletionOptions, CompletionResult, ChatMessage } 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 { 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 { HttpConfig, HttpServerDeps } from './http/index.js';
|
||||||
export type {
|
export type {
|
||||||
JsonRpcRequest,
|
JsonRpcRequest,
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
import { describe, it, expect, afterAll, afterEach } from 'vitest';
|
import { describe, it, expect, afterAll, afterEach } from 'vitest';
|
||||||
import http from 'node:http';
|
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.
|
* Create a local HTTP server for testing McpdClient behavior.
|
||||||
@@ -85,7 +85,7 @@ describe('McpdClient', () => {
|
|||||||
|
|
||||||
// ── Timeout behavior ──
|
// ── 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) => {
|
const { server, url } = await createTestServer((_req, _res) => {
|
||||||
// Never respond — simulates a hanging upstream tool call
|
// Never respond — simulates a hanging upstream tool call
|
||||||
});
|
});
|
||||||
@@ -96,7 +96,7 @@ describe('McpdClient', () => {
|
|||||||
|
|
||||||
const start = Date.now();
|
const start = Date.now();
|
||||||
await expect(client.post('/api/v1/mcp/proxy', { serverId: 's1' })).rejects.toThrow(
|
await expect(client.post('/api/v1/mcp/proxy', { serverId: 's1' })).rejects.toThrow(
|
||||||
/timed out/,
|
/did not respond within/,
|
||||||
);
|
);
|
||||||
const elapsed = Date.now() - start;
|
const elapsed = Date.now() - start;
|
||||||
|
|
||||||
@@ -105,7 +105,7 @@ describe('McpdClient', () => {
|
|||||||
expect(elapsed).toBeLessThan(3000);
|
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) => {
|
const { server, url } = await createTestServer((_req, _res) => {
|
||||||
// Never respond
|
// Never respond
|
||||||
});
|
});
|
||||||
@@ -117,8 +117,12 @@ describe('McpdClient', () => {
|
|||||||
await client.get('/test');
|
await client.get('/test');
|
||||||
expect.unreachable('Should have thrown');
|
expect.unreachable('Should have thrown');
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
expect(err).toBeInstanceOf(ConnectionError);
|
// Reporting a timeout as "cannot connect" is what sent a previous
|
||||||
expect((err as Error).message).toContain('Request timed out after 200ms');
|
// 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 derived = client.withHeaders({ 'X-Custom': 'val' });
|
||||||
|
|
||||||
const start = Date.now();
|
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;
|
const elapsed = Date.now() - start;
|
||||||
expect(elapsed).toBeLessThan(2000);
|
expect(elapsed).toBeLessThan(2000);
|
||||||
});
|
});
|
||||||
|
|||||||
255
src/mcplocal/tests/proxy-long-running.test.ts
Normal file
255
src/mcplocal/tests/proxy-long-running.test.ts
Normal file
@@ -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<string> {
|
||||||
|
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<FastifyInstance> {
|
||||||
|
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<boolean>((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<string>((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/);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -18,8 +18,12 @@ import { describe, it, expect, beforeAll, afterAll } from 'vitest';
|
|||||||
import http from 'node:http';
|
import http from 'node:http';
|
||||||
import https from 'node:https';
|
import https from 'node:https';
|
||||||
import { spawnSync, execSync } from 'node:child_process';
|
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 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_URL = process.env.MCPCTL_SMOKE_LLM_URL;
|
||||||
const LLM_MODEL = process.env.MCPCTL_SMOKE_LLM_MODEL ?? 'qwen3-thinking';
|
const LLM_MODEL = process.env.MCPCTL_SMOKE_LLM_MODEL ?? 'qwen3-thinking';
|
||||||
const LLM_KEY = process.env.MCPCTL_SMOKE_LLM_KEY;
|
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 SECRET_NAME = `smoke-chat-sec-${SUFFIX}`;
|
||||||
const LLM_NAME = `smoke-chat-llm-${SUFFIX}`;
|
const LLM_NAME = `smoke-chat-llm-${SUFFIX}`;
|
||||||
const AGENT_NAME = `smoke-chat-agent-${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 }
|
interface CliResult { code: number; stdout: string; stderr: string }
|
||||||
|
|
||||||
@@ -99,6 +107,7 @@ describe('agent chat smoke (live LLM)', () => {
|
|||||||
afterAll(() => {
|
afterAll(() => {
|
||||||
if (!liveLlmConfigured || !mcpdUp) return;
|
if (!liveLlmConfigured || !mcpdUp) return;
|
||||||
run(`delete agent ${AGENT_NAME}`);
|
run(`delete agent ${AGENT_NAME}`);
|
||||||
|
run(`delete agent ${STREAM_AGENT_NAME}`);
|
||||||
run(`delete llm ${LLM_NAME}`);
|
run(`delete llm ${LLM_NAME}`);
|
||||||
run(`delete secret ${SECRET_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]+/);
|
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<void>((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', () => {
|
it('streaming `mcpctl chat` emits text deltas', () => {
|
||||||
if (!liveLlmConfigured || !mcpdUp) return;
|
if (!liveLlmConfigured || !mcpdUp) return;
|
||||||
// Default mode is streaming. Pipe stdout/stderr separately.
|
// Default mode is streaming. Pipe stdout/stderr separately.
|
||||||
|
|||||||
Reference in New Issue
Block a user