From dfff5c7e6b9e2589181114ef37cc0dfe2f38988f Mon Sep 17 00:00:00 2001 From: Jordan Ritter Date: Thu, 1 Oct 2026 21:26:04 -0700 Subject: [PATCH 1/4] Answer 404 to an unknown or expired Mcp-Session-Id on POST, GET and DELETE /mcp so clients re-initialize Classify each /mcp request by its session header, look sessions up by own key so prototype-named ids are unknown, let an initialize with a stale header start a new session, and echo the request id in the POST 404. --- src/__tests__/mcp-unknown-session-404.test.ts | 231 ++++++++++++++++++ src/server.ts | 130 +++++++++- 2 files changed, 349 insertions(+), 12 deletions(-) create mode 100644 src/__tests__/mcp-unknown-session-404.test.ts diff --git a/src/__tests__/mcp-unknown-session-404.test.ts b/src/__tests__/mcp-unknown-session-404.test.ts new file mode 100644 index 0000000..28cb73e --- /dev/null +++ b/src/__tests__/mcp-unknown-session-404.test.ts @@ -0,0 +1,231 @@ +/** + * Unknown or expired Mcp-Session-Id must answer 404 (JSON-RPC -32001) so + * clients re-initialize. These tests cover the helpers server.ts exports for + * that: the request classifier, the 404 writer, and the live-transport lookup. + */ +import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; +import type { Response } from "express"; + +vi.mock("../config.js", () => ({ + getConfig: vi.fn().mockReturnValue({ + port: 0, + databaseUrl: "pglite:///tmp/test", + openaiApiKey: "", + githubToken: "", + githubWebhookSecret: "", + nodeEnv: "test", + logLevel: "info", + cloneDir: "/tmp/test", + slackBotToken: "", + slackSigningSecret: "", + discordBotToken: "", + discordPublicKey: "", + notionToken: "", + mcpJwtSecret: "e".repeat(64), + p2pTelemetryUrl: undefined, + p2pTelemetryDisabled: false, + packageVersion: "test", + }), + getServerConfig: vi.fn().mockReturnValue({ + server: { name: "pathfinder-test", version: "0.0.0" }, + sources: [], + tools: [], + }), + getAnalyticsConfig: vi.fn().mockReturnValue(undefined), + hasSearchTools: vi.fn().mockReturnValue(false), + hasKnowledgeTools: vi.fn().mockReturnValue(false), + hasCollectTools: vi.fn().mockReturnValue(false), + hasBashSemanticSearch: vi.fn().mockReturnValue(false), +})); + +type Method = "POST" | "GET" | "DELETE"; +const METHODS: Method[] = ["POST", "GET", "DELETE"]; +const SID = "deadbeef-0000-4000-8000-000000000000"; + +describe("classifyMcpSessionRequest", () => { + it.each(METHODS)( + "routes %s with a known session id (initialize or not)", + async (method) => { + const { classifyMcpSessionRequest } = await import("../server.js"); + for (const isInitialize of [false, true]) { + expect( + classifyMcpSessionRequest({ + method, + sessionId: SID, + hasTransport: true, + isInitialize, + }), + ).toBe("route"); + } + }, + ); + + it("starts a new session for POST initialize with an unknown (stale) session id", async () => { + const { classifyMcpSessionRequest } = await import("../server.js"); + expect( + classifyMcpSessionRequest({ + method: "POST", + sessionId: SID, + hasTransport: false, + isInitialize: true, + }), + ).toBe("new-session"); + }); + + it.each(METHODS)( + "reports unknown-session for %s with an unknown id and no initialize", + async (method) => { + const { classifyMcpSessionRequest } = await import("../server.js"); + expect( + classifyMcpSessionRequest({ + method, + sessionId: SID, + hasTransport: false, + isInitialize: false, + }), + ).toBe("unknown-session"); + }, + ); + + it("starts a new session for POST initialize with no session id", async () => { + const { classifyMcpSessionRequest } = await import("../server.js"); + expect( + classifyMcpSessionRequest({ + method: "POST", + sessionId: undefined, + hasTransport: false, + isInitialize: true, + }), + ).toBe("new-session"); + }); + + it.each(METHODS)( + "reports no-session for %s with no session id and no initialize", + async (method) => { + const { classifyMcpSessionRequest } = await import("../server.js"); + expect( + classifyMcpSessionRequest({ + method, + sessionId: undefined, + hasTransport: false, + isInitialize: false, + }), + ).toBe("no-session"); + }, + ); +}); + +describe("writeUnknownSession404", () => { + const NOT_FOUND_BODY = { + jsonrpc: "2.0", + error: { code: -32001, message: "Session not found" }, + id: null, + }; + let warn: ReturnType; + let log: ReturnType; + let error: ReturnType; + /** The value Date.now() returns. Each test starts a new 404 window at T0. */ + const T0 = 1_800_000_000_000; + let clock = T0; + + function fakeRes() { + const json = vi.fn(); + const status = vi.fn().mockReturnValue({ json }); + return { res: { status } as unknown as Response, status, json }; + } + + /** + * The writer's options for a body-less unknown-session request. The test + * that no session id or ip bytes are logged drives the real route, in + * mcp-unknown-session-404-routes.test.ts. + */ + function unknownReq(method: Method) { + return { method }; + } + + /** Every string written to console.warn/log/error since the last reset. */ + function consoleOutput(): string[] { + return [warn, log, error].flatMap((spy) => + spy.mock.calls.map((args: unknown[]) => args.map(String).join(" ")), + ); + } + + beforeEach(async () => { + warn = vi.spyOn(console, "warn").mockImplementation(() => {}); + log = vi.spyOn(console, "log").mockImplementation(() => {}); + error = vi.spyOn(console, "error").mockImplementation(() => {}); + clock = T0; + vi.spyOn(Date, "now").mockImplementation(() => clock); + warn.mockClear(); + log.mockClear(); + error.mockClear(); + }); + afterEach(() => { + vi.restoreAllMocks(); + }); + + it("sends 404 with the exact JSON-RPC body and writes no per-request log line", async () => { + const { writeUnknownSession404 } = await import("../server.js"); + const { res, status, json } = fakeRes(); + writeUnknownSession404(res, unknownReq("GET")); + expect(status).toHaveBeenCalledWith(404); + expect(json).toHaveBeenCalledTimes(1); + expect(json).toHaveBeenCalledWith(NOT_FOUND_BODY); + expect(consoleOutput()).toEqual([]); + }); + + it.each([ + [{ jsonrpc: "2.0", id: 42, method: "tools/list" }, 42], + [{ jsonrpc: "2.0", id: "abc", method: "tools/list" }, "abc"], + [{ jsonrpc: "2.0", id: 0, method: "tools/list" }, 0], + [{ jsonrpc: "2.0", method: "notifications/initialized" }, null], + [{ jsonrpc: "2.0", id: null, method: "tools/list" }, null], + [{ jsonrpc: "2.0", id: { x: 1 }, method: "tools/list" }, null], + [[{ jsonrpc: "2.0", id: 7, method: "tools/list" }], null], + ["not json", null], + [undefined, null], + ])( + "POST body %j gets response id %j", + async (body: unknown, expectedId: string | number | null) => { + const { writeUnknownSession404 } = await import("../server.js"); + const { res, json } = fakeRes(); + writeUnknownSession404(res, { method: "POST", body }); + expect(json).toHaveBeenCalledWith({ ...NOT_FOUND_BODY, id: expectedId }); + }, + ); +}); + +describe("getLiveTransport", () => { + it("returns the transport registered under a session id", async () => { + const { getLiveTransport } = await import("../server.js"); + const transport = { handleRequest: vi.fn() }; + expect(getLiveTransport({ [SID]: transport }, SID)).toBe(transport); + }); + + it("returns undefined for a missing or empty session id", async () => { + const { getLiveTransport } = await import("../server.js"); + expect(getLiveTransport({}, undefined)).toBeUndefined(); + expect(getLiveTransport({}, "")).toBeUndefined(); + expect(getLiveTransport({}, SID)).toBeUndefined(); + }); + + it.each(["constructor", "__proto__", "toString", "hasOwnProperty"])( + "does not treat the inherited key %s as a live session", + async (sid) => { + const { getLiveTransport, classifyMcpSessionRequest } = + await import("../server.js"); + const transports: Record void }> = {}; + expect(getLiveTransport(transports, sid)).toBeUndefined(); + for (const method of METHODS) { + expect( + classifyMcpSessionRequest({ + method, + sessionId: sid, + hasTransport: !!getLiveTransport(transports, sid), + isInitialize: false, + }), + ).toBe("unknown-session"); + } + }, + ); +}); diff --git a/src/server.ts b/src/server.ts index 7db2a8c..9db3524 100644 --- a/src/server.ts +++ b/src/server.ts @@ -1530,6 +1530,72 @@ export async function handleExistingSessionRequest(opts: { opts.sessionLastActivity[opts.sid] = now(); } +/** + * Look up the transport for a client-sent Mcp-Session-Id. Own keys only: + * the session maps are plain objects, so a bare `map[sid]` resolves ids + * such as "constructor" or "__proto__" to inherited Object.prototype + * members. Those ids get undefined here, so the routes answer them as + * unknown sessions. + */ +export function getLiveTransport( + map: Record, + sessionId: string | undefined, +): T | undefined { + return sessionId && Object.hasOwn(map, sessionId) + ? map[sessionId] + : undefined; +} + +/** + * Classify an /mcp request by its Mcp-Session-Id header. Rules, in order + * (an empty header counts as absent): + * 1. header present and a transport exists -> "route" + * 2. POST initialize (header absent or unknown) -> "new-session" + * 3. header present but unknown -> "unknown-session" (the route answers 404) + * 4. otherwise -> "no-session" (400 for POST and DELETE, 405 for GET) + * + * POST routes a live session before it calls this, so it always passes + * hasTransport: false. Exported so tests can cover the decision table + * without Express. + */ +export function classifyMcpSessionRequest(opts: { + method: "POST" | "GET" | "DELETE"; + sessionId: string | undefined; + hasTransport: boolean; + isInitialize: boolean; +}): "route" | "new-session" | "unknown-session" | "no-session" { + if (opts.sessionId && opts.hasTransport) return "route"; + if (opts.method === "POST" && opts.isInitialize) return "new-session"; + if (opts.sessionId) return "unknown-session"; + return "no-session"; +} + +/** + * Write the 404 JSON-RPC "Session not found" response for an /mcp request + * that carried an unknown or expired Mcp-Session-Id. MCP Streamable HTTP, + * session management: a client that gets this 404 must start a new session + * with an initialize request that carries no session id. + */ +export function writeUnknownSession404( + res: Response, + opts: { method: "POST" | "GET" | "DELETE"; body?: unknown }, +): void { + // JSON-RPC 2.0: the response id echoes the request id, and is null only + // when it cannot be determined. Only a single message with a string or + // number id qualifies; GET/DELETE carry no body, so they stay null. + const bodyId = + typeof opts.body === "object" && opts.body !== null + ? (opts.body as { id?: unknown }).id + : undefined; + const id = + typeof bodyId === "string" || typeof bodyId === "number" ? bodyId : null; + res.status(404).json({ + jsonrpc: "2.0", + error: { code: -32001, message: "Session not found" }, + id, + }); +} + /** * Read and normalize the request-origin tag from the X-Pathfinder-Source * header. Captured ONCE at MCP-session init and closed over for the lifetime @@ -1611,7 +1677,8 @@ app.post("/mcp", bearerMiddleware, async (req: Request, res: Response) => { const ip = clientIp(req, isTrustingProxy()); // Existing session — route to its transport - if (sessionId && transports[sessionId]) { + const existingTransport = getLiveTransport(transports, sessionId); + if (sessionId && existingTransport) { const method = req.body?.method as string | undefined; if (method === "tools/call") { const params = req.body?.params as Record | undefined; @@ -1648,7 +1715,7 @@ app.post("/mcp", bearerMiddleware, async (req: Request, res: Response) => { // success-vs-throw contract unit-testable without spinning up Express. await handleExistingSessionRequest({ sid: sessionId, - transport: transports[sessionId], + transport: existingTransport, req, res, sessionLastActivity, @@ -1656,8 +1723,24 @@ app.post("/mcp", bearerMiddleware, async (req: Request, res: Response) => { return; } - // New session — must be an initialize request - if (!sessionId && isInitializeRequest(req.body)) { + const disposition = classifyMcpSessionRequest({ + method: "POST", + sessionId, + hasTransport: false, + isInitialize: isInitializeRequest(req.body), + }); + + // Unknown or expired session id on a non-initialize request: 404 (see + // writeUnknownSession404). + if (disposition === "unknown-session") { + writeUnknownSession404(res, { method: "POST", body: req.body }); + return; + } + + // New session — must be an initialize request. An initialize that carries + // a stale Mcp-Session-Id also lands here. That is safe: the SDK skips + // session validation for initialize, and nothing below reads the header. + if (disposition === "new-session") { // Global session cap — reject before doing any per-IP work. if (isAtGlobalCapacity(transports, sseTransports, MAX_SESSIONS)) { const total = getTotalSessionCount(transports, sseTransports); @@ -2003,9 +2086,10 @@ app.post("/mcp", bearerMiddleware, async (req: Request, res: Response) => { }); // SSE stream for server-initiated notifications. -// Returns 405 when no valid session — the SDK interprets this as -// "server doesn't offer SSE at GET" which is the expected no-auth path. -// Returning 400 instead would cause the SDK to throw and trigger auth flow. +// Returns 405 when the Mcp-Session-Id header is missing or empty — the SDK +// interprets this as "server doesn't offer SSE at GET" which is the expected +// no-auth path. Returning 400 instead would cause the SDK to throw and trigger +// auth flow. An unknown session id gets 404 (see writeUnknownSession404). // // Intentionally NOT wrapped in `bearerMiddleware`. An unauthenticated // client that probes GET /mcp without a Mcp-Session-Id expects a @@ -2020,8 +2104,19 @@ app.get("/mcp", async (req: Request, res: Response) => { // than escaping to Express's default error handler. try { const sessionId = req.headers["mcp-session-id"] as string | undefined; - if (sessionId && transports[sessionId]) { - await transports[sessionId].handleRequest(req, res); + const transport = getLiveTransport(transports, sessionId); + const disposition = classifyMcpSessionRequest({ + method: "GET", + sessionId, + hasTransport: !!transport, + isInitialize: false, + }); + if (disposition === "unknown-session") { + writeUnknownSession404(res, { method: "GET" }); + return; + } + if (disposition === "route" && transport) { + await transport.handleRequest(req, res); } else { res.status(405).json({ jsonrpc: "2.0", @@ -2048,12 +2143,23 @@ app.delete("/mcp", bearerMiddleware, async (req: Request, res: Response) => { // structured 500 than leak the throw. try { const sessionId = req.headers["mcp-session-id"] as string | undefined; - if (sessionId && transports[sessionId]) { - await transports[sessionId].handleRequest(req, res); + const transport = getLiveTransport(transports, sessionId); + const disposition = classifyMcpSessionRequest({ + method: "DELETE", + sessionId, + hasTransport: !!transport, + isInitialize: false, + }); + if (disposition === "unknown-session") { + writeUnknownSession404(res, { method: "DELETE" }); + return; + } + if (disposition === "route" && transport) { + await transport.handleRequest(req, res); } else { res.status(400).json({ jsonrpc: "2.0", - error: { code: -32000, message: "Invalid or missing session ID" }, + error: { code: -32000, message: "Missing session ID" }, id: null, }); } From 22b77e7821e5d0c426714e897fa20a223d042477 Mon Sep 17 00:00:00 2001 From: Jordan Ritter Date: Thu, 1 Oct 2026 21:26:07 -0700 Subject: [PATCH 2/4] Count unknown-session 404s per method and log the totals on the reaper tick and at shutdown One line per flush carries window_s and the GET, POST and DELETE counts, and nothing client-supplied. --- src/__tests__/mcp-unknown-session-404.test.ts | 85 ++++++++++++++++++- src/server.ts | 73 +++++++++++++++- 2 files changed, 153 insertions(+), 5 deletions(-) diff --git a/src/__tests__/mcp-unknown-session-404.test.ts b/src/__tests__/mcp-unknown-session-404.test.ts index 28cb73e..18bf29c 100644 --- a/src/__tests__/mcp-unknown-session-404.test.ts +++ b/src/__tests__/mcp-unknown-session-404.test.ts @@ -1,7 +1,8 @@ /** * Unknown or expired Mcp-Session-Id must answer 404 (JSON-RPC -32001) so * clients re-initialize. These tests cover the helpers server.ts exports for - * that: the request classifier, the 404 writer, and the live-transport lookup. + * that: the request classifier, the 404 writer with its per-method counter + * flush, and the live-transport lookup. */ import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; import type { Response } from "express"; @@ -115,7 +116,7 @@ describe("classifyMcpSessionRequest", () => { ); }); -describe("writeUnknownSession404", () => { +describe("writeUnknownSession404 and flushUnknownSession404Counts", () => { const NOT_FOUND_BODY = { jsonrpc: "2.0", error: { code: -32001, message: "Session not found" }, @@ -156,6 +157,8 @@ describe("writeUnknownSession404", () => { error = vi.spyOn(console, "error").mockImplementation(() => {}); clock = T0; vi.spyOn(Date, "now").mockImplementation(() => clock); + const { flushUnknownSession404Counts } = await import("../server.js"); + flushUnknownSession404Counts(); warn.mockClear(); log.mockClear(); error.mockClear(); @@ -193,6 +196,84 @@ describe("writeUnknownSession404", () => { expect(json).toHaveBeenCalledWith({ ...NOT_FOUND_BODY, id: expectedId }); }, ); + + it("counts every unknown-session 404 in the flush window, with no per-id state", async () => { + const { writeUnknownSession404, flushUnknownSession404Counts } = + await import("../server.js"); + for (let i = 0; i < 60; i++) { + const { res, status } = fakeRes(); + writeUnknownSession404(res, unknownReq("GET")); + expect(status).toHaveBeenCalledWith(404); + } + const { res, status, json } = fakeRes(); + writeUnknownSession404(res, unknownReq("POST")); + expect(status).toHaveBeenCalledWith(404); + expect(json).toHaveBeenCalledWith(NOT_FOUND_BODY); + + flushUnknownSession404Counts(); + expect(consoleOutput()).toEqual([ + "[mcp] 404 unknown-session-id window_s=0 total=61 GET=60 POST=1 DELETE=0", + ]); + }); + + it("counts each method separately", async () => { + const { writeUnknownSession404, flushUnknownSession404Counts } = + await import("../server.js"); + const plan: Record = { GET: 2, POST: 3, DELETE: 4 }; + for (const method of METHODS) { + for (let i = 0; i < plan[method]; i++) { + writeUnknownSession404(fakeRes().res, unknownReq(method)); + } + } + flushUnknownSession404Counts(); + expect(consoleOutput()).toEqual([ + "[mcp] 404 unknown-session-id window_s=0 total=9 GET=2 POST=3 DELETE=4", + ]); + }); + + it("resets the counts on flush and emits nothing when the count is 0", async () => { + const { writeUnknownSession404, flushUnknownSession404Counts } = + await import("../server.js"); + flushUnknownSession404Counts(); + expect(consoleOutput()).toEqual([]); + + writeUnknownSession404(fakeRes().res, unknownReq("DELETE")); + flushUnknownSession404Counts(); + expect(consoleOutput()).toEqual([ + "[mcp] 404 unknown-session-id window_s=0 total=1 GET=0 POST=0 DELETE=1", + ]); + + flushUnknownSession404Counts(); + expect(consoleOutput()).toHaveLength(1); + + writeUnknownSession404(fakeRes().res, unknownReq("GET")); + flushUnknownSession404Counts(); + expect(consoleOutput()).toEqual([ + "[mcp] 404 unknown-session-id window_s=0 total=1 GET=0 POST=0 DELETE=1", + "[mcp] 404 unknown-session-id window_s=0 total=1 GET=1 POST=0 DELETE=0", + ]); + }); + + it("reports the whole seconds since the last reset, not the reaper period", async () => { + const { writeUnknownSession404, flushUnknownSession404Counts } = + await import("../server.js"); + // A partial window, as on shutdown() or stop(). + writeUnknownSession404(fakeRes().res, unknownReq("GET")); + clock = T0 + 137_400; + flushUnknownSession404Counts(); + + // An empty flush writes nothing but still starts a new window. + clock = T0 + 200_000; + flushUnknownSession404Counts(); + + writeUnknownSession404(fakeRes().res, unknownReq("POST")); + clock = T0 + 211_600; + flushUnknownSession404Counts(); + expect(consoleOutput()).toEqual([ + "[mcp] 404 unknown-session-id window_s=137 total=1 GET=1 POST=0 DELETE=0", + "[mcp] 404 unknown-session-id window_s=12 total=1 GET=0 POST=1 DELETE=0", + ]); + }); }); describe("getLiveTransport", () => { diff --git a/src/server.ts b/src/server.ts index 9db3524..cf10a1d 100644 --- a/src/server.ts +++ b/src/server.ts @@ -895,8 +895,12 @@ export function reapIdleSessionsTickForTesting(opts: { // Session reaper tick — started from startServer() so importing this module // (including from tests) doesn't leak a 5-minute setInterval into the loop. -// Thin closure over module state that delegates to the parametric form. +// Flushes the unknown-session 404 counts, warns at 80% session capacity, then +// runs the parametric form over module state. The parametric form does not +// flush. function reapIdleSessionsTick(): void { + flushUnknownSession404Counts(); + // 80% capacity warning if (MAX_SESSIONS !== undefined) { const total = getTotalSessionCount(transports, sseTransports); @@ -1570,9 +1574,66 @@ export function classifyMcpSessionRequest(opts: { return "no-session"; } +/** Period of the idle-session reaper tick, which also flushes the 404 counts. */ +export const SESSION_REAPER_INTERVAL_MS = 5 * 60 * 1000; + +/** + * Unknown-session 404s per method since the last reset (module load, boot, + * a flush, or the test seam). The count holds no client-supplied bytes, so + * the flushed line cannot carry them. + */ +const unknownSession404Counts: Record<"POST" | "GET" | "DELETE", number> = { + POST: 0, + GET: 0, + DELETE: 0, +}; + +/** + * Date.now() when unknownSession404Counts were last reset (module load, boot, + * a flush, or the test seam). The flushed line reports the time since then + * as its window. + */ +let unknownSession404WindowStartMs = Date.now(); + +/** Zero the unknown-session 404 counts and start a new window at `now`. */ +function resetUnknownSession404Counts(now: number): void { + unknownSession404Counts.GET = 0; + unknownSession404Counts.POST = 0; + unknownSession404Counts.DELETE = 0; + unknownSession404WindowStartMs = now; +} + +/** + * Write one `[mcp] 404 unknown-session-id window_s= ...` line with the + * per-method counts of unknown-session 404s since the last reset, then reset + * the counts. N is the elapsed time since that reset, rounded to the + * nearest second. It is about the reaper period for each tick, and shorter + * for the flushes on shutdown() and stop(). Writes nothing when the total is + * 0, but still starts a new window. Exported so tests can flush without + * waiting for the tick. + */ +export function flushUnknownSession404Counts(): void { + const now = Date.now(); + const windowS = Math.round((now - unknownSession404WindowStartMs) / 1000); + const { GET, POST, DELETE } = unknownSession404Counts; + const total = GET + POST + DELETE; + resetUnknownSession404Counts(now); + if (total > 0) { + console.warn( + `[mcp] 404 unknown-session-id window_s=${windowS} total=${total} GET=${GET} POST=${POST} DELETE=${DELETE}`, + ); + } +} + +/** @internal — test seam: zeroes the unknown-session 404 counts and starts a new window now. */ +export function __resetUnknownSession404CountsForTesting(): void { + resetUnknownSession404Counts(Date.now()); +} + /** * Write the 404 JSON-RPC "Session not found" response for an /mcp request - * that carried an unknown or expired Mcp-Session-Id. MCP Streamable HTTP, + * that carried an unknown or expired Mcp-Session-Id, and count it under its + * method for {@link flushUnknownSession404Counts}. MCP Streamable HTTP, * session management: a client that gets this 404 must start a new session * with an initialize request that carries no session id. */ @@ -1580,6 +1641,7 @@ export function writeUnknownSession404( res: Response, opts: { method: "POST" | "GET" | "DELETE"; body?: unknown }, ): void { + unknownSession404Counts[opts.method] += 1; // JSON-RPC 2.0: the response id echoes the request id, and is null only // when it cannot be determined. Only a single message with a string or // number id qualifies; GET/DELETE carry no body, so they stay null. @@ -4386,7 +4448,11 @@ async function startServerInner(options?: ServerOptions): Promise { // Start the idle-session reaper. Running it from here (rather than at // module import) keeps test imports free of leaked timers. if (!sessionReaperInterval) { - sessionReaperInterval = setInterval(reapIdleSessionsTick, 5 * 60 * 1000); + resetUnknownSession404Counts(Date.now()); + sessionReaperInterval = setInterval( + reapIdleSessionsTick, + SESSION_REAPER_INTERVAL_MS, + ); } console.log( `[startup] IP rate limit: ${maxSessionsPerIp} sessions/IP, TTL: ${serverCfg.server.session_ttl_minutes ?? 30}m, unused TTL: ${serverCfg.server.session_unused_ttl_minutes ?? 15}m, global cap: ${MAX_SESSIONS}`, @@ -4565,6 +4631,7 @@ async function startServerInner(options?: ServerOptions): Promise { clearInterval(telemetryFlushInterval); telemetryFlushInterval = undefined; } + flushUnknownSession404Counts(); if (sessionReaperInterval) { clearInterval(sessionReaperInterval); sessionReaperInterval = undefined; From 996eb0edfa57b498a4b5696aa3e756d0a73b01f3 Mon Sep 17 00:00:00 2001 From: Jordan Ritter Date: Thu, 1 Oct 2026 21:26:09 -0700 Subject: [PATCH 3/4] Return {server, stop} from startServer and add an in-process test server helper stop() closes the HTTP server, flushes the 404 counts, clears the reaper and telemetry intervals and removes the signal listeners. Production callers ignore the result. --- src/__tests__/helpers/inProcessServer.ts | 62 +++++++ src/__tests__/started-server-stop.test.ts | 200 ++++++++++++++++++++++ src/server.ts | 64 ++++++- 3 files changed, 322 insertions(+), 4 deletions(-) create mode 100644 src/__tests__/helpers/inProcessServer.ts create mode 100644 src/__tests__/started-server-stop.test.ts diff --git a/src/__tests__/helpers/inProcessServer.ts b/src/__tests__/helpers/inProcessServer.ts new file mode 100644 index 0000000..b5fd4af --- /dev/null +++ b/src/__tests__/helpers/inProcessServer.ts @@ -0,0 +1,62 @@ +// Boot the real app in-process for route-level tests. +// +// The helper calls startServer({ port: 0 }), so the OS picks a free port, and +// reads the port back from the bound server's address. It resolves once the +// server is listening. It rejects on the server's first "error" event, and on +// any other failure after startServer() returns. On those failures it calls +// stop() and then rethrows the failure. If stop() also rejects, it throws an +// AggregateError of [failure, stop error] with the failure as its cause. It +// does not intercept process.exit: the server's own "error" handler still +// runs shutdown(), which ends the process. With port 0, a bind conflict is +// not a realistic case. +// +// stop() is StartedServer.stop(). It closes this listener and its +// connections, writes out the pending unknown-session 404 counts, clears the +// module's session reaper and telemetry flush intervals, and removes the +// SIGINT/SIGTERM listeners that this boot added. It does not exit the +// process. It does NOT close the DB pool, live MCP transports or the +// nightly-reindex interval, and it does not unmount the routes that +// startServer() added to the module-level app. +// +// That module state is shared by every boot in the module instance. Run one +// in-process server at a time per test file, and stop it before the file +// boots another. +// +// The calling test file must vi.mock src/config.js before this runs, the same +// way it would for a direct startServer() call. From a test file in +// src/__tests__/, that path is "../config.js" (vi.mock paths are relative to +// the file that calls vi.mock, not to this helper). +import { once } from "node:events"; +import { startServer } from "../../server.js"; + +export interface InProcessServer { + baseUrl: string; + stop(): Promise; +} + +export async function startInProcessServer(): Promise { + const started = await startServer({ port: 0 }); + try { + if (!started.server.listening) await once(started.server, "listening"); + const address = started.server.address(); + if (address === null || typeof address === "string") { + throw new Error(`expected a TCP address, got ${String(address)}`); + } + return { + baseUrl: `http://127.0.0.1:${address.port}`, + stop: () => started.stop(), + }; + } catch (err) { + try { + await started.stop(); + } catch (stopErr) { + // Keep the original failure: it is the root cause. + throw new AggregateError( + [err, stopErr], + "startInProcessServer failed, and stop() also failed", + { cause: err }, + ); + } + throw err; + } +} diff --git a/src/__tests__/started-server-stop.test.ts b/src/__tests__/started-server-stop.test.ts new file mode 100644 index 0000000..07329b7 --- /dev/null +++ b/src/__tests__/started-server-stop.test.ts @@ -0,0 +1,200 @@ +/** + * StartedServer.stop(): what it closes, and how startInProcessServer() uses + * it on a failure. + * + * stop() must close the server even when the bind has not finished yet, so + * a late bind cannot leave a listening server behind. startServer() calls + * app.listen(port) with no host, and Node binds that synchronously, so the + * server is already listening when startServer() resolves. To get a bind + * that is still pending when stop() runs, the first test adds the host + * "localhost" to that one listen() call. Node then looks the host up with + * dns.lookup() before it binds. That test calls startServer() directly, not + * startInProcessServer(), because the helper waits for "listening". + */ +import dns from "node:dns"; +import net from "node:net"; +import { describe, it, expect, vi, afterAll, afterEach } from "vitest"; + +vi.mock("../config.js", async (importOriginal) => ({ + ...(await importOriginal()), + getConfig: vi.fn().mockReturnValue({ + port: 0, + databaseUrl: "pglite:///tmp/test-started-server-stop", + openaiApiKey: "", + githubToken: "", + githubWebhookSecret: "", + nodeEnv: "test", + logLevel: "info", + cloneDir: "/tmp/test-started-server-stop", + slackBotToken: "", + slackSigningSecret: "", + discordBotToken: "", + discordPublicKey: "", + notionToken: "", + mcpJwtSecret: "f".repeat(64), + p2pTelemetryUrl: undefined, + p2pTelemetryDisabled: true, + packageVersion: "test", + }), + getServerConfig: vi.fn().mockReturnValue({ + server: { + name: "pathfinder-started-server-stop", + version: "0.0.0", + max_sessions_per_ip: 50, + session_ttl_minutes: 30, + allowlist: [], + trust_proxy: false, + }, + sources: [], + tools: [], + }), + getAnalyticsConfig: vi.fn().mockReturnValue(undefined), + hasSearchTools: vi.fn().mockReturnValue(false), + hasKnowledgeTools: vi.fn().mockReturnValue(false), + hasCollectTools: vi.fn().mockReturnValue(false), + hasBashSemanticSearch: vi.fn().mockReturnValue(false), +})); + +import { startServer } from "../server.js"; +import { startInProcessServer } from "./helpers/inProcessServer.js"; + +/** Make the next listen() on any server bind to "localhost", so it is pending. */ +function deferNextBindToLocalhost(): void { + const realListen = net.Server.prototype.listen; + vi.spyOn(net.Server.prototype, "listen").mockImplementationOnce(function ( + this: net.Server, + ...args: unknown[] + ): net.Server { + const [port, ...rest] = args; + return Reflect.apply(realListen, this, [port, "localhost", ...rest]); + }); +} + +/** + * Resolve once the next dns.lookup("localhost") has called back and the + * work that callback scheduled (a bind, then a nextTick "listening") has run. + */ +function nextLocalhostLookupSettled(): Promise { + const realLookup = dns.lookup; + return new Promise((resolve) => { + const spy = vi.spyOn(dns, "lookup").mockImplementation((( + ...args: unknown[] + ) => { + const last = args.length - 1; + const callback = args[last]; + if (args[0] === "localhost" && typeof callback === "function") { + spy.mockRestore(); + args[last] = (...result: unknown[]) => { + Reflect.apply(callback, undefined, result); + setImmediate(resolve); + }; + } + return Reflect.apply(realLookup, dns, args); + }) as typeof dns.lookup); + }); +} + +/** Close `server` for real if a test left it listening. */ +async function closeIfListening(server: net.Server): Promise { + if (server.listening) { + await new Promise((resolve) => server.close(() => resolve())); + } +} + +describe("StartedServer.stop()", () => { + afterAll(() => { + vi.restoreAllMocks(); + }); + + it("closes a server whose bind is still pending, so it never listens", async () => { + vi.spyOn(console, "log").mockImplementation(() => {}); + const lookupSettled = nextLocalhostLookupSettled(); + deferNextBindToLocalhost(); + const started = await startServer({ port: 0 }); + vi.mocked(net.Server.prototype.listen).mockRestore(); + let listened = false; + started.server.on("listening", () => (listened = true)); + try { + expect(started.server.listening).toBe(false); + await started.stop(); + // Wait for the pending lookup to call back. A bind that stop() did not + // cancel happens in that callback. + await lookupSettled; + expect(listened).toBe(false); + expect(started.server.listening).toBe(false); + } finally { + await closeIfListening(started.server); + } + }); +}); + +describe("startInProcessServer() failure path", () => { + afterEach(() => { + vi.restoreAllMocks(); + }); + + /** Spy on listen() to capture the server that startServer() creates. */ + function captureServer(): () => net.Server { + let captured: net.Server | undefined; + const realListen = net.Server.prototype.listen; + vi.spyOn(net.Server.prototype, "listen").mockImplementationOnce(function ( + this: net.Server, + ...args: unknown[] + ): net.Server { + captured = this; + return Reflect.apply(realListen, this, args); + }); + return () => { + if (!captured) throw new Error("startServer() did not call listen()"); + return captured; + }; + } + + it("calls stop() and rethrows the original error", async () => { + vi.spyOn(console, "log").mockImplementation(() => {}); + const baseSigterm = process.listenerCount("SIGTERM"); + const server = captureServer(); + vi.spyOn(net.Server.prototype, "address").mockReturnValueOnce(null); + try { + await expect(startInProcessServer()).rejects.toThrow( + "expected a TCP address, got null", + ); + // stop() ran: the listener is closed and the signal listener is gone. + expect(server().listening).toBe(false); + expect(process.listenerCount("SIGTERM")).toBe(baseSigterm); + } finally { + await closeIfListening(server()); + } + }); + + it("keeps the original error when stop() also fails", async () => { + vi.spyOn(console, "log").mockImplementation(() => {}); + const server = captureServer(); + vi.spyOn(net.Server.prototype, "address").mockReturnValueOnce(null); + const closeErr = new Error("close failed"); + vi.spyOn(net.Server.prototype, "close").mockImplementationOnce(function ( + this: net.Server, + callback?: (err?: Error) => void, + ): net.Server { + callback?.(closeErr); + return this; + }); + try { + const err: unknown = await startInProcessServer().then( + () => undefined, + (e: unknown) => e, + ); + expect(err).toBeInstanceOf(AggregateError); + if (!(err instanceof AggregateError)) return; + expect(err.errors).toHaveLength(2); + expect(String(err.errors[0])).toContain( + "expected a TCP address, got null", + ); + expect(err.errors[1]).toBe(closeErr); + expect(err.cause).toBe(err.errors[0]); + } finally { + vi.mocked(net.Server.prototype.close).mockRestore(); + await closeIfListening(server()); + } + }); +}); diff --git a/src/server.ts b/src/server.ts index cf10a1d..c2e139c 100644 --- a/src/server.ts +++ b/src/server.ts @@ -7,6 +7,7 @@ import express, { import compression from "compression"; import cors from "cors"; import { randomUUID, timingSafeEqual } from "node:crypto"; +import type { Server as HttpServer } from "node:http"; import { Bash } from "just-bash"; import { StreamableHTTPServerTransport } from "@modelcontextprotocol/sdk/server/streamableHttp.js"; import type { SSEServerTransport } from "@modelcontextprotocol/sdk/server/sse.js"; @@ -118,6 +119,28 @@ export interface ServerOptions { configPath?: string; } +/** + * What startServer() resolves to. Production callers ignore it. `stop` is a + * test seam with a narrow scope. It closes this HTTP server and its + * connections (also when the bind has not finished), writes out the pending + * unknown-session 404 counts, clears the module's session reaper and + * telemetry flush intervals, and removes the SIGINT/SIGTERM listeners that + * this startServer() call added. It does not call process.exit. + * + * It does NOT undo the rest of startServer(): the routes mounted on the + * module-level app, the nightly-reindex interval, the DB pool, live MCP + * transports, workspace and telemetry state, and the server's "error" + * listener (which calls shutdown() and process.exit) all remain. + * + * The intervals, the 404 counts and the other state above are shared by the + * whole module instance, not owned by one server. Run one in-process server + * at a time per test file, and stop it before the file boots another. + */ +export interface StartedServer { + server: HttpServer; + stop(): Promise; +} + const app = express(); app.use( cors({ @@ -4293,7 +4316,9 @@ export function registerAdminOpsRoutes( // Startup // --------------------------------------------------------------------------- -export async function startServer(options?: ServerOptions): Promise { +export async function startServer( + options?: ServerOptions, +): Promise { // Top-level try/catch around the entire startup sequence so synchronous // throws from getConfig/getServerConfig AND async failures from // initializeSchema/checkAndIndex carry a uniform '[startup] fatal:' log @@ -4310,7 +4335,9 @@ export async function startServer(options?: ServerOptions): Promise { } } -async function startServerInner(options?: ServerOptions): Promise { +async function startServerInner( + options?: ServerOptions, +): Promise { if (options?.configPath) { process.env.PATHFINDER_CONFIG = options.configPath; } @@ -4665,6 +4692,35 @@ async function startServerInner(options?: ServerOptions): Promise { process.exit(0); } - process.on("SIGINT", () => shutdown("SIGINT")); - process.on("SIGTERM", () => shutdown("SIGTERM")); + const onSigint = () => shutdown("SIGINT"); + const onSigterm = () => shutdown("SIGTERM"); + process.on("SIGINT", onSigint); + process.on("SIGTERM", onSigterm); + + return { + server, + async stop() { + process.off("SIGINT", onSigint); + process.off("SIGTERM", onSigterm); + if (telemetryFlushInterval) { + clearInterval(telemetryFlushInterval); + telemetryFlushInterval = undefined; + } + flushUnknownSession404Counts(); + if (sessionReaperInterval) { + clearInterval(sessionReaperInterval); + sessionReaperInterval = undefined; + } + // close() also cancels a bind that is still pending. Its callback then + // gets ERR_SERVER_NOT_RUNNING, which is not a failure here. + server.closeAllConnections(); + await new Promise((resolve, reject) => + server.close((err?: NodeJS.ErrnoException) => + !err || err.code === "ERR_SERVER_NOT_RUNNING" + ? resolve() + : reject(err), + ), + ); + }, + }; } From 3e9b7f53615c1b4ad46cf7b44351f4d1b8702a90 Mon Sep 17 00:00:00 2001 From: Jordan Ritter Date: Thu, 1 Oct 2026 21:26:11 -0700 Subject: [PATCH 4/4] Add route-level tests for the /mcp unknown-session 404 --- .../mcp-unknown-session-404-routes.test.ts | 811 ++++++++++++++++++ 1 file changed, 811 insertions(+) create mode 100644 src/__tests__/mcp-unknown-session-404-routes.test.ts diff --git a/src/__tests__/mcp-unknown-session-404-routes.test.ts b/src/__tests__/mcp-unknown-session-404-routes.test.ts new file mode 100644 index 0000000..b5dc646 --- /dev/null +++ b/src/__tests__/mcp-unknown-session-404-routes.test.ts @@ -0,0 +1,811 @@ +/** + * Route-level coverage for the unknown-session 404 on /mcp. Boots the real + * app in-process with startInProcessServer() and drives POST/GET/DELETE /mcp over + * HTTP, so deleting or inverting a route's 404 branch, or reverting the POST + * initialize gate, fails a test. The helper-level tests live in + * mcp-unknown-session-404.test.ts. + */ +import { + describe, + it, + expect, + vi, + beforeAll, + afterAll, + beforeEach, + type MockInstance, +} from "vitest"; +import http from "node:http"; +import net from "node:net"; + +vi.mock("../config.js", async (importOriginal) => ({ + ...(await importOriginal()), + getConfig: vi.fn().mockReturnValue({ + port: 0, + databaseUrl: "pglite:///tmp/test-mcp-unknown-session-404-routes", + openaiApiKey: "", + githubToken: "", + githubWebhookSecret: "", + nodeEnv: "test", + logLevel: "info", + cloneDir: "/tmp/test-mcp-unknown-session-404-routes", + slackBotToken: "", + slackSigningSecret: "", + discordBotToken: "", + discordPublicKey: "", + notionToken: "", + mcpJwtSecret: "e".repeat(64), + p2pTelemetryUrl: undefined, + p2pTelemetryDisabled: true, + packageVersion: "test", + }), + getServerConfig: vi.fn().mockReturnValue({ + server: { + name: "pathfinder-unknown-session-routes", + version: "0.0.0", + max_sessions_per_ip: 50, + session_ttl_minutes: 30, + allowlist: [], + trust_proxy: false, + }, + sources: [], + tools: [], + }), + getAnalyticsConfig: vi.fn().mockReturnValue(undefined), + hasSearchTools: vi.fn().mockReturnValue(false), + hasKnowledgeTools: vi.fn().mockReturnValue(false), + hasCollectTools: vi.fn().mockReturnValue(false), + hasBashSemanticSearch: vi.fn().mockReturnValue(false), +})); + +import type { Response } from "express"; +import { + __resetUnknownSession404CountsForTesting, + flushUnknownSession404Counts, + SESSION_REAPER_INTERVAL_MS, + writeUnknownSession404, +} from "../server.js"; +import { + startInProcessServer, + type InProcessServer, +} from "./helpers/inProcessServer.js"; + +type Method = "POST" | "GET" | "DELETE"; +const METHODS: Method[] = ["POST", "GET", "DELETE"]; +const SID = "deadbeef-0000-4000-8000-000000000000"; + +const PROTOTYPE_KEY_SIDS = [ + "constructor", + "__proto__", + "toString", + "hasOwnProperty", +]; + +/** The 404 body. POST echoes a string or number request id; GET and DELETE carry no body, so their id is null. */ +function notFoundBody(id: string | number | null) { + return { + jsonrpc: "2.0", + error: { code: -32001, message: "Session not found" }, + id, + }; +} + +/** The answer to a request with no (or an empty) session header. */ +function noSessionBody(message: string) { + return { jsonrpc: "2.0", error: { code: -32000, message }, id: null }; +} + +const NO_SESSION_ANSWERS = [ + [ + "POST", + 400, + "Bad Request: No valid session. Send an initialize request first.", + ], + ["GET", 405, "Method Not Allowed"], + ["DELETE", 400, "Missing session ID"], +] as const; + +function expectedIdFor(method: Method): number | null { + return method === "POST" ? TOOLS_LIST.id : null; +} + +const TOOLS_LIST = { jsonrpc: "2.0", id: 1, method: "tools/list" }; +const INITIALIZE = { + jsonrpc: "2.0", + id: 1, + method: "initialize", + params: { + protocolVersion: "2025-03-26", + capabilities: {}, + clientInfo: { name: "route-test", version: "0.0.0" }, + }, +}; + +let running: InProcessServer | undefined; +let port = 0; +let warnSpy: MockInstance | undefined; +let logSpy: MockInstance | undefined; +let errorSpy: MockInstance | undefined; +/** The callback startServer() scheduled as the session reaper tick. */ +let reaperTick: (() => void) | undefined; + +type HttpResult = { + status: number; + headers: http.IncomingHttpHeaders; + body: string; +}; + +/** Send one request to /mcp and collect its status, headers and body. */ +function mcpRequest( + method: Method, + headers: Record, + body?: unknown, +): Promise { + return new Promise((resolve, reject) => { + const payload = body === undefined ? undefined : JSON.stringify(body); + const req = http.request( + { + hostname: "127.0.0.1", + port, + path: "/mcp", + method, + headers: { + Accept: "application/json, text/event-stream", + ...(payload ? { "Content-Type": "application/json" } : {}), + ...headers, + }, + }, + (res) => { + let data = ""; + res.on("data", (chunk) => (data += chunk)); + res.on("end", () => + resolve({ + status: res.statusCode ?? 0, + headers: res.headers, + body: data, + }), + ); + }, + ); + req.on("error", reject); + if (payload) req.write(payload); + req.end(); + }); +} + +/** + * Resolve with the status as soon as the response headers arrive, then drop + * the connection. A GET on a live session opens an SSE stream that never + * ends, so mcpRequest() would hang on it. + */ +function mcpStatusOnly( + method: Method, + headers: Record, +): Promise { + return new Promise((resolve, reject) => { + const req = http.request( + { + hostname: "127.0.0.1", + port, + path: "/mcp", + method, + headers: { + Accept: "application/json, text/event-stream", + ...headers, + }, + }, + (res) => { + const status = res.statusCode ?? 0; + res.resume(); + req.destroy(); + resolve(status); + }, + ); + req.on("error", (err: NodeJS.ErrnoException) => { + // The destroy() above can surface as a reset after resolve(); ignore it. + if (err.code !== "ECONNRESET") reject(err); + }); + req.end(); + }); +} + +/** + * Send one raw HTTP/1.1 request whose Mcp-Session-Id header carries `sid` + * byte for byte (latin1), which http.request() refuses to do for control + * characters. A POST carries bodyFor("POST") as its JSON body. Resolves with + * the response status. + */ +function rawSessionRequest(method: Method, sid: string): Promise { + const body = bodyFor(method); + const payload = body === undefined ? "" : JSON.stringify(body); + const bodyHeaders = payload + ? "Content-Type: application/json\r\n" + + `Content-Length: ${Buffer.byteLength(payload)}\r\n` + : ""; + return new Promise((resolve, reject) => { + const socket = net.connect(port, "127.0.0.1", () => { + socket.write( + Buffer.from( + `${method} /mcp HTTP/1.1\r\nHost: 127.0.0.1\r\n` + + "Accept: application/json, text/event-stream\r\n" + + bodyHeaders + + `Mcp-Session-Id: ${sid}\r\nConnection: close\r\n\r\n` + + payload, + "latin1", + ), + ); + }); + let data = ""; + socket.on("data", (chunk: Buffer) => (data += chunk.toString("latin1"))); + socket.on("end", () => { + const m = /^HTTP\/1\.1 (\d{3}) /.exec(data); + if (m) resolve(Number(m[1])); + else reject(new Error(`no status line in ${JSON.stringify(data)}`)); + }); + socket.on("error", reject); + }); +} + +/** Every string written to console.warn/log/error since their last clear. */ +function consoleOutput(): string[] { + return [warnSpy, logSpy, errorSpy].flatMap((spy) => + (spy?.mock.calls ?? []).map((args: unknown[]) => + args.map(String).join(" "), + ), + ); +} + +function unknownSessionLines(spy: MockInstance): string[] { + return spy.mock.calls + .map((args) => String(args[0])) + .filter((l) => l.startsWith("[mcp] 404 unknown-session-id")); +} + +/** + * Check that each 404 flush line reports a real elapsed window: an integer + * number of seconds below 60, far below the 5 minute reaper period, since + * every window here starts at boot or at a reset in the same test. Return + * the lines with the number replaced by "N", so a test can compare the + * counts exactly. + */ +function withWindowN(lines: string[]): string[] { + return lines.map((line) => { + const m = /^\[mcp\] 404 unknown-session-id window_s=(\d+) /.exec(line); + if (!m) throw new Error(`no window_s field in ${JSON.stringify(line)}`); + expect(Number(m[1])).toBeLessThan(60); + return line.replace(/window_s=\d+/, "window_s=N"); + }); +} + +/** Initialize a new session and return its id. */ +async function initializeSession( + headers: Record = {}, +): Promise<{ res: HttpResult; sid: string | undefined }> { + const res = await mcpRequest("POST", headers, INITIALIZE); + const raw = res.headers["mcp-session-id"]; + return { res, sid: typeof raw === "string" ? raw : undefined }; +} + +function bodyFor(method: Method): unknown { + return method === "POST" ? TOOLS_LIST : undefined; +} + +describe("/mcp routes: unknown session id", () => { + beforeAll(async () => { + warnSpy = vi.spyOn(console, "warn").mockImplementation(() => {}); + logSpy = vi.spyOn(console, "log").mockImplementation(() => {}); + errorSpy = vi.spyOn(console, "error").mockImplementation(() => {}); + const realSetInterval = globalThis.setInterval; + // Capture the reaper tick so a test can run it on demand. + vi.spyOn(globalThis, "setInterval").mockImplementation((( + ...args: Parameters + ) => { + const [callback, ms] = args; + if (ms === SESSION_REAPER_INTERVAL_MS && typeof callback === "function") { + reaperTick = () => callback(); + } + return realSetInterval(...args); + }) as typeof setInterval); + running = await startInProcessServer(); + port = Number(new URL(running.baseUrl).port); + }); + + afterAll(async () => { + try { + await running?.stop(); + } finally { + running = undefined; + vi.restoreAllMocks(); + } + }); + + beforeEach(() => { + __resetUnknownSession404CountsForTesting(); + }); + + it.each(METHODS)( + "%s with an unknown session id returns 404 and the -32001 body", + async (method) => { + const res = await mcpRequest( + method, + { "Mcp-Session-Id": SID }, + bodyFor(method), + ); + expect(res.status).toBe(404); + expect(JSON.parse(res.body)).toEqual(notFoundBody(expectedIdFor(method))); + }, + ); + + describe.each(METHODS)("%s with a prototype-key session id", (method) => { + it.each(PROTOTYPE_KEY_SIDS)("'%s' returns 404", async (sid) => { + const res = await mcpRequest( + method, + { "Mcp-Session-Id": sid }, + bodyFor(method), + ); + expect(res.status).toBe(404); + expect(JSON.parse(res.body)).toEqual(notFoundBody(expectedIdFor(method))); + }); + }); + + it.each([ + ["a string id", "req-7", "req-7"], + ["a number id", 42, 42], + ["no id", undefined, null], + ["a null id", null, null], + ["an object id", { nested: 1 }, null], + ] as const)( + "POST 404 for an unknown session with %s answers with id %j", + async (_label, id, expected) => { + const body = + id === undefined + ? { jsonrpc: "2.0", method: "tools/list" } + : { jsonrpc: "2.0", id, method: "tools/list" }; + const res = await mcpRequest("POST", { "Mcp-Session-Id": SID }, body); + expect(res.status).toBe(404); + expect(JSON.parse(res.body)).toEqual(notFoundBody(expected)); + }, + ); + + it("a live session routes POST and GET, and DELETE tears it down", async () => { + const { res: init, sid } = await initializeSession(); + expect(init.status).toBe(200); + expect(typeof sid).toBe("string"); + let deleted = false; + try { + const live = { "Mcp-Session-Id": String(sid) }; + + const list = await mcpRequest("POST", live, TOOLS_LIST); + expect(list.status).toBe(200); + expect(list.body).not.toContain("Session not found"); + // A 200 JSON-RPC reply with the request id and no "Session not found" + // comes from the live transport, not the unknown-session 404. + expect(JSON.parse(list.body)).toMatchObject({ + jsonrpc: "2.0", + id: TOOLS_LIST.id, + }); + + // A live GET opens the SSE stream: 200, never the unknown-session 404. + const getStatus = await mcpStatusOnly("GET", live); + expect(getStatus).toBe(200); + + const del = await mcpRequest("DELETE", live); + expect(del.status).toBe(200); + deleted = true; + + // Once deleted, the id is unknown again. + const after = await mcpRequest("POST", live, TOOLS_LIST); + expect(after.status).toBe(404); + } finally { + if (!deleted && typeof sid === "string") { + await mcpRequest("DELETE", { "Mcp-Session-Id": sid }); + } + } + }); + + it("POST initialize with a stale session id starts a new session", async () => { + const { res, sid: newSid } = await initializeSession({ + "Mcp-Session-Id": SID, + }); + expect(res.status).toBe(200); + let deleted = false; + try { + expect(typeof newSid).toBe("string"); + expect(newSid).not.toBe(SID); + expect(res.body).not.toContain("Session not found"); + + // The new session is live, and the stale id was not registered. + const live = await mcpRequest( + "POST", + { "Mcp-Session-Id": String(newSid) }, + TOOLS_LIST, + ); + expect(live.status).toBe(200); + const stale = await mcpRequest( + "POST", + { "Mcp-Session-Id": SID }, + TOOLS_LIST, + ); + expect(stale.status).toBe(404); + expect(JSON.parse(stale.body)).toEqual(notFoundBody(TOOLS_LIST.id)); + + // Tear the new session down through the real DELETE route. + const del = await mcpRequest("DELETE", { + "Mcp-Session-Id": String(newSid), + }); + expect(del.status).toBe(200); + deleted = true; + } finally { + if (!deleted && typeof newSid === "string") { + await mcpRequest("DELETE", { "Mcp-Session-Id": newSid }); + } + } + }); + + it("the flush writes one line with the per-method unknown-session 404 counts", async () => { + const counts: Record = { POST: 3, GET: 2, DELETE: 1 }; + for (const method of METHODS) { + for (let i = 0; i < counts[method]; i++) { + const res = await mcpRequest( + method, + { "Mcp-Session-Id": `${SID}-${method}-${i}` }, + bodyFor(method), + ); + expect(res.status).toBe(404); + } + } + // A request with no session header is not an unknown-session 404. + await mcpRequest("GET", {}); + + const spy = warnSpy; + if (!spy) throw new Error("console.warn spy not installed"); + spy.mockClear(); + flushUnknownSession404Counts(); + expect(withWindowN(unknownSessionLines(spy))).toEqual([ + "[mcp] 404 unknown-session-id window_s=N total=6 GET=2 POST=3 DELETE=1", + ]); + + // The flush resets the counts, so a second flush writes nothing. + spy.mockClear(); + flushUnknownSession404Counts(); + expect(spy).not.toHaveBeenCalled(); + }); + + it("the reaper tick flushes the unknown-session 404 counts", async () => { + const tick = reaperTick; + const spy = warnSpy; + if (!tick) throw new Error("reaper interval was not scheduled at boot"); + if (!spy) throw new Error("console.warn spy not installed"); + for (const method of METHODS) { + const res = await mcpRequest( + method, + { "Mcp-Session-Id": SID }, + bodyFor(method), + ); + expect(res.status).toBe(404); + } + spy.mockClear(); + tick(); + expect(withWindowN(unknownSessionLines(spy))).toEqual([ + "[mcp] 404 unknown-session-id window_s=N total=3 GET=1 POST=1 DELETE=1", + ]); + }); + + it("logs no session id bytes or client ip, even for ids with newlines or control characters", async () => { + const MARKER = "E1MARKERSID"; + // Node's HTTP parser answers 400 to a header value with LF, CR or a C0 + // control other than tab, so those never reach the route. Tab and + // obs-text bytes (0x80-0xff) do, and get the unknown-session 404. + const rejectedByParser = [ + `${MARKER}\n[mcp] forged line sid=x`, + `${MARKER}\r\n folded`, + `${MARKER}\u0001\u007fkey=value`, + `${MARKER}\u001b[31mred`, + ]; + const reachRoute = [`${MARKER}\tkey=value`, `${MARKER}\u00ff\u00fezz`]; + warnSpy?.mockClear(); + logSpy?.mockClear(); + errorSpy?.mockClear(); + for (const sid of rejectedByParser) { + expect(await rawSessionRequest("GET", sid)).toBe(400); + } + for (const sid of reachRoute) { + for (const method of METHODS) { + expect(await rawSessionRequest(method, sid)).toBe(404); + } + } + flushUnknownSession404Counts(); + + const output = consoleOutput(); + expect( + withWindowN( + output.filter((l) => l.startsWith("[mcp] 404 unknown-session-id")), + ), + ).toEqual([ + "[mcp] 404 unknown-session-id window_s=N total=6 GET=2 POST=2 DELETE=2", + ]); + const joined = output.join("\n"); + for (const needle of [ + MARKER, + "forged", + "key=value", + "folded", + "127.0.0.1", + "\t", + "\r", + "\u0001", + "\u001b", + "\u00ff", + ]) { + expect(joined).not.toContain(needle); + } + }); + + it.each(NO_SESSION_ANSWERS)( + "%s with no session header keeps its %i", + async (method, status, message) => { + const res = await mcpRequest(method, {}, bodyFor(method)); + expect(res.status).toBe(status); + expect(JSON.parse(res.body)).toEqual(noSessionBody(message)); + }, + ); + + // An empty header is falsy, so it is treated the same as no header. + it.each(NO_SESSION_ANSWERS)( + "%s with an empty session id is handled as no header (%i)", + async (method, status, message) => { + const res = await mcpRequest( + method, + { "Mcp-Session-Id": "" }, + bodyFor(method), + ); + expect(res.status).toBe(status); + expect(JSON.parse(res.body)).toEqual(noSessionBody(message)); + }, + ); + + // POST and DELETE /mcp sit behind bearerMiddleware, which answers 401 to a + // Bearer token that is present but invalid before the route looks at the + // session id. GET /mcp has no auth middleware, so it reaches the 404. + it.each([ + ["POST", 401], + ["GET", 404], + ["DELETE", 401], + ] as const)( + "%s with an invalid bearer token and an unknown session id answers %i", + async (method, status) => { + const res = await mcpRequest( + method, + { Authorization: "Bearer not-a-jwt", "Mcp-Session-Id": SID }, + bodyFor(method), + ); + expect(res.status).toBe(status); + const spy = warnSpy; + if (!spy) throw new Error("console.warn spy not installed"); + spy.mockClear(); + flushUnknownSession404Counts(); + expect(withWindowN(unknownSessionLines(spy))).toEqual( + status === 404 + ? [ + "[mcp] 404 unknown-session-id window_s=N total=1 GET=1 POST=0 DELETE=0", + ] + : [], + ); + }, + ); +}); + +// Self-contained: boots its own server, so it passes when run alone. It runs +// after the suite above has stopped its server, so one server runs at a time. +describe("teardown hygiene", () => { + it("stop() removes its signal listeners, clears its intervals and writes the pending 404 counts", async () => { + const baseSigint = process.listenerCount("SIGINT"); + const baseSigterm = process.listenerCount("SIGTERM"); + const created: unknown[] = []; + const cleared = new Set(); + const realSetInterval = globalThis.setInterval; + const realClearInterval = globalThis.clearInterval; + const warn = vi.spyOn(console, "warn").mockImplementation(() => {}); + vi.spyOn(console, "log").mockImplementation(() => {}); + vi.spyOn(console, "error").mockImplementation(() => {}); + // Both spies stay in place until after stop(), so intervals created after + // boot (by fire-and-forget boot work or by requests) are recorded too. + vi.spyOn(globalThis, "setInterval").mockImplementation((( + ...args: Parameters + ) => { + const h = realSetInterval(...args); + created.push(h); + return h; + }) as typeof setInterval); + vi.spyOn(globalThis, "clearInterval").mockImplementation((( + h: Parameters[0], + ) => { + cleared.add(h); + realClearInterval(h); + }) as typeof clearInterval); + try { + __resetUnknownSession404CountsForTesting(); + const server = await startInProcessServer(); + try { + expect(process.listenerCount("SIGINT")).toBe(baseSigint + 1); + expect(process.listenerCount("SIGTERM")).toBe(baseSigterm + 1); + expect(created.length).toBeGreaterThan(0); + // Leave unflushed unknown-session 404 counts for stop() to write. + port = Number(new URL(server.baseUrl).port); + expect( + (await mcpRequest("GET", { "Mcp-Session-Id": `${SID}-stop-1` })) + .status, + ).toBe(404); + expect( + (await mcpRequest("DELETE", { "Mcp-Session-Id": `${SID}-stop-2` })) + .status, + ).toBe(404); + warn.mockClear(); + } finally { + await server.stop(); + } + expect(withWindowN(unknownSessionLines(warn))).toEqual([ + "[mcp] 404 unknown-session-id window_s=N total=2 GET=1 POST=0 DELETE=1", + ]); + expect(process.listenerCount("SIGINT")).toBe(baseSigint); + expect(process.listenerCount("SIGTERM")).toBe(baseSigterm); + expect(created.filter((h) => !cleared.has(h))).toEqual([]); + } finally { + vi.restoreAllMocks(); + } + }); + + it("boot starts a new 404 window and drops the counts from before boot", async () => { + const warn = vi.spyOn(console, "warn").mockImplementation(() => {}); + vi.spyOn(console, "log").mockImplementation(() => {}); + vi.spyOn(console, "error").mockImplementation(() => {}); + try { + // Start a window 1000 s ago and leave one count in it, before boot. + const realNow = Date.now.bind(Date); + const nowSpy = vi + .spyOn(Date, "now") + .mockImplementation(() => realNow() - 1_000_000); + __resetUnknownSession404CountsForTesting(); + nowSpy.mockRestore(); + const fakeRes = { status: () => fakeRes, json: () => fakeRes }; + writeUnknownSession404(fakeRes as unknown as Response, { + method: "POST", + }); + + const server = await startInProcessServer(); + try { + port = Number(new URL(server.baseUrl).port); + expect( + (await mcpRequest("GET", { "Mcp-Session-Id": `${SID}-boot` })).status, + ).toBe(404); + warn.mockClear(); + flushUnknownSession404Counts(); + // withWindowN() also checks that window_s is far below 1000. + expect(withWindowN(unknownSessionLines(warn))).toEqual([ + "[mcp] 404 unknown-session-id window_s=N total=1 GET=1 POST=0 DELETE=0", + ]); + } finally { + await server.stop(); + } + } finally { + vi.restoreAllMocks(); + } + }); + + it("stop() closes a connection that still has a response in flight", async () => { + vi.spyOn(console, "warn").mockImplementation(() => {}); + vi.spyOn(console, "log").mockImplementation(() => {}); + vi.spyOn(console, "error").mockImplementation(() => {}); + let stream: http.ClientRequest | undefined; + try { + __resetUnknownSession404CountsForTesting(); + const server = await startInProcessServer(); + let stopped = false; + try { + port = Number(new URL(server.baseUrl).port); + const { sid } = await initializeSession(); + if (typeof sid !== "string") throw new Error("no session id"); + // A live GET opens an SSE stream that never ends on its own. Its + // connection is not idle, so server.close() alone waits for it. + const res = await new Promise( + (resolve, reject) => { + stream = http.request( + { + hostname: "127.0.0.1", + port, + path: "/mcp", + method: "GET", + headers: { + Accept: "text/event-stream", + "Mcp-Session-Id": sid, + }, + }, + resolve, + ); + stream.on("error", reject); + stream.end(); + }, + ); + expect(res.statusCode).toBe(200); + // When stop() cuts the stream, the response emits "aborted" as an + // error and then "close". The cut is the expected outcome here. + stream?.removeAllListeners("error"); + stream?.on("error", () => {}); + let aborted = false; + res.on("error", () => (aborted = true)); + const streamClosed = new Promise((resolve) => + res.on("close", () => resolve()), + ); + + const outcome = await Promise.race([ + server.stop().then(() => "stopped" as const), + new Promise<"timed out">((resolve) => + setTimeout(() => resolve("timed out"), 2000), + ), + ]); + stopped = outcome === "stopped"; + expect(outcome).toBe("stopped"); + await streamClosed; + expect(aborted).toBe(true); + } finally { + if (!stopped) { + stream?.destroy(); + await server.stop(); + } + } + } finally { + stream?.destroy(); + vi.restoreAllMocks(); + } + }); + + it("shutdown() writes the pending 404 counts before it exits", async () => { + const warn = vi.spyOn(console, "warn").mockImplementation(() => {}); + vi.spyOn(console, "log").mockImplementation(() => {}); + vi.spyOn(console, "error").mockImplementation(() => {}); + const before = new Set(process.listeners("SIGTERM")); + let exitSpy: MockInstance | undefined; + try { + __resetUnknownSession404CountsForTesting(); + const server = await startInProcessServer(); + try { + // The SIGTERM listener that this boot added calls shutdown(). + const added = process + .listeners("SIGTERM") + .filter((l) => !before.has(l)); + expect(added).toHaveLength(1); + port = Number(new URL(server.baseUrl).port); + for (const method of ["POST", "DELETE"] as const) { + const res = await mcpRequest( + method, + { "Mcp-Session-Id": `${SID}-shutdown` }, + bodyFor(method), + ); + expect(res.status).toBe(404); + } + // shutdown() ends with process.exit(0). Stub it in this test only. + let onExit: ( + code: number | string | null | undefined, + ) => void = () => {}; + const exited = new Promise( + (resolve) => (onExit = resolve), + ); + exitSpy = vi + .spyOn(process, "exit") + .mockImplementation(((code) => onExit(code)) as typeof process.exit); + warn.mockClear(); + added[0]("SIGTERM"); + expect(await exited).toBe(0); + // Checked before stop(), which flushes too. + expect(withWindowN(unknownSessionLines(warn))).toEqual([ + "[mcp] 404 unknown-session-id window_s=N total=2 GET=0 POST=1 DELETE=1", + ]); + } finally { + exitSpy?.mockRestore(); + await server.stop(); + } + } finally { + vi.restoreAllMocks(); + } + }); +});