diff --git a/src/runner/server.session-lifetime.test.ts b/src/runner/server.session-lifetime.test.ts new file mode 100644 index 00000000..74943a4d --- /dev/null +++ b/src/runner/server.session-lifetime.test.ts @@ -0,0 +1,166 @@ +import http from 'http'; +import type { AddressInfo } from 'net'; +import { expect, onTestFinished, test } from 'vitest'; +import { runServerConformanceTest } from './server'; + +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + +test('closes a connection completed after timeout before the abandoned scenario can use it', async () => { + const initialized = deferred(); + const lateActivity = deferred(); + const liveSessions = new Set(); + const events: Array<{ method: string; sessionId?: string }> = []; + let issuedSessions = 0; + + const respond = (res: http.ServerResponse, body: unknown) => { + res.writeHead(200, { 'Content-Type': 'application/json' }); + res.end(JSON.stringify(body)); + }; + + // Real SDK Client + StreamableHTTPClientTransport, driven by the actual + // runner. This wire fixture has room for one live session at a time. + const server = http.createServer((req, res) => { + const sessionId = req.headers['mcp-session-id'] as string | undefined; + if (req.method === 'DELETE') { + events.push({ method: 'DELETE', sessionId }); + if (sessionId) liveSessions.delete(sessionId); + res.writeHead(204).end(); + if (sessionId === 's1') lateActivity.resolve(); + return; + } + if (req.method === 'GET') { + events.push({ method: 'GET', sessionId }); + res.writeHead(200, { 'Content-Type': 'text/event-stream' }); + res.write(': open\n\n'); + return; + } + + let raw = ''; + req.on('data', (chunk) => (raw += chunk)); + req.on('end', () => { + const msg = JSON.parse(raw); + if (msg.method === 'initialize') { + if (liveSessions.size > 0) { + events.push({ method: 'REJECT initialize' }); + res.writeHead(503).end('Session capacity exhausted'); + return; + } + const sessionId = `s${++issuedSessions}`; + liveSessions.add(sessionId); + events.push({ method: 'POST initialize', sessionId }); + res.setHeader('mcp-session-id', sessionId); + respond(res, { + jsonrpc: '2.0', + id: msg.id, + result: { + protocolVersion: msg.params.protocolVersion, + capabilities: { prompts: {} }, + serverInfo: { name: 'late-connect-probe', version: '0.0.1' } + } + }); + return; + } + + events.push({ method: `POST ${msg.method}`, sessionId }); + if (msg.method === 'notifications/initialized') { + if (sessionId === 's1') { + // The SDK awaits this HTTP acknowledgment before connect() resolves. + // A real session already exists, but ctx.connect() is still pending. + initialized.resolve(res); + return; + } + res.writeHead(202).end(); + return; + } + if (msg.method === 'prompts/list' && sessionId === 's2') { + respond(res, { + jsonrpc: '2.0', + id: msg.id, + result: { prompts: [{ name: 'p', description: 'healthy prompt' }] } + }); + return; + } + + respond(res, { + jsonrpc: '2.0', + id: msg.id, + error: { code: -32601, message: `Method not found: ${msg.method}` } + }); + // The unfixed runner hands s1 to the abandoned scenario. Its failed + // prompts/list skips the scenario's close(), leaving no later sweeper. + if (sessionId === 's1') lateActivity.resolve(); + }); + }); + onTestFinished(async () => { + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + }); + + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const url = `http://127.0.0.1:${(server.address() as AddressInfo).port}/mcp`; + + const pendingRun = runServerConformanceTest( + url, + 'prompts-list', + undefined, + '2025-06-18', + false, + 1000 + ); + const acknowledgment = await initialized.promise; + const timedOut = await pendingRun; + expect(timedOut.checks).toContainEqual( + expect.objectContaining({ id: 'scenario-timeout', status: 'FAILURE' }) + ); + expect([...liveSessions]).toEqual(['s1']); + expect(events.some((event) => event.method === 'DELETE')).toBe(false); + + events.push({ method: 'runner returned scenario-timeout' }); + acknowledgment.writeHead(202).end(); + // Positive network barrier: either cleanup terminates s1 or the abandoned + // scenario incorrectly probes it. No quiet-time sleeps decide the result. + await lateActivity.promise; + + const healthy = await runServerConformanceTest( + url, + 'prompts-list', + undefined, + '2025-06-18' + ); + const observation = { + lateSessionDeletes: events.filter( + (event) => event.sessionId === 's1' && event.method === 'DELETE' + ).length, + lateSessionProbes: events.filter( + (event) => + event.sessionId === 's1' && event.method === 'POST prompts/list' + ).length, + healthyPromptStatus: healthy.checks.find( + (check) => check.id === 'prompts-list' + )?.status, + healthyFailures: healthy.checks.filter( + (check) => check.status === 'FAILURE' + ).length, + healthySessionDeletes: events.filter( + (event) => event.sessionId === 's2' && event.method === 'DELETE' + ).length, + liveSessions: [...liveSessions] + }; + // Capture state before any fixture teardown can hide a leaked session. + console.log('late connection wire events:', JSON.stringify(events)); + console.log('late connection observation:', JSON.stringify(observation)); + expect(observation).toEqual({ + lateSessionDeletes: 1, + lateSessionProbes: 0, + healthyPromptStatus: 'SUCCESS', + healthyFailures: 0, + healthySessionDeletes: 1, + liveSessions: [] + }); +}); diff --git a/src/runner/server.test.ts b/src/runner/server.test.ts index 2f6ac645..6ecccc46 100644 --- a/src/runner/server.test.ts +++ b/src/runner/server.test.ts @@ -187,3 +187,150 @@ describe('runServerConformanceTest wire selection for draft-only scenarios', () }); }, 30000); }); + +describe('runServerConformanceTest session teardown', () => { + // A minimal stateful Streamable HTTP server that records the methods it is + // sent, so a test can assert on the DELETE that terminates the session. + // `promptsList` decides which exit path the prompts-list scenario takes: + // 'ok' its success path, 'error' its catch, 'hang' the runner's timeout. + let server: http.Server; + let url: string; + let methods: string[]; + let liveSessions: Set; + let promptsList: 'ok' | 'error' | 'hang'; + + const sse = (res: http.ServerResponse, body: unknown) => { + res.writeHead(200, { 'Content-Type': 'text/event-stream' }); + res.write(`event: message\ndata: ${JSON.stringify(body)}\n\n`); + res.end(); + }; + + beforeEach(async () => { + methods = []; + liveSessions = new Set(); + promptsList = 'ok'; + + server = http.createServer((req, res) => { + const sessionId = req.headers['mcp-session-id'] as string | undefined; + + if (req.method === 'DELETE') { + methods.push('DELETE'); + if (sessionId) liveSessions.delete(sessionId); + res.writeHead(204).end(); + return; + } + if (req.method === 'GET') { + // The standalone SSE stream: held open, as a real server holds it. + methods.push('GET'); + res.writeHead(200, { 'Content-Type': 'text/event-stream' }); + res.write(': open\n\n'); + return; + } + + let raw = ''; + req.on('data', (c) => (raw += c)); + req.on('end', () => { + const msg = JSON.parse(raw || '{}'); + methods.push(`POST ${msg.method}`); + + if (msg.method === 'initialize') { + const sessionId = `s${liveSessions.size + 1}`; + liveSessions.add(sessionId); + res.setHeader('mcp-session-id', sessionId); + sse(res, { + jsonrpc: '2.0', + id: msg.id, + result: { + protocolVersion: msg.params.protocolVersion, + capabilities: { prompts: {} }, + serverInfo: { name: 'teardown-probe', version: '0.0.1' } + } + }); + return; + } + if (msg.id === undefined) { + res.writeHead(202).end(); // notification + return; + } + if (msg.method === 'prompts/list' && promptsList === 'hang') { + return; // accepted, never answered + } + if (msg.method === 'prompts/list' && promptsList === 'ok') { + sse(res, { + jsonrpc: '2.0', + id: msg.id, + result: { prompts: [{ name: 'p', description: 'd' }] } + }); + return; + } + sse(res, { + jsonrpc: '2.0', + id: msg.id, + error: { code: -32601, message: `Method not found: ${msg.method}` } + }); + }); + }); + + await new Promise((resolve) => + server.listen(0, '127.0.0.1', resolve) + ); + url = `http://127.0.0.1:${(server.address() as AddressInfo).port}/mcp`; + }); + + afterEach(async () => { + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + }); + + const deleteCount = () => methods.filter((m) => m === 'DELETE').length; + + test('terminates the session when the scenario fails', async () => { + // prompts/list answers -32601, so the scenario throws and its own + // close() — the last statement of its try block — never runs. + promptsList = 'error'; + const result = await runServerConformanceTest( + url, + 'prompts-list', + undefined, + '2025-06-18' + ); + + expect(result.checks.some((c) => c.status === 'FAILURE')).toBe(true); + expect(deleteCount()).toBe(1); + expect([...liveSessions]).toEqual([]); + }, 60000); + + test('terminates the session when the scenario times out', async () => { + // The timeout path abandons the scenario with its promise left pending, so + // no `finally` inside the scenario can ever run: only the runner can close. + promptsList = 'hang'; + const result = await runServerConformanceTest( + url, + 'prompts-list', + undefined, + '2025-06-18', + false, + 2000 + ); + + expect(result.checks.some((c) => c.id === 'scenario-timeout')).toBe(true); + expect(deleteCount()).toBe(1); + expect([...liveSessions]).toEqual([]); + }, 60000); + + test('does not terminate a session twice when the scenario closed it', async () => { + // A second close() re-sends the DELETE, so cleanup must skip connections + // the scenario already closed. + promptsList = 'ok'; + const result = await runServerConformanceTest( + url, + 'prompts-list', + undefined, + '2025-06-18' + ); + + expect(result.checks.every((c) => c.status !== 'FAILURE')).toBe(true); + expect(deleteCount()).toBe(1); + expect([...liveSessions]).toEqual([]); + }, 60000); +}); diff --git a/src/runner/server.ts b/src/runner/server.ts index 21e80ec4..dba5ecff 100644 --- a/src/runner/server.ts +++ b/src/runner/server.ts @@ -7,7 +7,7 @@ import { DRAFT_PROTOCOL_VERSION } from '../types'; import { getClientScenario, isScenarioApplicableAt } from '../scenarios'; -import { connectFor, type RunContext } from '../connection'; +import { connectFor, type Connection, type RunContext } from '../connection'; import { resetWireValidation, wireSchemaChecks @@ -73,6 +73,69 @@ async function runScenarioBounded( ]; } +/** + * Hand a connection to the scenario while remembering it, so the runner can + * terminate whatever the scenario did not. + * + * An explicit `conn.close()` removes it from `open` before delegating, so a + * scenario that closes on its own path is not closed twice — a second + * `close()` sends a second HTTP DELETE for the same session, which is exactly + * the wire noise this is meant to avoid. + * + * `notifications` is exposed as a getter rather than copied: it is appended to + * for the connection's lifetime, and scenarios read it after the requests that + * produce the notifications. + */ +function trackConnection(open: Set, conn: Connection): Connection { + const tracked: Connection = { + request: ( + method: string, + params?: Record, + extraHeaders?: Record + ) => conn.request(method, params, extraHeaders), + get notifications() { + return conn.notifications; + }, + discover: () => conn.discover(), + close: async () => { + open.delete(tracked); + await conn.close(); + } + }; + open.add(tracked); + return tracked; +} + +/** + * Terminate every session the scenario left open. + * + * Scenarios overwhelmingly call `close()` as the last statement of a `try`, + * so a scenario that throws — a `-32601` for a method the server does not + * implement, a failed assertion — leaves its session and its standalone GET + * stream open until the process exits. A server that caps concurrent sessions + * per client then refuses the sessions of later scenarios, and their results + * describe the harness rather than the server. + * + * This lives in the runner, not in a `finally` inside each scenario, because + * `runScenarioBounded` abandons a timed-out scenario with its promise left + * pending on purpose: that scenario's own `finally` never runs, so the runner + * is the only place that can still close the connection. + * + * Failures are swallowed: the cleanup is best effort, and a server that has + * already dropped the session must not turn into a scenario failure. + */ +async function closeOpenConnections(open: Set): Promise { + const leaked = [...open]; + open.clear(); + for (const conn of leaked) { + try { + await conn.close(); + } catch { + // best-effort teardown; the scenario's own result already stands + } + } +} + export async function runServerConformanceTest( serverUrl: string, scenarioName: string, @@ -139,10 +202,24 @@ export async function runServerConformanceTest( `Running client scenario '${scenarioName}' against server: ${serverUrl}` ); + const openConnections = new Set(); + let scenarioFinished = false; const ctx: RunContext = { serverUrl, specVersion: resolvedSpecVersion, - connect: (opts) => connectFor(resolvedSpecVersion)(serverUrl, opts) + connect: async (opts) => { + if (scenarioFinished) { + throw new Error(`Scenario '${scenarioName}' has already finished`); + } + const conn = await connectFor(resolvedSpecVersion)(serverUrl, opts); + // A handshake can finish after the timeout's cleanup sweep. Close it + // here and keep the abandoned scenario from issuing any more probes. + if (scenarioFinished) { + await conn.close(); + throw new Error(`Scenario '${scenarioName}' has already finished`); + } + return trackConnection(openConnections, conn); + } }; resetWireValidation(); const checks = await runScenarioBounded( @@ -150,6 +227,8 @@ export async function runServerConformanceTest( scenarioName, timeout ); + scenarioFinished = true; + await closeOpenConnections(openConnections); checks.push(...wireSchemaChecks(resolvedSpecVersion)); if (resultDir) {