-
-
Notifications
You must be signed in to change notification settings - Fork 577
fix(mcp): harden HTTP connection lifecycle #371
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -2,6 +2,7 @@ import assert from "node:assert/strict"; | |
| import { execFile } from "node:child_process"; | ||
| import { createHash } from "node:crypto"; | ||
| import { access, mkdtemp, mkdir, readFile, rm, symlink, writeFile } from "node:fs/promises"; | ||
| import { createServer as createHttpServer } from "node:http"; | ||
| import { platform, tmpdir } from "node:os"; | ||
| import { join } from "node:path"; | ||
| import test, { type TestContext } from "node:test"; | ||
|
|
@@ -14,7 +15,14 @@ import { buildLocalAgentProviderStatuses } from "./local-agent-catalog.js"; | |
| import type { SubagentsConfig } from "./local-agent-config.js"; | ||
| import { createReviewCheckpointManager } from "./review-checkpoints.js"; | ||
| import { ProcessSessionManager } from "./process-sessions.js"; | ||
| import { createMcpServer, createServer } from "./server.js"; | ||
| import { shutdownHttpServer } from "./server-shutdown.js"; | ||
| import { | ||
| configureHttpServer, | ||
| createMcpServer, | ||
| createServer, | ||
| DEVSPACE_HTTP_HEADERS_TIMEOUT_MS, | ||
| DEVSPACE_HTTP_KEEP_ALIVE_TIMEOUT_MS, | ||
| } from "./server.js"; | ||
| import { SqliteWorkspaceStore } from "./workspace-store.js"; | ||
| import { WorkspaceRegistry } from "./workspaces.js"; | ||
| import { writeTestDevspaceConfig } from "./test-support/config.test.js"; | ||
|
|
@@ -470,6 +478,14 @@ test("open_workspace scopes checkout reuse to OpenAI session metadata", async (t | |
| assert.ok(Array.isArray(structuredContent(unscoped).agents_files)); | ||
| }); | ||
|
|
||
| test("HTTP listener advertises a tunnel-safe keep-alive lifetime", () => { | ||
| const httpServer = createHttpServer(); | ||
| configureHttpServer(httpServer); | ||
|
|
||
| assert.equal(httpServer.keepAliveTimeout, DEVSPACE_HTTP_KEEP_ALIVE_TIMEOUT_MS); | ||
| assert.equal(httpServer.headersTimeout, DEVSPACE_HTTP_HEADERS_TIMEOUT_MS); | ||
| }); | ||
|
|
||
| test("HTTP endpoint serves modern MCP and stateless legacy clients", async (t) => { | ||
| const { root, localBaseUrl, accessToken } = await httpServerFixture( | ||
| t, | ||
|
|
@@ -503,6 +519,8 @@ test("HTTP endpoint serves modern MCP and stateless legacy clients", async (t) = | |
| {}, | ||
| ); | ||
| assert.equal(listed.status, 200, await listed.clone().text()); | ||
| assert.equal(listed.headers.get("connection"), "keep-alive"); | ||
| assert.match(listed.headers.get("keep-alive") ?? "", /timeout=300/); | ||
| const listBody = await listed.json() as { | ||
| result?: { tools?: Array<{ name?: string }> }; | ||
| }; | ||
|
|
@@ -583,6 +601,52 @@ test("HTTP endpoint serves modern MCP and stateless legacy clients", async (t) = | |
| assert.match(await legacyTools.text(), /"open_workspace"/); | ||
| }); | ||
|
|
||
| test("aborting one MCP request does not poison the next stateless request", async (t) => { | ||
| const { root, localBaseUrl, accessToken } = await httpServerFixture( | ||
| t, | ||
| "devspace-aborted-http-test-", | ||
| ); | ||
| const opened = await postModernMcp( | ||
| localBaseUrl, | ||
| accessToken, | ||
| "tools/call", | ||
| { | ||
| name: "open_workspace", | ||
| arguments: { path: root }, | ||
| _meta: { "openai/session": "aborted-http-test" }, | ||
| }, | ||
| ); | ||
| const openBody = await opened.json() as { | ||
| result?: { structuredContent?: { workspace_id?: string } }; | ||
| }; | ||
| const workspaceId = openBody.result?.structuredContent?.workspace_id; | ||
| assert.equal(typeof workspaceId, "string"); | ||
|
|
||
| const controller = new AbortController(); | ||
| const toolCall = postModernMcp( | ||
| localBaseUrl, | ||
| accessToken, | ||
| "tools/call", | ||
| { | ||
| name: "exec_command", | ||
| arguments: { | ||
| workspace_id: workspaceId, | ||
| cmd: "node -e \"const fs=require('node:fs');fs.writeFileSync('started','');setTimeout(()=>fs.writeFileSync('finished',''),500)\"", | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win Keep the command active until the abort occurs. If the test does not observe 🤖 Prompt for AI Agents |
||
| yield_time_ms: 1_000, | ||
| }, | ||
| }, | ||
| { signal: controller.signal }, | ||
| ); | ||
| await waitForFile(join(root, "started")); | ||
| controller.abort(); | ||
| await assert.rejects(toolCall, /abort/i); | ||
|
|
||
| const listed = await postModernMcp(localBaseUrl, accessToken, "tools/list", {}); | ||
| assert.equal(listed.status, 200, await listed.clone().text()); | ||
| await listed.text(); | ||
| await waitForFile(join(root, "finished")); | ||
| }); | ||
|
|
||
| test("server shutdown waits for an active MCP tool call", async (t) => { | ||
| const { root, localBaseUrl, accessToken, running } = await httpServerFixture( | ||
| t, | ||
|
|
@@ -694,12 +758,10 @@ async function httpServerFixture( | |
| const running = createServer(config, { incomingArtifactAdapters: [] }); | ||
| const httpServer = running.app.listen(0, "127.0.0.1"); | ||
| await new Promise<void>((resolve) => httpServer.once("listening", resolve)); | ||
| configureHttpServer(httpServer); | ||
|
|
||
| t.after(async () => { | ||
| await new Promise<void>((resolve, reject) => { | ||
| httpServer.close((error) => error ? reject(error) : resolve()); | ||
| }); | ||
| await running.close(); | ||
| await shutdownHttpServer(httpServer, running.close); | ||
| await rm(root, { recursive: true, force: true }); | ||
| }); | ||
|
|
||
|
|
@@ -909,6 +971,7 @@ function postModernMcp( | |
| accessToken: string | undefined, | ||
| method: string, | ||
| params: Record<string, unknown>, | ||
| options: { signal?: AbortSignal } = {}, | ||
| ): Promise<Response> { | ||
| const mcpName = typeof params.name === "string" | ||
| ? params.name | ||
|
|
@@ -924,6 +987,7 @@ function postModernMcp( | |
| "mcp-protocol-version": "2026-07-28", | ||
| ...(mcpName ? { "mcp-name": mcpName } : {}), | ||
| }, | ||
| signal: options.signal, | ||
| body: JSON.stringify({ | ||
| jsonrpc: "2.0", | ||
| id: `modern-${method}`, | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,6 +1,7 @@ | ||
| import { randomUUID } from "node:crypto"; | ||
| import { readFileSync } from "node:fs"; | ||
| import { access, realpath } from "node:fs/promises"; | ||
| import type { Server as HttpServer } from "node:http"; | ||
| import { fileURLToPath } from "node:url"; | ||
| import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; | ||
| import { createMcpExpressApp } from "@modelcontextprotocol/sdk/server/express.js"; | ||
|
|
@@ -92,6 +93,16 @@ interface RunningServer { | |
| close(): Promise<void>; | ||
| } | ||
|
|
||
| // Keep the local origin alive longer than the reverse proxy's pooled connection. | ||
| // Node must also advertise this timeout itself, so MCP responses drop hop-by-hop headers below. | ||
| export const DEVSPACE_HTTP_KEEP_ALIVE_TIMEOUT_MS = 5 * 60 * 1_000; | ||
| export const DEVSPACE_HTTP_HEADERS_TIMEOUT_MS = DEVSPACE_HTTP_KEEP_ALIVE_TIMEOUT_MS + 5_000; | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
If untrusted clients can reach the Node listener directly, they can leave request headers incomplete for up to 305 seconds before the server rejects the connection, compared with the previous 60-second deadline. This increases resource-exhaustion exposure. Keep the header deadline separate from the desired keep-alive lifetime, or enforce a shorter deadline at the ingress. How this was verified: The configured deadline increased, and an incomplete-header connection stayed open longer under a controlled shorter deadline. ArtifactsSource for the local HTTP and incomplete-header TCP probe
Base revision HTTP and TCP probe output
PR-head HTTP and TCP probe output
|
||
|
|
||
| export function configureHttpServer(httpServer: HttpServer): void { | ||
| httpServer.keepAliveTimeout = DEVSPACE_HTTP_KEEP_ALIVE_TIMEOUT_MS; | ||
| httpServer.headersTimeout = DEVSPACE_HTTP_HEADERS_TIMEOUT_MS; | ||
| } | ||
|
|
||
| type TrackToolActivity = <T>(operation: () => Promise<T>) => Promise<T>; | ||
|
|
||
| class ToolActivityTracker { | ||
|
|
@@ -218,6 +229,32 @@ function requestLogFields(req: Request, config: ServerConfig): Record<string, un | |
| }; | ||
| } | ||
|
|
||
| function rpcRequestLogFields(body: unknown): Record<string, unknown> { | ||
| if (!body || typeof body !== "object" || Array.isArray(body)) return {}; | ||
| const request = body as { id?: unknown; method?: unknown }; | ||
| return { | ||
| rpcId: request.id, | ||
| rpcMethod: request.method, | ||
| }; | ||
| } | ||
|
|
||
| function nodeMcpResponse(response: globalThis.Response): globalThis.Response { | ||
| const headers = new Headers(response.headers); | ||
|
|
||
| // Connection is hop-by-hop state owned by Node's HTTP server. Preserving an | ||
| // application-supplied value prevents Node from advertising its real socket | ||
| // lifetime and can make a proxy reuse a connection while the server closes it. | ||
| headers.delete("connection"); | ||
| headers.delete("keep-alive"); | ||
| headers.delete("transfer-encoding"); | ||
|
|
||
| return new globalThis.Response(response.body, { | ||
| status: response.status, | ||
| statusText: response.statusText, | ||
| headers, | ||
| }); | ||
| } | ||
|
|
||
| function assetBaseUrl(config: ServerConfig): string { | ||
| return `${config.publicBaseUrl.replace(/\/+$/, "")}/mcp-app-assets`; | ||
| } | ||
|
|
@@ -857,7 +894,9 @@ export function createServer( | |
| legacy: "stateless", | ||
| onerror: logMcpHandlerError, | ||
| }); | ||
| const mcpNodeHandler = toNodeHandler(mcpHandler, { | ||
| const mcpNodeHandler = toNodeHandler({ | ||
| fetch: async (request, options) => nodeMcpResponse(await mcpHandler.fetch(request, options)), | ||
| }, { | ||
| onerror: logMcpHandlerError, | ||
| }); | ||
|
|
||
|
|
@@ -868,12 +907,25 @@ export function createServer( | |
| app.use((req, res, next) => { | ||
| const requestId = randomUUID(); | ||
| const startedAt = performance.now(); | ||
| const path = requestPath(req); | ||
| const shouldLogRequest = config.logging.requests | ||
| && (config.logging.assets || !path.startsWith("/mcp-app-assets")); | ||
| let finished = false; | ||
| res.locals.requestId = requestId; | ||
|
|
||
| if (shouldLogRequest) { | ||
| logEvent(config.logging, "debug", "http_request_start", { | ||
| requestId, | ||
| method: req.method, | ||
| path, | ||
| ...requestLogFields(req, config), | ||
| ...(path === "/mcp" ? rpcRequestLogFields(req.body) : {}), | ||
| }); | ||
| } | ||
|
|
||
| res.on("finish", () => { | ||
| const path = requestPath(req); | ||
| if (!config.logging.requests) return; | ||
| if (!config.logging.assets && path.startsWith("/mcp-app-assets")) return; | ||
| finished = true; | ||
| if (!shouldLogRequest) return; | ||
|
|
||
| logEvent(config.logging, "info", "http_request", { | ||
| requestId, | ||
|
|
@@ -882,6 +934,21 @@ export function createServer( | |
| status: res.statusCode, | ||
| durationMs: Math.round(performance.now() - startedAt), | ||
| ...requestLogFields(req, config), | ||
| ...(path === "/mcp" ? rpcRequestLogFields(req.body) : {}), | ||
| }); | ||
| }); | ||
|
|
||
| res.on("close", () => { | ||
| if (finished || !shouldLogRequest) return; | ||
| logEvent(config.logging, "warn", "http_request_aborted", { | ||
| requestId, | ||
| method: req.method, | ||
| path, | ||
| headersSent: res.headersSent, | ||
| ...(res.headersSent ? { status: res.statusCode } : {}), | ||
| durationMs: Math.round(performance.now() - startedAt), | ||
| ...requestLogFields(req, config), | ||
| ...(path === "/mcp" ? rpcRequestLogFields(req.body) : {}), | ||
| }); | ||
| }); | ||
|
|
||
|
|
@@ -944,6 +1011,10 @@ export function createServer( | |
| logEvent(config.logging, "debug", "mcp_request", { | ||
| requestId, | ||
| method: req.method, | ||
| protocolVersion: req.header("mcp-protocol-version"), | ||
| mcpMethod: req.header("mcp-method"), | ||
| mcpName: req.header("mcp-name"), | ||
| ...rpcRequestLogFields(req.body), | ||
| }); | ||
|
|
||
| try { | ||
|
|
@@ -1011,6 +1082,7 @@ if (await isMainModule()) { | |
| console.log(`native artifact download: ${artifactDownloadStatus}`); | ||
| console.log(`subagent providers: ${formatLocalAgentProviderStatusSummary(localAgentProviders)}`); | ||
| }); | ||
| configureHttpServer(httpServer); | ||
|
|
||
| let shuttingDown = false; | ||
| const shutdown = async () => { | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
If shutdown begins during an
/mcp-app-assetsdownload, application cleanup does not wait for the response to finish. This call destroys the active connection, leaving the client with a truncated asset. The client may need to retry the download before the workspace app can load.Artifacts
Authored HTTP asset-download and shutdown repro
Download with force-close omitted
Download with actual force-close behavior